You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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. 我的实现存在哪些问题?
  2. 哪里能找到相关参考文档?

解答

方案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操作,否则无法消费消息;更新订阅时要确保线程安全,避免并发问题。

可行实现步骤

  1. 将CustomKafkaConsumer改为Quarkus的@ApplicationScoped Bean,利用@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) {}
}
  1. 注意事项:
    • 订阅更新时必须加锁,避免消费线程和更新线程同时操作消费者导致并发问题。
    • 确保TopicListApiClient能正确处理API调用异常,避免更新任务崩溃。
    • 可以根据需求调整poll的超时时间和更新周期。

参考文档提示

Quarkus官方文档中关于Kafka消费者的原生使用、Bean生命周期管理、定时任务的章节是核心参考:

  • Quarkus的Bean生命周期:@PostConstruct、@PreDestroy的用法
  • Quarkus定时任务(如果不用ScheduledExecutorService,也可以用Quarkus的@Scheduled注解)
  • Apache Kafka官方文档中关于KafkaConsumer.subscribe()动态更新订阅的说明

内容的提问来源于stack exchange,提问作者Lorenzo De Francesco

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.28 20:05:02