Kafka消费者运行一段时间后停止接收消息问题求助
核心问题分析
从你的代码和配置来看,有几个关键问题会导致消费者逐渐停止接收消息:
MAX_POLL_INTERVAL_MS配置错误
你设置了MAX_POLL_INTERVAL_MS_CONFIG=60000(1分钟),但注释标注的是5分钟。如果单次拉取的150条消息处理时间超过1分钟,消费者会因为超过最大轮询间隔被Kafka集群判定为失效,触发消费组重平衡。频繁的重平衡会导致分区分配异常,最终出现消费者无法获取到任何分区的情况。消息确认时机不合理
你在消息处理前就调用了acknowledgment.acknowledge(),这种做法会带来两个问题:- 如果消息处理失败且DLT发送失败,这条消息会彻底丢失,因为offset已经提交,消费者不会重新消费。
- 如果处理过程中消费者崩溃,已确认但未处理完成的消息会丢失,重启后会从下一个offset开始消费。
潜在的重平衡Listener问题
你自定义了RebalanceListener但未提供实现代码,如果Listener在重平衡回调(如onPartitionsRevoked)中执行耗时操作或未正确处理offset提交,会导致重平衡超时,进而引发分区分配失败。
具体修复方案
1. 修正MAX_POLL_INTERVAL_MS配置
将配置值改为300000(5分钟),匹配你的注释说明:
props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 300000); // 正确设置为5分钟
同时,根据实际消息处理耗时,调整MAX_POLL_RECORDS_CONFIG的值,确保单次拉取的消息能在最大轮询间隔内处理完成。
2. 调整消息确认时机
将acknowledge()调用移到消息处理成功之后,处理失败时确保只有在DLT发送成功后才确认offset(或根据业务需求选择是否确认):
@Override public void onMessage(ConsumerRecord<String, Request> consumerRecord, Acknowledgment acknowledgment) { try { // process the message // 处理成功后再确认offset acknowledgment.acknowledge(); } catch (Exception e) { // publish it to DLT topic // 确保DLT发送成功后再确认,避免消息丢失 if (sendToDltSuccessfully(consumerRecord)) { acknowledgment.acknowledge(); } // 如果DLT发送失败,不要确认,让消费者重新消费该消息 } }
3. 检查自定义RebalanceListener
确保RebalanceListener的回调方法中没有耗时操作,在onPartitionsRevoked时正确提交当前已处理的offset(如果需要),避免阻塞重平衡流程。例如:
public class RebalanceListener implements ConsumerRebalanceListener { @Override public void onPartitionsRevoked(Collection<TopicPartition> partitions) { // 仅执行必要的offset提交操作,避免耗时逻辑 // 如果使用手动提交,这里可以提交已确认的offset } @Override public void onPartitionsAssigned(Collection<TopicPartition> partitions) { // 初始化分区相关逻辑,避免阻塞 } }
4. 调整消费者并发数
根据Topic的分区数量,设置ConcurrentKafkaListenerContainerFactory的concurrency属性,提升消费能力,避免消息堆积:
factory.setConcurrency(3); // 建议值不超过Topic的分区数
5. 考虑升级spring-kafka版本
spring-kafka 2.8.11虽然是稳定版,但后续的2.8.x小版本可能修复了一些重平衡和消费者稳定性相关的bug,升级到最新的2.8.x版本可以降低此类问题的发生概率。
额外排查点
- 检查Kafka集群的状态,确保broker没有出现频繁重启或网络波动,这会导致消费者连接不稳定。
- 监控消费者的JVM状态,查看是否存在内存泄漏或线程阻塞,导致消费者无法正常轮询消息。
内容的提问来源于stack exchange,提问作者Vega

