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