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")