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

如何实现单个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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 07:22:39