Kafka消费者接收约1万条消息后挂起问题咨询(含代码)
嘿,看你遇到的这个Kafka消费卡壳的问题,结合你描述的现象(消费1万条后挂死、重启恢复再循环,日志疑似rebalance后无法继续)和提供的代码,我整理了几个针对性的解决方案,你可以挨个试试:
1. 别让消费线程被数据库操作卡死,保住心跳
你现在的代码里,writeMessageToDatabase是同步执行的——如果数据库写入慢(比如网络卡、批量插入耗时久),消费线程会一直被占着,没法给ZooKeeper发心跳。旧版High-level Consumer全靠ZK维持会话,一旦心跳超时,ZK就会认为这个consumer挂了,触发rebalance;要是rebalance的时候消费线程还在阻塞,新分配的分区就可能没法正常启动消费。
试试这么做:
把数据库写入扔到异步线程池里,让消费线程只负责快速读消息、处理offset:
// 先初始化一个线程池,大小根据你的数据库承载能力调 ExecutorService dbExecutor = Executors.newFixedThreadPool(4); for (MessageAndMetadata<byte[], byte[]> message : stream0) { try { String messageReceived = new String(message.message(), "UTF-8"); logger.info("partition = " + message.partition() + ", offset=" + message.offset() + " => " + messageReceived); // 异步处理数据库写入,不阻塞消费线程 dbExecutor.submit(() -> writeMessageToDatabase(messageReceived)); } catch (UnsupportedEncodingException e) { e.printStackTrace(); } }
2. 手动控制offset提交,避免rebalance时的offset混乱
你现在开了自动提交(auto.commit.enable=true),间隔1秒,但如果rebalance刚好卡在自动提交前,很可能出现offset没提交、重复消费,甚至直接把消费流搞断的情况。手动提交能精准控制时机,确保只有消息处理完了才提交offset。
试试这么做:
- 先关掉自动提交:
props.put("auto.commit.enable", "false"); - 每处理一批消息就手动提交一次,比如每100条提交一次:
int batchSize = 100; int processedCount = 0; for (MessageAndMetadata<byte[], byte[]> message : stream0) { try { String messageReceived = new String(message.message(), "UTF-8"); logger.info("partition = " + message.partition() + ", offset=" + message.offset() + " => " + messageReceived); writeMessageToDatabase(messageReceived); processedCount++; // 每攒够一批就提交offset if (processedCount % batchSize == 0) { consumer.commitOffsets(); logger.info("已提交第" + processedCount + "条消息的offset"); } } catch (UnsupportedEncodingException e) { e.printStackTrace(); } }
3. 调整ZK会话超时参数,适配你的消费速度
旧版High-level Consumer的zookeeper.session.timeout.ms默认是60秒,要是你的单批消息处理时间接近甚至超过这个值,ZK会直接判定consumer死亡并触发rebalance。你现在只调了连接超时(200秒),但会话超时没改,很容易出现会话提前过期的情况。
试试这么做:
在props里加这两个参数:
props.put("zookeeper.session.timeout.ms", "200000"); // 和连接超时保持一致 props.put("zookeeper.sync.time.ms", "20000"); // 同步时间设为会话超时的1/10,符合ZK的最佳实践
4. 赶紧升级到新版Kafka Consumer API(最推荐)
你现在用的ConsumerConnector是已经被废弃的旧版High-level Consumer,这个API在rebalance机制、可靠性上有不少已知bug。Kafka官方早就推荐用新版的org.apache.kafka.clients.consumer.KafkaConsumer(从0.9版本开始引入),它靠Kafka集群自己的协调器,不依赖ZK,rebalance逻辑更稳定,调试也更方便。
给你个简单的新版代码示例:
Properties props = new Properties(); props.put("bootstrap.servers", "你的Kafka broker地址"); props.put("group.id", "Tornado"); props.put("auto.offset.reset", "earliest"); props.put("enable.auto.commit", "false"); props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props); consumer.subscribe(Collections.singletonList(TwitterConstant.Kafka.TWITTER_STREAMING_TOPIC)); try { while (true) { // 拉取消息,超时时间设100毫秒 ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100)); for (ConsumerRecord<String, String> record : records) { logger.info("partition = " + record.partition() + ", offset=" + record.offset() + " => " + record.value()); writeMessageToDatabase(record.value()); } // 批量提交offset,确保消息都处理完了再提交 consumer.commitSync(); } } finally { consumer.close(); }
针对KafkaSpout的补充建议
如果用KafkaSpout也出现类似的停服问题,本质原因差不多:要么是底层消费线程被堵,要么是rebalance异常。你可以试试:
- 调小
max.spout.pending参数,别让Spout同时处理太多消息导致阻塞 - 确保Spout用的是新版Kafka客户端,或者直接换基于新版Consumer API的Spout实现
- 调整Spout的ack超时时间,避免因为没及时ack导致消息堆积、停止消费
内容的提问来源于stack exchange,提问作者keduadoi

