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

SpringBoot Kafka消费者无法接收消息问题求助

SpringBoot Kafka消费者无法接收消息排查方案

问题描述

当前采用多微服务架构,包含生产者、消费者及Kafka组件。生产者接收HTTP请求后将消息转发至Kafka,Kafka控制台可正常接收消息,但SpringBoot消费者无法接收消息,请求协助排查解决。

相关配置信息

KafkaConsumerConfig

@Configuration
@EnableKafka
public class KafkaConfig {
    @Value("${kafka.bootstrap-servers}")
    private String bootstrapServers;

    @Bean
    public Map<String, Object> consumerConfigs() {
        Map<String, Object> props = new HashMap<>();
        // list of host:port pairs used for establishing the initial connections to the Kafka cluster
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        // allows a pool of processes to divide the work of consuming and processing records
        props.put(ConsumerConfig.GROUP_ID_CONFIG, "testTopicGroup");
        // automatically reset the offset to the earliest offset
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");

        return props;
    }

    @Bean
    public ConsumerFactory<String, String> consumerFactory() {
        return new DefaultKafkaConsumerFactory<>(consumerConfigs());
    }

    @Bean
    @Primary
    public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());

        return factory;
    }

    @Bean(name = "kafkaListenerContainerFactoryWith6Consumer")
    public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> kafkaListenerContainerFactoryWith6Consumer() {
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConcurrency(6); //3 partition -> 6 thread in parallel in a single consumer
        factory.setConsumerFactory(consumerFactory());
        return factory;
    }

    @Bean(name = "kafkaListenerContainerFactoryForBatchConsumer")
    public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> kafkaListenerContainerFactoryForBatchConsumer() {
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConcurrency(1);
        factory.setBatchListener(true);
        factory.setConsumerFactory(consumerFactory());
        return factory;
    }
}

Listener

@Slf4j
public class KafkaListener {
    @org.springframework.kafka.annotation.KafkaListener(topics = "${kafka.topic.messageTopic}")
    public void receive(String payload){
        log.info("Message is received from Kafka: " + payload);
    }

}

ConsumerApplicationProperties

kafka.bootstrap-servers=localhost:9092
server.port=8081
kafka.topic.messageTopic=test

排查步骤与解决方案

  • 修复监听器的Spring容器注册问题:当前KafkaListener类未添加@Component(或@Service、@Repository)注解,Spring无法扫描并初始化该监听器。修改监听器类,添加@Component注解:

    @Slf4j
    @Component
    public class KafkaListener {
        @org.springframework.kafka.annotation.KafkaListener(topics = "${kafka.topic.messageTopic}")
        public void receive(String payload){
            log.info("Message is received from Kafka: " + payload);
        }
    }
    
  • 验证消费者组偏移量状态:如果消费者组testTopicGroup之前已提交过偏移量且已追平最新消息位置,即使配置了auto.offset.reset=earliest也不会重新消费旧消息。可通过Kafka命令行工具查看偏移量:

    kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group testTopicGroup
    

    若偏移量已处于最新位置,可发送新消息测试,或手动重置偏移量:

    kafka-consumer-groups.sh --bootstrap-server localhost:9092 --reset-offsets --to-earliest --topic test --group testTopicGroup --execute
    
  • 检查主题分区与并发数匹配:确认test主题的分区数量,若使用6并发的容器工厂,分区数需大于等于并发数才能充分利用线程。查看主题信息:

    kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic test
    
  • 开启调试日志定位问题:将Spring Kafka及Kafka客户端日志级别设为DEBUG,查看消费者连接、订阅、拉取消息的详细流程,定位具体异常:

    logging.level.org.springframework.kafka=DEBUG
    logging.level.kafka=DEBUG
    
  • 确认序列化/反序列化一致性:确保生产者使用的序列化器与消费者的反序列化器完全匹配(当前配置为StringSerializer/StringDeserializer),若生产者使用其他序列化方式(如JSON),需同步修改消费者的反序列化配置。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 16:37:50