可能是云消息队列 Kafka 版客户端版本过低或者Consumer没有独立线程维持心跳。
问题现象
使用云消息队列 Kafka 版时,消费客户端频繁出现Rebalance。
可能原因
可能导致故障的原因包括:
v0.10.2之前版本的客户端:Consumer没有独立线程维持心跳,而是把心跳维持与poll接口耦合在一起。其结果就是,如果用户消费出现卡顿,就会导致Consumer心跳超时,引发Rebalance。
v0.10.2及之后版本的客户端:如果消费时间过慢,超过一定时间(
max.poll.interval.ms设置的值,默认5分钟)未进行poll拉取消息,则会导致客户端主动离开队列,而引发Rebalance。NOT_COORDINATOR / join group failed:消费者组协调器(Group Coordinator)发生变更,或客户端心跳超时后重新寻找 Coordinator 时,会触发 Rebalance。客户端日志中可见
NOT_COORDINATOR或join group failed关键词。MemberIdRequiredException / consumer poll timeout:v0.10.2 之前版本的客户端,心跳与 poll 接口耦合,消费卡顿导致心跳超时时可能出现此报错;v0.10.2 及之后版本中,若消费速度过慢,超过
max.poll.interval.ms配置的时间未进行 poll,客户端会主动离开消费组,客户端日志中可见consumer poll timeout相关关键词。CommitFailedException:消费者处理消息耗时过长,或未能及时发送心跳,导致被消费者组踢出(触发 Rebalance)后仍尝试提交 Offset,引发该异常。
解决方案
首先您需要了解以下几点信息:
session.timeout.ms:心跳超时时间(可以由客户端自行设置)。max.poll.records:每次poll返回的最大消息数量。v0.10.2之前版本的客户端:心跳是通过poll接口来实现的,没有内置的独立线程。
v0.10.2及之后版本的客户端:为了防止客户端长时间不进行消费,Kafka客户端在v0.10.2及之后的版本中引入了
max.poll.interval.ms配置参数。heartbeat.interval.ms:心跳发送间隔(由客户端自行设置)。
参考以下说明调整参数值:
session.timeout.ms:v0.10.2之前的版本可适当提高该参数值,需要大于消费一批数据的时间,但不要超过30s,建议设置为25s;而v0.10.2及其之后的版本,保持默认值10s即可。
max.poll.records:降低该参数值,建议远远小于
<单个线程每秒消费的条数> * <消费线程的个数> *<max.poll.interval.ms>的积。max.poll.interval.ms:该值要大于
<max.poll.records> / (<单个线程每秒消费的条数> * <消费线程的个数>)的值。heartbeat.interval.ms:建议设置为不超过 10000 ms(10 s),且需小于
session.timeout.ms的三分之一,以确保心跳在会话超时前能发送至服务端。
如果您使用 Spring Kafka 框架,可以在项目的
application.yml文件中通过以下属性配置上述参数:spring: kafka: consumer: properties: session.timeout.ms: 25000 max.poll.interval.ms: 300000 max.poll.records: 50 heartbeat.interval.ms: 8000尽量提高客户端的消费速度,消费逻辑另起线程进行处理。
减少Group订阅Topic的数量,一个Group订阅的Topic最好不要超过5个,建议一个Group只订阅一个Topic。
将客户端升级至0.10.2以上版本。
常见问题
如何临时减少 Rebalance 以加快堆积消息消费?
在消息出现堆积的紧急场景下,可参考以下措施快速恢复消费:
确保客户端版本 ≥ 0.10.2,以使用独立心跳线程,减少因心跳与 poll 耦合导致的超时 Rebalance。
适当调大
max.poll.interval.ms,使其大于单次完整消费逻辑的最大耗时,避免客户端主动离组。降低
max.poll.records值,减少单次拉取的消息量,确保每次 poll 拉取的消息能在超时时间内处理完毕。减少 Group 订阅的 Topic 数量(建议不超过 5 个,最好一个 Group 只订阅一个 Topic)。
避免消费线程阻塞,将消费逻辑异步处理(另起独立线程执行业务逻辑,主线程只负责 poll 和提交 Offset)。