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
相关产品推荐
相关产品推荐

