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

Spring Boot使用KafkaListener消费2个Topic仅处理启用Topic的方法

Spring Boot 整合Kafka实现多Broker分Topic按配置启停消费方案

配置层定义

首先在Spring Boot配置文件(application.yml/application.properties)中维护两个Kafka集群的连接参数、Topic名称、启停标记:

kafka:
  cluster-a:
    bootstrap-servers: <集群A地址>:9092
    topic: <集群A业务Topic名>
    enabled: true # 标记为启用则消费该集群消息
  cluster-b:
    bootstrap-servers: <集群B地址>:9092
    topic: <集群B业务Topic名>
    enabled: false # 标记为禁用则不处理该集群消息

通过@ConfigurationProperties将两类集群配置绑定为独立的配置对象,避免硬编码。

多集群消费者容器配置

由于两个Topic归属不同独立Kafka Broker(跨集群),默认的Kafka消费者工厂仅能对接单个集群,需要为每个集群单独创建对应的监听容器工厂,加载对应集群的连接、消费者参数:

@Configuration
public class MultiKafkaConfig {

    @Bean
    @ConfigurationProperties(prefix = "kafka.cluster-a")
    public KafkaProperties clusterAKafkaProps() {
        return new KafkaProperties();
    }

    @Bean
    @ConfigurationProperties(prefix = "kafka.cluster-b")
    public KafkaProperties clusterBKafkaProps() {
        return new KafkaProperties();
    }

    @Bean("clusterAFactory")
    public ConcurrentKafkaListenerContainerFactory<?, ?> clusterAContainerFactory(
            @Qualifier("clusterAKafkaProps") KafkaProperties props) {
        ConcurrentKafkaListenerContainerFactory<Object, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(new DefaultKafkaConsumerFactory<>(props.buildConsumerProperties()));
        // 如需手动提交偏移量、并发消费配置,可在此处统一设置
        return factory;
    }

    @Bean("clusterBFactory")
    public ConcurrentKafkaListenerContainerFactory<?, ?> clusterBContainerFactory(
            @Qualifier("clusterBKafkaProps") KafkaProperties props) {
        ConcurrentKafkaListenerContainerFactory<Object, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(new DefaultKafkaConsumerFactory<>(props.buildConsumerProperties()));
        return factory;
    }
}

消费逻辑实现

根据配置是否需要热更新,可选择两种实现方式:

方式1:配置静态生效(无无效连接)

直接通过@KafkaListener的autoStartup属性读取配置中的启停标记,标记为禁用的监听器不会启动,不会和对应Broker建立消费连接,资源占用最低,适合配置不需要频繁变动的场景:

@Component
@Slf4j
public class BizMsgConsumer {

    @KafkaListener(
            topics = "${kafka.cluster-a.topic}",
            containerFactory = "clusterAFactory",
            groupId = "biz-sync-group",
            autoStartup = "${kafka.cluster-a.enabled:false}"
    )
    public void consumeClusterA(ConsumerRecord<?, String> record) {
        doProcess(record.value(), "clusterA");
    }

    @KafkaListener(
            topics = "${kafka.cluster-b.topic}",
            containerFactory = "clusterBFactory",
            groupId = "biz-sync-group",
            autoStartup = "${kafka.cluster-b.enabled:false}"
    )
    public void consumeClusterB(ConsumerRecord<?, String> record) {
        doProcess(record.value(), "clusterB");
    }

    private void doProcess(String msg, String clusterTag) {
        // 统一业务处理逻辑
        log.info("处理来自{}的业务消息:{}", clusterTag, msg);
    }
}

两个独立集群的消费组即使重名也不会冲突,消费偏移量由各自集群独立维护。

方式2:配置热生效(支持运行时动态调整)

如果使用配置中心、需要不重启应用切换消费的Topic,就保持两个监听器始终启动,在消费逻辑内判断启停标记,标记为禁用时直接提交偏移量跳过处理即可:

@Component
@Slf4j
public class BizMsgConsumer {
    // 配置项支持配置中心热更新
    @Value("${kafka.cluster-a.enabled:false}")
    private boolean clusterAEnable;
    @Value("${kafka.cluster-b.enabled:false}")
    private boolean clusterBEnable;

    @KafkaListener(
            topics = "${kafka.cluster-a.topic}",
            containerFactory = "clusterAFactory",
            groupId = "biz-sync-group"
    )
    public void consumeClusterA(ConsumerRecord<?, String> record, Acknowledgment ack) {
        if (!clusterAEnable) {
            ack.acknowledge();
            return;
        }
        doProcess(record.value(), "clusterA");
        ack.acknowledge();
    }

    @KafkaListener(
            topics = "${kafka.cluster-b.topic}",
            containerFactory = "clusterBFactory",
            groupId = "biz-sync-group"
    )
    public void consumeClusterB(ConsumerRecord<?, String> record, Acknowledgment ack) {
        if (!clusterBEnable) {
            ack.acknowledge();
            return;
        }
        doProcess(record.value(), "clusterB");
        ack.acknowledge();
    }

    private void doProcess(String msg, String clusterTag) {
        log.info("处理来自{}的业务消息:{}", clusterTag, msg);
    }
}

特殊场景说明

如果两个Topic归属同一个Kafka集群,不需要创建多套容器工厂,直接在单个@KafkaListener的topics属性中配置两个Topic名称,消费时通过record.topic()获取当前消息所属Topic,结合配置的启用标记判断是否处理即可。


内容的提问来源于stack exchange,提问作者Mohammad Gouse

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 03:42:17