From d5ef061cca6aacb1d945341b96468123ae63f461 Mon Sep 17 00:00:00 2001 From: lingniu Date: Wed, 1 Jul 2026 04:15:43 +0800 Subject: [PATCH] refactor: keep history consumer tdengine only --- ...icleHistoryKafkaConsumerConfiguration.java | 10 +-- .../VehicleHistoryAppCompositionTest.java | 74 +++++++++++++++++++ 2 files changed, 76 insertions(+), 8 deletions(-) diff --git a/modules/apps/vehicle-history-app/src/main/java/com/lingniu/ingest/historyapp/VehicleHistoryKafkaConsumerConfiguration.java b/modules/apps/vehicle-history-app/src/main/java/com/lingniu/ingest/historyapp/VehicleHistoryKafkaConsumerConfiguration.java index ac993058..0226b3de 100644 --- a/modules/apps/vehicle-history-app/src/main/java/com/lingniu/ingest/historyapp/VehicleHistoryKafkaConsumerConfiguration.java +++ b/modules/apps/vehicle-history-app/src/main/java/com/lingniu/ingest/historyapp/VehicleHistoryKafkaConsumerConfiguration.java @@ -2,7 +2,6 @@ package com.lingniu.ingest.historyapp; import com.lingniu.ingest.api.consumer.EnvelopeConsumerProcessor; import com.lingniu.ingest.api.consumer.EnvelopeDeadLetterSink; -import com.lingniu.ingest.eventfilestore.EventFileStore; import com.lingniu.ingest.eventhistory.EventHistoryEnvelopeIngestor; import com.lingniu.ingest.eventhistory.TelemetryEnvelopeRecordMapper; import com.lingniu.ingest.sink.mq.KafkaEnvelopeConsumerFactory; @@ -33,17 +32,12 @@ public class VehicleHistoryKafkaConsumerConfiguration { @Bean @ConditionalOnMissingBean - public EventHistoryEnvelopeIngestor eventHistoryEnvelopeIngestor(ObjectProvider store, - TelemetryEnvelopeRecordMapper mapper, + public EventHistoryEnvelopeIngestor eventHistoryEnvelopeIngestor(TelemetryEnvelopeRecordMapper mapper, ObjectProvider writer, @Value("${lingniu.ingest.tdengine-history.telemetry-fields-enabled:false}") boolean telemetryFieldsEnabled) { TdengineHistoryWriter tdengineWriter = writer.getIfAvailable(); - EventFileStore eventFileStore = store.getIfAvailable(); - if (eventFileStore == null) { - return new EventHistoryEnvelopeIngestor(mapper, tdengineWriter, telemetryFieldsEnabled); - } - return new EventHistoryEnvelopeIngestor(eventFileStore, mapper, tdengineWriter, telemetryFieldsEnabled); + return new EventHistoryEnvelopeIngestor(mapper, tdengineWriter, telemetryFieldsEnabled); } @Bean diff --git a/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/VehicleHistoryAppCompositionTest.java b/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/VehicleHistoryAppCompositionTest.java index 01f22eb0..e3749e84 100644 --- a/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/VehicleHistoryAppCompositionTest.java +++ b/modules/apps/vehicle-history-app/src/test/java/com/lingniu/ingest/historyapp/VehicleHistoryAppCompositionTest.java @@ -1,6 +1,7 @@ package com.lingniu.ingest.historyapp; import com.lingniu.ingest.api.consumer.EnvelopeDeadLetterSink; +import com.lingniu.ingest.eventfilestore.EventFileStore; import com.lingniu.ingest.eventhistory.EventHistoryEnvelopeIngestor; import com.lingniu.ingest.eventhistory.Gb32960DecodedFrameService; import com.lingniu.ingest.eventhistory.Gb32960FrameController; @@ -17,9 +18,14 @@ import com.lingniu.ingest.sink.archive.RawArchiveEventSink; import com.lingniu.ingest.sink.archive.config.SinkArchiveAutoConfiguration; import com.lingniu.ingest.sink.mq.KafkaEnvelopeDeadLetterSink; import com.lingniu.ingest.sink.mq.KafkaEventSink; +import com.lingniu.ingest.sink.mq.proto.ParseStatusProto; +import com.lingniu.ingest.sink.mq.proto.RawArchiveRef; +import com.lingniu.ingest.sink.mq.proto.RawFrameFactPayload; +import com.lingniu.ingest.sink.mq.proto.VehicleEnvelope; import com.lingniu.ingest.sink.mq.SinkMqAutoConfiguration; import com.lingniu.ingest.tdenginehistory.TdengineHistorySchema; import com.lingniu.ingest.tdenginehistory.TdengineHistoryWriter; +import com.lingniu.ingest.tdenginehistory.TdengineRawFrameRow; import com.lingniu.ingest.tdenginehistory.config.TdengineHistoryAutoConfiguration; import org.apache.kafka.clients.producer.KafkaProducer; import org.junit.jupiter.api.Test; @@ -28,9 +34,14 @@ import org.springframework.boot.autoconfigure.AutoConfigurations; import org.springframework.boot.test.context.runner.ApplicationContextRunner; import java.nio.file.Path; +import java.util.ArrayList; +import java.util.List; import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; class VehicleHistoryAppCompositionTest { @@ -99,8 +110,71 @@ class VehicleHistoryAppCompositionTest { }); } + @Test + void historyIngestorIgnoresEventFileStoreBeanAndWritesTdengineOnly() { + EventFileStore eventFileStore = mock(EventFileStore.class); + CapturingTdengineWriter writer = new CapturingTdengineWriter(); + + new ApplicationContextRunner() + .withUserConfiguration(VehicleHistoryKafkaConsumerConfiguration.class) + .withBean(EventFileStore.class, () -> eventFileStore) + .withBean(TdengineHistoryWriter.class, () -> writer) + .withBean(EnvelopeDeadLetterSink.class, () -> mock(EnvelopeDeadLetterSink.class)) + .withPropertyValues("lingniu.ingest.sink.mq.consumer.enabled=false") + .run(context -> { + EventHistoryEnvelopeIngestor ingestor = context.getBean(EventHistoryEnvelopeIngestor.class); + + ingestor.tryIngest(rawFrameEnvelope().toByteArray()); + + verify(eventFileStore, never()).appendAll(any()); + assertThat(writer.rawFrames) + .extracting(TdengineRawFrameRow::frameId) + .containsExactly("raw-frame-1"); + }); + } + @SuppressWarnings("unchecked") private static KafkaProducer kafkaProducer() { return mock(KafkaProducer.class); } + + private static VehicleEnvelope rawFrameEnvelope() { + return VehicleEnvelope.newBuilder() + .setSchemaVersion("1.0") + .setEventId("raw-event-1") + .setVin("VINRAW001") + .setSource("JT808") + .setEventTimeMs(1_782_112_400_000L) + .setIngestTimeMs(1_782_112_401_000L) + .setRawArchive(RawArchiveRef.newBuilder() + .setUri("archive://jt808/2026/06/29/raw-frame-1.bin") + .setSizeBytes(68)) + .setRawFrameFact(RawFrameFactPayload.newBuilder() + .setFrameId("raw-frame-1") + .setVehicleKey("jt808:013800000000") + .setVin("VINRAW001") + .setPhone("013800000000") + .setMessageId(0x0200) + .setRawUri("archive://jt808/2026/06/29/raw-frame-1.bin") + .setRawSizeBytes(68) + .setParseStatus(ParseStatusProto.PARSE_STATUS_SUCCEEDED)) + .build(); + } + + private static final class CapturingTdengineWriter implements TdengineHistoryWriter { + private final List rawFrames = new ArrayList<>(); + + @Override + public void appendRawFrames(List rows) { + rawFrames.addAll(rows); + } + + @Override + public void appendLocations(List rows) { + } + + @Override + public void appendTelemetryFields(List rows) { + } + } }