From b841141fedb3998f682632f38521e70d931c146f Mon Sep 17 00:00:00 2001 From: lingniu Date: Wed, 1 Jul 2026 04:54:11 +0800 Subject: [PATCH] refactor: centralize history raw consumer processor --- .../VehicleHistoryKafkaConsumerConfiguration.java | 9 --------- .../historyapp/VehicleHistoryAppCompositionTest.java | 3 ++- .../config/EventHistoryAutoConfiguration.java | 8 ++++++++ .../config/EventHistoryAutoConfigurationTest.java | 8 ++++++-- 4 files changed, 16 insertions(+), 12 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 e38e7577..53f53529 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 @@ -1,8 +1,6 @@ package com.lingniu.ingest.historyapp; import com.lingniu.ingest.api.consumer.EnvelopeConsumerProcessor; -import com.lingniu.ingest.api.consumer.EnvelopeDeadLetterSink; -import com.lingniu.ingest.eventhistory.EventHistoryEnvelopeIngestor; import com.lingniu.ingest.sink.mq.KafkaEnvelopeConsumerFactory; import com.lingniu.ingest.sink.mq.KafkaEnvelopeConsumerRunner; import com.lingniu.ingest.sink.mq.KafkaEnvelopeConsumerWorker; @@ -20,13 +18,6 @@ import java.util.Map; @Configuration(proxyBeanMethods = false) public class VehicleHistoryKafkaConsumerConfiguration { - @Bean - @ConditionalOnMissingBean(name = "eventHistoryRawEnvelopeConsumerProcessor") - public EnvelopeConsumerProcessor eventHistoryRawEnvelopeConsumerProcessor(EventHistoryEnvelopeIngestor ingestor, - EnvelopeDeadLetterSink deadLetterSink) { - return new EnvelopeConsumerProcessor("event-history-raw", ingestor, deadLetterSink); - } - @Bean @ConditionalOnMissingBean @ConditionalOnProperty(prefix = "lingniu.ingest.sink.mq.consumer", name = "enabled", havingValue = "true") 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 bc5daab0..8b804d08 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 @@ -61,7 +61,8 @@ class VehicleHistoryAppCompositionTest { assertThat(beanMethods) .doesNotContain("telemetryEnvelopeRecordMapper") .doesNotContain("eventHistoryEnvelopeIngestor") - .doesNotContain("eventHistoryEnvelopeConsumerProcessor"); + .doesNotContain("eventHistoryEnvelopeConsumerProcessor") + .doesNotContain("eventHistoryRawEnvelopeConsumerProcessor"); } @Test 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 dfe1b504..77f8a94d 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 @@ -87,6 +87,14 @@ public class EventHistoryAutoConfiguration { return new EnvelopeConsumerProcessor("event-history", ingestor, deadLetterSink); } + @Bean + @ConditionalOnBean({EventHistoryEnvelopeIngestor.class, EnvelopeDeadLetterSink.class}) + @ConditionalOnMissingBean(name = "eventHistoryRawEnvelopeConsumerProcessor") + public EnvelopeConsumerProcessor eventHistoryRawEnvelopeConsumerProcessor(EventHistoryEnvelopeIngestor ingestor, + EnvelopeDeadLetterSink deadLetterSink) { + return new EnvelopeConsumerProcessor("event-history-raw", ingestor, deadLetterSink); + } + @Bean @ConditionalOnBean(EventFileStore.class) @ConditionalOnMissingBean diff --git a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/config/EventHistoryAutoConfigurationTest.java b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/config/EventHistoryAutoConfigurationTest.java index 6fbf55b0..6b52d61d 100644 --- a/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/config/EventHistoryAutoConfigurationTest.java +++ b/modules/services/event-history-service/src/test/java/com/lingniu/ingest/eventhistory/config/EventHistoryAutoConfigurationTest.java @@ -43,7 +43,9 @@ class EventHistoryAutoConfigurationTest { assertThat(context).hasSingleBean(TelemetryEnvelopeRecordMapper.class); assertThat(context).hasSingleBean(EventHistoryEnvelopeIngestor.class); assertThat(context.getBeansOfType(EnvelopeConsumerProcessor.class)) - .containsOnlyKeys("eventHistoryEnvelopeConsumerProcessor"); + .containsOnlyKeys( + "eventHistoryEnvelopeConsumerProcessor", + "eventHistoryRawEnvelopeConsumerProcessor"); assertThat(context).hasSingleBean(EventHistoryController.class); assertThat(context).doesNotHaveBean(Gb32960DecodedFrameService.class); assertThat(context).doesNotHaveBean(Gb32960FrameController.class); @@ -104,7 +106,9 @@ class EventHistoryAutoConfigurationTest { assertThat(context).hasSingleBean(TelemetryEnvelopeRecordMapper.class); assertThat(context).hasSingleBean(EventHistoryEnvelopeIngestor.class); assertThat(context.getBeansOfType(EnvelopeConsumerProcessor.class)) - .containsOnlyKeys("eventHistoryEnvelopeConsumerProcessor"); + .containsOnlyKeys( + "eventHistoryEnvelopeConsumerProcessor", + "eventHistoryRawEnvelopeConsumerProcessor"); }); }