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

Kafka多消费者实例无法接收消息问题排查求助

问题:多实例Kafka消费者无法共同接收消息(同Group ID下)

我有一个同时作为生产者和消费者的应用,希望启动多个实例时每个实例都能接收消息。但目前只有使用不同Group ID的实例才能接收消息;我已设置分区数多于实例数,但问题仍未解决。相关配置代码如下:

生产者配置

@EnableKafka
@Configuration
public class KafkaProducerConfigForDepartment {


    @Value(value = "${kafka.bootstrapAddress}")
    private String bootstrapAddress;

    @Bean
    public ProducerFactory<String, MessageEventForDepartment> producerFactoryForDepartment() {
        Map<String, Object> configProps = new HashMap<>();
        configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress);
        configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
        return new DefaultKafkaProducerFactory<>(configProps);
    }

    @Bean
    public NewTopic topic1() {
        return TopicBuilder.name("MARCEL")
                .partitions(10)
                .compact()
                .build();
    }

    @Bean
    public KafkaTemplate<String, MessageEventForDepartment> kafkaTemplate() {
        return new KafkaTemplate<>(producerFactoryForDepartment());
    }

}

消费者容器配置

@Configuration
public class KafkaTopicConfig {

    @Value(value = "${kafka.bootstrapAddress}")
    private String bootstrapAddress;

/*    @Value(value = "${kafka.configId}")
    private String configId;*/

    @Bean
    public ConsumerFactory<String, MessageEventForDepartment> consumerFactoryForDepartments() {
        Map<String, Object> props = new HashMap<>();
        props.put(JsonDeserializer.TRUSTED_PACKAGES, "*");
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress);
        //props.put(ConsumerConfig.GROUP_ID_CONFIG, configId);
        return new DefaultKafkaConsumerFactory<>(props, new StringDeserializer(), new JsonDeserializer<>(MessageEventForDepartment.class));
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, MessageEventForDepartment>
    kafkaListenerContainerFactoryForDepartments() {

        ConcurrentKafkaListenerContainerFactory<String, MessageEventForDepartment> factory =
                new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactoryForDepartments());
        return factory;
    }

}

消费者监听代码

@Component
@Slf4j
public class DepartmentKafkaService {

    @Autowired
    private DepartmentRepository departmentRepository;

    @KafkaListener(topics = "MARCEL" , groupId = "ren",containerFactory = "kafkaListenerContainerFactoryForDepartments")
    public void listenGroupFoo(MessageEventForDepartment message) {
    ...
}

问题原因与解决方案

核心原因

  1. 消费线程数未配置:当前消费者容器工厂默认只有1个消费线程,Kafka规则是同一个消费组内每个分区只能被一个消费者(或线程)消费。如果消息集中在少数分区,就会出现只有部分实例能收到消息的情况。
  2. 消息分区分配不均:如果生产者发送消息时使用固定Key,所有消息会被路由到同一个分区,此时即使有多个实例,也只有一个实例能消费该分区的消息。

解决方案

1. 配置消费者容器并发数

在KafkaTopicConfig的容器工厂中添加并发数设置,建议值不超过分区数(你的分区数是10,可根据实例数调整):

@Bean
public ConcurrentKafkaListenerContainerFactory<String, MessageEventForDepartment>
kafkaListenerContainerFactoryForDepartments() {

    ConcurrentKafkaListenerContainerFactory<String, MessageEventForDepartment> factory =
            new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactoryForDepartments());
    // 设置并发线程数,比如实例数为3时设为3,确保每个实例至少分配到一个分区
    factory.setConcurrency(3); 
    return factory;
}

2. 优化生产者的分区路由策略

  • 若不需要按Key分区,发送消息时不指定Key,Kafka会自动轮询将消息分配到各个分区,保证多实例均匀接收消息;
  • 若必须使用Key,确保Key取值多样,让哈希值均匀映射到不同分区,避免消息集中在少数分区。

3. 验证分区分配状态

使用Kafka命令行工具查看消费组的分区分配情况,确认每个实例是否分配到了分区:

kafka-consumer-groups.sh --bootstrap-server ${kafka.bootstrapAddress} --describe --group ren

输出会显示每个分区对应的消费者实例,以此验证分配是否正常。


内容的提问来源于stack exchange,提问作者Renzo Redoran

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 22:31:15