From 8fe0adcac904f78298fa7427d5847e53f52f5fe4 Mon Sep 17 00:00:00 2001 From: lingniu Date: Mon, 29 Jun 2026 14:58:19 +0800 Subject: [PATCH] fix: create history consumer processor in app runtime --- ...icleHistoryKafkaConsumerConfiguration.java | 27 +++++++++++++++++++ 1 file changed, 27 insertions(+) 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 8003813c..754747f7 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,11 +1,17 @@ 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; import com.lingniu.ingest.sink.mq.KafkaEnvelopeConsumerRunner; import com.lingniu.ingest.sink.mq.KafkaEnvelopeConsumerWorker; import com.lingniu.ingest.sink.mq.SinkMqProperties; +import com.lingniu.ingest.tdenginehistory.TdengineHistoryWriter; import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.beans.factory.ObjectProvider; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.context.annotation.Bean; @@ -18,6 +24,27 @@ import java.util.Map; @Configuration(proxyBeanMethods = false) public class VehicleHistoryKafkaConsumerConfiguration { + @Bean + @ConditionalOnMissingBean + public TelemetryEnvelopeRecordMapper telemetryEnvelopeRecordMapper() { + return new TelemetryEnvelopeRecordMapper(); + } + + @Bean + @ConditionalOnMissingBean + public EventHistoryEnvelopeIngestor eventHistoryEnvelopeIngestor(EventFileStore store, + TelemetryEnvelopeRecordMapper mapper, + ObjectProvider writer) { + return new EventHistoryEnvelopeIngestor(store, mapper, writer.getIfAvailable()); + } + + @Bean + @ConditionalOnMissingBean(name = "eventHistoryEnvelopeConsumerProcessor") + public EnvelopeConsumerProcessor eventHistoryEnvelopeConsumerProcessor(EventHistoryEnvelopeIngestor ingestor, + EnvelopeDeadLetterSink deadLetterSink) { + return new EnvelopeConsumerProcessor("event-history", ingestor, deadLetterSink); + } + @Bean @ConditionalOnMissingBean @ConditionalOnProperty(prefix = "lingniu.ingest.sink.mq.consumer", name = "enabled", havingValue = "true")