fix: start history kafka consumer after processor setup

This commit is contained in:
lingniu
2026-06-29 14:55:44 +08:00
parent 06423db8d4
commit 00eee9a769
2 changed files with 41 additions and 0 deletions

View File

@@ -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<KafkaEnvelopeConsumerWorker> 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());
}
}

View File

@@ -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<String, EnvelopeConsumerProcessor> processors,
SinkMqProperties props) {