低流量微服务KafkaConsumer#poll耗时过长致Offset提交失败问题
Kafka Consumer低流量场景下poll耗时过长导致被踢出消费组问题
问题现象
遇到Kafka Consumer报错:
Offset commit cannot be completed since the consumer is not part of an active group for auto partition assignment; it is likely that the consumer was kicked out of the group.
排查后新增interval between poll指标,发现两次KafkaConsumer#poll调用间隔平均30-40秒,与设计的100ms轮询频率不符。核心消费代码如下:
while (true) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100)); // Process records... Thread.sleep(100); }
进一步定位发现,耗时几乎全部来自KafkaConsumer#poll方法本身,且仅低流量微服务存在该问题(每5分钟接收1条小于1KiB的消息),处理高流量、1MiB左右消息的服务无此现象。
当前消费者配置:
heartbeat.interval.ms=10000 session.timeout.ms=30000 max.poll.records=100 max.poll.interval.ms=300000
消费者日志
agent.MyAgent | 11/10/2023 09:45:27.267 ERROR [MessageService on MyAgent-b834fc68614a] ConsumerCoordinator - [Consumer instanceId=MyAgent_b834fc68614a, clientId=consumer-MyAgent-MyAgent_b834fc68614a, groupId=MyAgent] Offset commit failed on partition DiscoveryRequest.Target-2 at offset 265: The coordinator is not aware of this member. agent.MyAgent | 11/10/2023 09:45:27.284 ERROR [MessageService on MyAgent-b834fc68614a] ConsumerCoordinator - [Consumer instanceId=MyAgent_b834fc68614a, clientId=consumer-MyAgent-MyAgent_b834fc68614a, groupId=MyAgent] Asynchronous auto-commit of offsets failed: Offset commit cannot be completed since the consumer is not part of an active group for auto partition assignment; it is likely that the consumer was kicked out of the group.. Will continue to join group. agent.MyAgent | 11/10/2023 09:45:27.860 WARN [MessageService on MyAgent-b834fc68614a] ConsumerCoordinator - [Consumer instanceId=MyAgent_b834fc68614a, clientId=consumer-MyAgent-MyAgent_b834fc68614a, groupId=MyAgent] Asynchronous auto-commit of offsets {...} failed: Offset commit cannot be completed since the consumer is not part of an active group for auto partition assignment; it is likely that the consumer was kicked out of the group.
Kafka日志
kafka | [2023-10-11 09:44:34,120] INFO [GroupCoordinator 1]: Member MyAgent_b834fc68614a-9891de54-09fa-42b0-9df3-7237953297c2 in group MyAgent has failed, removing it from the group (kafka.coordinator.group.GroupCoordinator) kafka | [2023-10-11 09:44:34,120] INFO [GroupCoordinator 1]: Preparing to rebalance group MyAgent in state PreparingRebalance with old generation 203 (__consumer_offsets-28) (reason: removing member MyAgent_b834fc68614a-9891de54-09fa-42b0-9df3-7237953297c2 on heartbeat expiration) (kafka.coordinator.group.GroupCoordinator) kafka | [2023-10-11 09:44:34,120] INFO [GroupCoordinator 1]: Group MyAgent with generation 204 is now empty (__consumer_offsets-28) (kafka.coordinator.group.GroupCoordinator) kafka | [2023-10-11 09:45:27,571] INFO [GroupCoordinator 1]: Static member with groupInstanceId=MyAgent_b834fc68614a and unknown member id joins group MyAgent in Empty state. Created a new member id MyAgent_b834fc68614a-80adc162-5dae-4769-8eea-cd2fd0ac5ad2 for this member and add to the group. (kafka.coordinator.group.GroupCoordinator) kafka | [2023-10-11 09:45:27,571] INFO [GroupCoordinator 1]: Preparing to rebalance group MyAgent in state PreparingRebalance with old generation 204 (__consumer_offsets-28) (reason: Adding new member MyAgent_b834fc68614a-80adc162-5dae-4769-8eea-cd2fd0ac5ad2 with group instance id Some(MyAgent_b834fc68614a); client reason: encountered UNKNOWN_MEMBER_ID from OFFSET_COMMIT response) (kafka.coordinator.group.GroupCoordinator)
核心疑问
为何低流量、小消息场景下KafkaConsumer#poll方法会如此耗时?
问题分析与解决方案
根本原因
- poll方法的后台任务阻塞:
KafkaConsumer#poll不仅负责拉取消息,还需处理心跳发送、消费组协调器交互、重平衡等后台逻辑。低流量场景下,没有消息可拉取时,poll会等待超时,但如果此时出现心跳超时被踢出组的情况,poll内部会触发重新加入消费组的流程,该流程包含与协调器通信、分区重新分配等操作,会导致poll耗时剧增。 - 心跳配置不合理:当前
heartbeat.interval.ms=10000、session.timeout.ms=30000,心跳间隔接近会话超时的1/3(Kafka推荐心跳间隔为会话超时的1/3以内),低流量下poll执行间隔被拉长,导致心跳发送不及时,触发会话超时被踢出组。 - 额外sleep加剧问题:代码中
Thread.sleep(100)进一步拉长了两次poll的间隔,减少了心跳发送和协调器交互的机会,加速了被踢出组的过程,而重加入组的开销又让poll耗时变得更长。
高流量服务无此问题的原因:高流量下poll能频繁获取消息快速返回,心跳和协调器交互可正常在poll间隙完成,不会出现长时间阻塞。
解决方案
- 调整心跳与会话超时配置:将心跳间隔设置为会话超时的1/3以内,比如:
heartbeat.interval.ms=3000 session.timeout.ms=10000
确保Consumer能及时发送心跳,避免被踢出组。
2. 移除不必要的sleep:删除Thread.sleep(100),让poll能更频繁执行,保证心跳和协调器交互的及时性。修改后代码:
while (true) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100)); // Process records... }
- 改用手动提交偏移量:自动提交在Consumer被踢出组时容易出现提交失败,手动提交可精准控制提交时机,减少异常交互。
- 确保Consumer单线程使用:
KafkaConsumer不是线程安全的,多线程操作会导致内部逻辑阻塞,需保证Consumer实例仅在单线程中调用。
内容的提问来源于stack exchange,提问作者John
相关产品推荐
相关产品推荐

