Spring Boot中单个KafkaListener能否配置多个consumer factory监听不同集群
问题结论
@KafkaListener 注解的 containerFactory 属性本身不支持传入数组,无法直接通过指定多工厂的方式将两个监听方法合并为单个注解声明的方法,但可以通过以下两种方案实现仅编写一次业务逻辑的需求:
方案1:抽取公共业务方法(改动最小,兼容性最好)
将重复的业务逻辑抽成独立的公共方法,两个监听方法直接调用该公共方法即可:
@KafkaListener(topics = "test-topic", containerFactory = "otherKafkaListenerContainerFactory") public void listenTestEvent(ConsumerRecord<String, String> record) { handleTestEvent(record); } @KafkaListener(topics = "test-topic", containerFactory = "kafkaListenerContainerFactory") public void listenTestEventServer2(ConsumerRecord<String, String> record) { handleTestEvent(record); } /** * 统一处理test-topic的消息逻辑 */ private void handleTestEvent(ConsumerRecord<String, String> record) { // 你的业务逻辑代码 }
方案2:使用@KafkaListeners容器注解(仅需声明一个方法)
Spring Kafka 2.2及以上版本提供了@KafkaListeners容器注解,可以在同一个方法上绑定多个@KafkaListener配置,分别指定不同的containerFactory即可:
@KafkaListeners({ @KafkaListener(topics = "test-topic", containerFactory = "otherKafkaListenerContainerFactory"), @KafkaListener(topics = "test-topic", containerFactory = "kafkaListenerContainerFactory") }) public void listenTestEvent(ConsumerRecord<String, String> record) { // 直接编写业务逻辑即可,两个集群的消息都会进入该方法处理 }
注意事项
- Spring Boot 2.2及以上版本默认依赖的Spring Kafka版本已支持
@KafkaListeners注解,低版本需要手动升级Spring Kafka依赖 - 两个Kafka集群的同名topic为独立资源,消息拉取互不干扰,都会正常进入方法处理
内容的提问来源于stack exchange,提问作者Milan Rathod
相关产品推荐
相关产品推荐

