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
相关产品推荐
相关产品推荐

