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 6cca1c95..17dcd880 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 @@ -17,6 +17,7 @@ import com.lingniu.ingest.protocol.gb32960.config.Gb32960AutoConfiguration; import com.lingniu.ingest.protocol.gb32960.inbound.Gb32960NettyServer; import com.lingniu.ingest.sink.mq.KafkaEnvelopeDeadLetterSink; import com.lingniu.ingest.sink.mq.KafkaEventSink; +import com.lingniu.ingest.sink.mq.KafkaEnvelopeConsumerRunner; import com.lingniu.ingest.sink.mq.SinkMqProperties; import com.lingniu.ingest.sink.mq.proto.ParseStatusProto; import com.lingniu.ingest.sink.mq.proto.RawArchiveRef; @@ -143,6 +144,35 @@ class VehicleHistoryAppCompositionTest { }); } + @Test + void createsKafkaConsumerRunnerWhenTdengineHistoryRuntimeConsumesKafka() { + new ApplicationContextRunner() + .withConfiguration(AutoConfigurations.of( + TdengineHistoryAutoConfiguration.class, + SinkMqAutoConfiguration.class, + EventHistoryAutoConfiguration.class)) + .withUserConfiguration(VehicleHistoryKafkaConsumerConfiguration.class) + .withAllowBeanDefinitionOverriding(true) + .withBean("kafkaProducer", KafkaProducer.class, VehicleHistoryAppCompositionTest::kafkaProducer) + .withPropertyValues( + "lingniu.ingest.tdengine-history.enabled=true", + "lingniu.ingest.tdengine-history.database=vehicle_history_test", + "lingniu.ingest.event-history.enabled=true", + "lingniu.ingest.sink.mq.enabled=true", + "lingniu.ingest.sink.mq.type=kafka", + "lingniu.ingest.sink.mq.bootstrap-servers=localhost:9092", + "lingniu.ingest.sink.mq.consumer.enabled=true", + "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryGb32960EnvelopeConsumerProcessor.enabled=true", + "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryGb32960EnvelopeConsumerProcessor.group-id=history-gb32960", + "lingniu.ingest.sink.mq.consumer.bindings.eventHistoryGb32960EnvelopeConsumerProcessor.topics[0]=vehicle.event.gb32960.v1") + .run(context -> { + assertThat(context).hasSingleBean(EventHistoryEnvelopeIngestor.class); + assertThat(context.getBeanNamesForType(EnvelopeConsumerProcessor.class)) + .contains("eventHistoryEnvelopeConsumerProcessor"); + assertThat(context).hasSingleBean(KafkaEnvelopeConsumerRunner.class); + }); + } + @Test void historyIngestorIgnoresEventFileStoreBeanAndWritesTdengineOnly() { EventFileStore eventFileStore = mock(EventFileStore.class); diff --git a/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/config/EventHistoryAutoConfiguration.java b/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/config/EventHistoryAutoConfiguration.java index 03b2b838..a72ac3b9 100644 --- a/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/config/EventHistoryAutoConfiguration.java +++ b/modules/services/event-history-service/src/main/java/com/lingniu/ingest/eventhistory/config/EventHistoryAutoConfiguration.java @@ -15,6 +15,7 @@ import com.lingniu.ingest.eventhistory.TelemetryFieldHistoryController; import com.lingniu.ingest.eventhistory.TelemetryEnvelopeRecordMapper; import com.lingniu.ingest.protocol.gb32960.codec.Gb32960MessageDecoder; import com.lingniu.ingest.protocol.gb32960.config.Gb32960AutoConfiguration; +import com.lingniu.ingest.sink.mq.SinkMqAutoConfiguration; import com.lingniu.ingest.tdenginehistory.config.TdengineHistoryAutoConfiguration; import com.lingniu.ingest.tdenginehistory.TdengineHistoryReader; import com.lingniu.ingest.tdenginehistory.TdengineHistoryWriter; @@ -45,6 +46,7 @@ import java.nio.file.Path; @AutoConfiguration @AutoConfigureAfter({ Gb32960AutoConfiguration.class, + SinkMqAutoConfiguration.class, TdengineHistoryAutoConfiguration.class }) @ConditionalOnProperty(prefix = "lingniu.ingest.event-history", name = "enabled", havingValue = "true")