Quarkus微服务实现Kafka多主题动态订阅问题求助
Quarkus实现Kafka动态订阅多Topic问题
我想用Java Quarkus开发微服务,实现运行时动态订阅多个Kafka Topic的功能:每日定期从API读取Topic列表,更新要消费的Topic。尝试了三种方案但都有问题,具体如下:
方案1:使用quarkus.kafka-streams
代码实现:
@ApplicationScoped @Slf4j public class KafkaTopology { @Inject IResourceFileService resourceFileService; @Produces public Topology buildTopology() { log.info("Start building topology - Read from file the topics"); List<TopicData> data = resourceFileService.getTopics(); List<String> topics = data.stream().map(single -> single.getTopic()).collect(Collectors.toList()); StreamsBuilder builder = new StreamsBuilder(); log.info("Topic list is" + topics); builder.stream(topics, Consumed.with(Serdes.String(), Serdes.String())).peek((topic, value) ->{ log.info("Peek " + topic + "Value: " + value); }); return builder.build(); } }
对应的application.properties配置:
quarkus.kafka-streams.bootstrap-servers=xxxx quarkus.kafka-streams.security.protocol=SASL_SSL quarkus.kafka-streams.sasl.mechanism=PLAIN quarkus.kafka-streams.sasl.jaas-config=xxxx
运行时错误:
2023-03-24 12:34:22,086 WARN [org.apa.kaf.str.int.met.ClientMetrics] (Quarkus Main Thread) Error while loading kafka-streams-version.properties: java.lang.NullPointerException: inStream parameter is null WARN [org.apa.kaf.cli.con.ConsumerConfig] (Quarkus Main Thread) These configurations '[streams.metadata.max.age.ms, bootstrap-servers, streams.commit.interval.ms, sasl.jaas-config, streams.auto.offset.reset, streams.cache.max.bytes.buffering, application-id]' were supplied but are not used yet. Exception in thread "pool-12-thread-1" java.lang.IllegalArgumentException: Missing list of topics
注:实际Topic列表不为空,但仍报错。
方案2:使用已废弃的旧版quarkus.kafka-streams
代码实现:
@ApplicationScoped public class MyTopology { // 该列表是动态的,来自事件总线 List<String> topics = List.of("topicA", "topicB"); @Produces public Topology buildTopology() { StreamsBuilder builder = new StreamsBuilder(); for(String topic : topics){ builder.stream(topic, Consumed.with(Serdes.String(), Serdes.String())) .mapValues(v -> { System.out.println("Topic message elaboration" + v); return v; }); } return builder.build(); } }
方案3:使用org.apache.kafka.clients.consumer.KafkaConsumer
代码实现:
public class CustomKafkaConsumer implements ConsumerRebalanceListener { List<String> topics = List.of("topicA", "topicB"); private Consumer<String, String> consumer; public void start() { Properties props = new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "xx"); props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "xx"); props.put("sasl.jaas.config", "22"); props.put(ConsumerConfig.GROUP_ID_CONFIG, "test"); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "true"); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer"); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringDeserializer"); this.consumer = new KafkaConsumer<>(props); this.consumer.subscribe(topics, this); } public void stop() { this.consumer.close(); } public void poll(Duration timeout) { this.consumer.poll(timeout).forEach(record -> { // 处理消息 System.out.println("TOPIC" + record.topic()); System.out.println("VALUE" + record.value()); System.out.println("KEY" + record.key()); }); } @Override public void onPartitionsRevoked(Collection<TopicPartition> partitions) { // 无操作 } @Override public void onPartitionsAssigned(Collection<TopicPartition> partitions) { // 无操作 } }
问题
目前未找到可行解决方案,请问:
- 我的实现存在哪些问题?
- 哪里能找到相关参考文档?
解答
方案1错误分析
报错Missing list of topics的核心原因:Quarkus启动Kafka Streams拓扑时,会在应用启动阶段提前初始化配置,而你的resourceFileService.getTopics()可能在拓扑构建时并未真正加载到Topic列表(比如文件读取延迟、依赖初始化顺序问题)。另外,Kafka Streams的拓扑一旦构建完成就无法动态修改,天生不支持运行时动态添加/移除Topic,这和你“每日更新订阅列表”的需求冲突,因此这个方案从设计上就不符合需求。
方案2问题
旧版quarkus.kafka-streams已废弃,且同样存在Kafka Streams拓扑不可动态修改的问题,即便能启动,后续也无法实现Topic列表的动态更新,直接排除。
方案3优化方向
使用原生KafkaConsumer是可行的方向,但你的实现缺少以下关键部分:
- Quarkus生命周期管理:没有把
CustomKafkaConsumer纳入Quarkus的Bean管理,无法利用Quarkus的启动/停止钩子,也无法依赖注入其他服务(比如读取Topic列表的API服务)。 - 动态更新订阅逻辑:没有实现定期从API拉取Topic列表并更新订阅的逻辑,当前代码只是启动时订阅固定列表。
- 线程安全与轮询机制:需要单独启动线程进行持续poll操作,否则无法消费消息;更新订阅时要确保线程安全,避免并发问题。
可行实现步骤
- 将
CustomKafkaConsumer改为Quarkus的@ApplicationScopedBean,利用@PostConstruct启动消费线程,@PreDestroy停止消费:
@ApplicationScoped public class CustomKafkaConsumer implements ConsumerRebalanceListener { private Consumer<String, String> consumer; private volatile List<String> currentTopics = Collections.emptyList(); private final ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor(); @Inject TopicListApiClient topicListApiClient; // 假设这是调用API获取Topic列表的服务 @PostConstruct public void init() { // 初始化消费者配置 Properties props = new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "xx"); props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "xx"); props.put("sasl.jaas.config", "22"); props.put(ConsumerConfig.GROUP_ID_CONFIG, "test"); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "true"); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); consumer = new KafkaConsumer<>(props); // 启动每日更新Topic的任务 scheduler.scheduleAtFixedRate(this::updateSubscribedTopics, 0, 1, TimeUnit.DAYS); // 启动消费线程 new Thread(this::consumeMessages).start(); } private void updateSubscribedTopics() { List<String> newTopics = topicListApiClient.fetchLatestTopics(); if (!newTopics.equals(currentTopics)) { synchronized (consumer) { consumer.subscribe(newTopics, this); currentTopics = newTopics; } } } private void consumeMessages() { while (!Thread.currentThread().isInterrupted()) { synchronized (consumer) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100)); records.forEach(record -> { // 处理消息逻辑 System.out.printf("TOPIC: %s, KEY: %s, VALUE: %s%n", record.topic(), record.key(), record.value()); }); } } } @PreDestroy public void shutdown() { scheduler.shutdown(); if (consumer != null) { consumer.close(); } } @Override public void onPartitionsRevoked(Collection<TopicPartition> partitions) {} @Override public void onPartitionsAssigned(Collection<TopicPartition> partitions) {} }
- 注意事项:
- 订阅更新时必须加锁,避免消费线程和更新线程同时操作消费者导致并发问题。
- 确保
TopicListApiClient能正确处理API调用异常,避免更新任务崩溃。 - 可以根据需求调整poll的超时时间和更新周期。
参考文档提示
Quarkus官方文档中关于Kafka消费者的原生使用、Bean生命周期管理、定时任务的章节是核心参考:
- Quarkus的Bean生命周期:
@PostConstruct、@PreDestroy的用法 - Quarkus定时任务(如果不用
ScheduledExecutorService,也可以用Quarkus的@Scheduled注解) - Apache Kafka官方文档中关于
KafkaConsumer.subscribe()动态更新订阅的说明
内容的提问来源于stack exchange,提问作者Lorenzo De Francesco
相关产品推荐
相关产品推荐

