From 00eee9a7698d9da1f7b67ca481b63ce307e7e6b4 Mon Sep 17 00:00:00 2001 From: lingniu Date: Mon, 29 Jun 2026 14:55:44 +0800 Subject: [PATCH] fix: start history kafka consumer after processor setup --- ...icleHistoryKafkaConsumerConfiguration.java | 39 +++++++++++++++++++ .../sink/mq/SinkMqAutoConfiguration.java | 2 + 2 files changed, 41 insertions(+) create mode 100644 modules/apps/vehicle-history-app/src/main/java/com/lingniu/ingest/historyapp/VehicleHistoryKafkaConsumerConfiguration.java 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 new file mode 100644 index 00000000..8003813c --- /dev/null +++ b/modules/apps/vehicle-history-app/src/main/java/com/lingniu/ingest/historyapp/VehicleHistoryKafkaConsumerConfiguration.java @@ -0,0 +1,39 @@ +package com.lingniu.ingest.historyapp; + +import com.lingniu.ingest.api.consumer.EnvelopeConsumerProcessor; +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 org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; + +import java.time.Duration; +import java.util.List; +import java.util.Map; + +@Configuration(proxyBeanMethods = false) +public class VehicleHistoryKafkaConsumerConfiguration { + + @Bean + @ConditionalOnMissingBean + @ConditionalOnProperty(prefix = "lingniu.ingest.sink.mq.consumer", name = "enabled", havingValue = "true") + public KafkaEnvelopeConsumerRunner vehicleHistoryKafkaEnvelopeConsumerRunner( + @Qualifier("eventHistoryEnvelopeConsumerProcessor") EnvelopeConsumerProcessor processor, + SinkMqProperties props) { + List workers = new KafkaEnvelopeConsumerFactory().createWorkers( + Map.of("eventHistoryEnvelopeConsumerProcessor", processor), + props); + if (workers.isEmpty()) { + throw new IllegalStateException("no vehicle history kafka consumer workers created; check consumer bindings"); + } + return new KafkaEnvelopeConsumerRunner( + workers, + Duration.ofMillis(props.getConsumer().getPollTimeoutMillis()), + Duration.ofMillis(props.getConsumer().getLoopBackoffMillis()), + props.getConsumer().isAutoStartup()); + } +} diff --git a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqAutoConfiguration.java b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqAutoConfiguration.java index 5763f6b4..53185392 100644 --- a/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqAutoConfiguration.java +++ b/modules/sinks/sink-mq/src/main/java/com/lingniu/ingest/sink/mq/SinkMqAutoConfiguration.java @@ -8,6 +8,7 @@ import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.common.serialization.ByteArraySerializer; import org.apache.kafka.common.serialization.StringSerializer; import org.springframework.boot.autoconfigure.AutoConfiguration; +import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.context.properties.EnableConfigurationProperties; @@ -101,6 +102,7 @@ public class SinkMqAutoConfiguration { @Bean @ConditionalOnMissingBean + @ConditionalOnBean(EnvelopeConsumerProcessor.class) @ConditionalOnProperty(prefix = "lingniu.ingest.sink.mq.consumer", name = "enabled", havingValue = "true") public KafkaEnvelopeConsumerRunner kafkaEnvelopeConsumerRunner(Map processors, SinkMqProperties props) {