如何实现单个Kafka监听器隶属于多个Kafka容器?
如何让一个Kafka监听器隶属于多个容器?
以下内容来自Spring Kafka官方文档:
获取消费者 group.id
当在多个容器中运行相同的监听器代码时,能够确定记录来自哪个容器(通过其group.id消费者属性标识)可能会很有用。
我原本认为一段Java Kafka监听器代码只能隶属于一个工厂,要么是默认工厂,要么是@KafkaListener注解中指定的工厂。但结合上述文档描述,显然还有更多实现方式。请问如何让一个监听器隶属于多个Kafka容器?
实现方式
1. 给同一个监听器方法添加多个@KafkaListener注解
直接在同一个方法上多次标注@KafkaListener,每个注解指定不同的groupId、监听主题,甚至不同的容器工厂。Spring会为每个注解创建独立的消费者容器,全部绑定到同一个监听器方法。
示例代码:
@Component public class MultiContainerListener { @KafkaListener(topics = "topic-1", groupId = "group-1") @KafkaListener(topics = "topic-2", groupId = "group-2", containerFactory = "customKafkaListenerContainerFactory") public void listen(ConsumerRecord<String, String> record) { // 获取当前容器的group.id String groupId = new String(record.headers().lastHeader("kafka_groupId").value()); System.out.printf("Received message from group %s: %s%n", groupId, record.value()); } }
2. 手动创建多个MessageListenerContainer并绑定同一个监听器实例
如果需要更灵活的自定义配置,可以手动创建多个消费者容器,将同一个监听器实例注入到每个容器中,每个容器配置不同的group.id和监听主题。
示例代码:
@Configuration public class MultiContainerConfig { @Bean public ConsumerFactory<String, String> consumerFactory() { return new DefaultKafkaConsumerFactory<>(Map.of( ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092", ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class, ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class )); } @Bean public MessageListenerContainer container1(ConsumerFactory<String, String> consumerFactory, MultiContainerListener listener) { ContainerProperties props = new ContainerProperties("topic-1"); props.setGroupId("group-1"); props.setMessageListener((MessageListener<String, String>) listener::listen); return new ConcurrentMessageListenerContainer<>(consumerFactory, props); } @Bean public MessageListenerContainer container2(ConsumerFactory<String, String> consumerFactory, MultiContainerListener listener) { ContainerProperties props = new ContainerProperties("topic-2"); props.setGroupId("group-2"); props.setMessageListener((MessageListener<String, String>) listener::listen); return new ConcurrentMessageListenerContainer<>(consumerFactory, props); } } @Component public class MultiContainerListener { public void listen(ConsumerRecord<String, String> record) { String groupId = new String(record.headers().lastHeader("kafka_groupId").value()); System.out.printf("Received message from group %s: %s%n", groupId, record.value()); } }
内容的提问来源于stack exchange,提问作者samshers
相关产品推荐
相关产品推荐

