Skip to content

Kafka 消费组重平衡:扩容时为什么反而会暂停消费 ​

给消费组加一台机器,本意是提高吞吐,结果所有分区一起停了 8 秒多。原因是组里另一个成员正在处理一批慢消息,而默认配置下的重平衡要等全体成员重新入组。换成协作式分配或 Kafka 4 的新协议后,同样的场景里,健康成员保留的分区只停了一百多毫秒。

本文在 Kafka 4.3.1(三节点 KRaft)和 kafka-clients 4.3.1 上对比三种重平衡方式,并复现重平衡导致的重复消费和静态成员身份的效果。消费组相关的服务端配置都是默认值,每个扩容场景运行两轮。

一、先说结论 ​

  • 客户端默认仍是经典协议加 eager 分配。4.3.1 客户端的 group.protocol 默认是 classic,partition.assignment.strategy 默认列表中第一个是 RangeAssignor。扩容时所有成员先撤销全部分区,再等全体重新入组。
  • eager 的代价取决于最慢的成员。批次处理得快时,实测全组只停了约 140ms;但只要有一个成员在处理慢批次,健康成员的分区也要陪着等,实测约 8.6 秒。
  • 协作式分配和新协议只移动需要移动的分区。代价是被移动的分区要停一段时间:协作式约 3.4 秒;新协议实测约 0.6 秒,但它取决于新成员的心跳落在 5 秒周期的哪里,最长可到 5 秒。新协议还去掉了全组同步,慢成员不会拖住别人的分区交接。
  • 重平衡本身不产生重复,没提交的进度才会。实测处理超时被移出组后,这一批 100 条被新成员再处理了一遍,提交抛出 CommitFailedException。
  • 滚动发布用静态成员身份(group.instance.id):实测成员重启期间,其他成员完全没有被打断。

二、三种重平衡方式 ​

经典协议 · eager(RangeAssignor)所有成员撤销全部分区等全部成员重新 JoinGroup重新分配全部分区经典协议 · 协作式(CooperativeStickyAssignor)其余分区照常消费只撤销要移走的分区全组同步一轮下一轮再分配给新成员新协议 · group.protocol = consumer(KIP-848)其余分区照常消费服务端计算目标分配随心跳逐个成员下发成员各自撤销、各自接收不等其他成员Kafka 4.3.1 客户端默认仍是 group.protocol = classic,默认分配策略列表中第一个是 eager 的 RangeAssignor
图 1 · eager 先撤销所有分区、等全体成员重新入组;协作式只撤销要移走的分区,但仍需要一轮全组同步;新协议由服务端逐个成员下发目标分配,不再有全组屏障
经典协议 · eager经典协议 · 协作式新协议(KIP-848)
启用方式默认partition.assignment.strategy=CooperativeStickyAssignorgroup.protocol=consumer
撤销范围所有成员的全部分区只撤销要移走的分区只撤销要移走的分区
全组同步需要,等所有成员重新 JoinGroup需要,并且可能要两轮不需要,服务端逐个成员下发
分配在哪里计算组内的 leader 消费者组内的 leader 消费者Broker 上的组协调器
心跳与会话超时客户端配置客户端配置服务端配置 group.consumer.*

新协议在 Kafka 4.0 正式可用(GA)。官方文档说明,它的增量设计「不再依赖全局同步屏障」。启用后,客户端的 partition.assignment.strategy、session.timeout.ms、heartbeat.interval.ms 都不能再设置(实测创建消费者时直接抛出 ConfigException: session.timeout.ms cannot be set when group.protocol=CONSUMER),分别由服务端的 group.consumer.assignors(默认 uniform,range)、group.consumer.session.timeout.ms、group.consumer.heartbeat.interval.ms(默认 5 秒)控制。目前不支持客户端自定义分配器,机架感知分配也还没有完整实现。

用 Spring Kafka 时也一样。在 Spring Boot 4.1.1 上实测,不做任何配置的 @KafkaListener 使用 group.protocol = classic 和 [RangeAssignor, CooperativeStickyAssignor],Spring 没有改动客户端的默认值;只加一个 spring.kafka.consumer.properties[group.protocol]=consumer 就能切换,监听器照常消费。如果配置里还留着 session.timeout.ms,应用在启动时失败,根因就是上面的 ConfigException。另外,Spring Boot 4.1.1 管理的 kafka-clients 是 4.2.1,和 broker 版本不必一致。

三、实测:给消费组扩容 ​

主题 6 个分区,生产者每个分区每秒发送约 10 条。两个消费者 C1、C2 稳定运行 12 秒后,加入第三个消费者 C3,记录每个成员的撤销与分配回调,以及每个分区相邻两条消息的最大处理间隔(下文称「中断」)。每条消息模拟 10ms 的处理耗时,消费能力有富余。

3.1 所有成员都健康时 ​

方式被撤销的分区最长中断中断超过 300ms 的分区
eager全部 6 个131~148 ms0
协作式2 个3.4 秒2
新协议2 个0.6 秒2

(每种方式运行两轮,取值范围来自两轮的结果;没有变化的分区在三种方式下最长中断都在 170ms 以内。)

这组结果和常见印象相反:eager 撤销了所有分区,但所有成员都在及时 poll(),很快就重新入组,全组只停了一百多毫秒。协作式和新协议只动了 2 个分区,但这 2 个分区停得更久。协作式要走两轮 JoinGroup:第一轮让原来的所有者交出分区,第二轮才把它分给新成员,每一轮都要等成员在心跳里发现需要重新入组(经典协议默认心跳间隔 3 秒)。新协议由服务端逐个成员下发变更,新成员要到下一次心跳才知道自己拿到了分区;默认心跳间隔是 5 秒,所以这段停顿在几百毫秒到 5 秒之间,调试时的多次运行里最短约 0.5 秒、最长约 5.2 秒。

所以,如果每批消息都处理得很快,eager 的停顿其实很短,这时换协议的收益有限。

3.2 有一个成员正在处理慢批次时 ​

真实系统里,扩容往往发生在消费跟不上的时候,这时很可能有成员正卡在一批慢消息上。重复上面的实验,让 C2 在 C3 加入前 200ms 进入一次 8 秒的处理:

eagerC1 保留的分区先被全部撤销8.6 秒换了所有者的分区等 C2 重新入组9.1 秒协作式C1 保留的分区150 msC1 移给 C3 的分区C2 结束后才交出3.1 秒新协议C1 保留的分区160 msC1 移给 C3 的分区C2 慢处理期间已交出1.3 秒C2 的慢处理:8 秒
图 2 · C3 加入时 C2 正在处理一批耗时 8 秒的消息;eager 让健康成员 C1 的分区也陪着等了约 8.6 秒,协作式和新协议下 C1 保留的分区基本不受影响(Kafka 4.3.1 三节点实测,两轮中的较大值)

eager 的事件时间线(节选):

text
 11800 ms  C2 开始一次 8 秒的慢处理
 12003 ms  C3 启动
 12180 ms  C1 撤销 [0, 1, 2]        ← 健康的 C1 立刻交出全部分区
 19801 ms  C2 慢处理结束
 19815 ms  C2 撤销 [3, 4, 5]
 19819 ms  C1 分配 [0, 1]           ← 等 C2 重新入组后才完成分配

C1 交出分区后,要等 C2 处理完这批消息、重新调用 poll() 并入组,整个组才能完成分配。C1 本来没有任何问题,它的分区却停了 7.7~8.6 秒。

方式C1 保留的分区最长中断C1 交出要移走的分区组里首次完成重新分配
eager7.7~8.6 秒(先被全部撤销)立刻交出全部分区C2 慢处理结束之后
协作式125~145 msC2 慢处理结束之后C2 慢处理结束之后
新协议129~155 msC2 还在慢处理时C2 还在慢处理时

经典协议下,C2 处理这批慢消息的 8 秒里,组里没有完成任何重新分配:eager 让 C1 先交出全部分区再干等;协作式让 C1 保留的分区照常消费,但要移走的分区也要等 C2 回来才交出。新协议下,服务端直接通知 C1 交出要移走的那个分区,新成员在自己的下一次心跳时接手,整个过程不需要 C2 参与。

C2 自己名下的分区在三种方式下都中断了约 8 秒。这是慢处理本身造成的,换协议解决不了,要靠控制单批处理时间。

四、重平衡导致的重复消费 ​

重复的根源是:业务处理完成和提交 offset 不是一个原子操作。处理完一批但还没提交时,分区被分配给了别人,新成员就会从上次提交的位置重新消费。

复现方法:两个分区、两个消费者,关闭自动提交,每批处理完调用 commitSync()。max.poll.interval.ms 设为 5 秒,让 C1 的第 3 批每条处理 80ms,一批 100 条共 8 秒:

text
 9138 ms  C2 撤销 [p1]
 9143 ms  C2 分配 [p0, p1]           ← C1 超过 max.poll.interval.ms,被移出组,p0 交给 C2
 9419 ms  C1 提交失败:CommitFailedException
 9423 ms  C1 丢失 [p0](已被移出组)
== classic:共 2000 条消息,其中 100 条被处理了两次

C1 处理完的那 100 条没能提交,C2 从上一次提交的位置开始,又处理了一遍。注意 C1 收到的是 onPartitionsLost 而不是 onPartitionsRevoked:成员已经被移出组,这时分区可能已经属于别人,不能再提交。

用新协议运行同样的程序,这一次没有出现重复:C1 被移出组后 0.1 秒就重新加入,又分到了 p0,这一批的处理和提交都由它自己完成。这是这次运行的时序决定的,不能当作新协议的保证。

应对方式:

  1. 消费逻辑要幂等:以消息的业务键或 offset 做去重,这是唯一可靠的办法。见 Kafka 不丢、不重与 Exactly Once。
  2. 在 onPartitionsRevoked 里提交已完成的进度,缩小重复的范围。它在分区被正常撤销前调用,可以同步提交。
  3. 控制单批处理时间:max.poll.records 乘以单条最坏耗时,要明显小于 max.poll.interval.ms。不要只是把超时调大,那样卡住的成员要更久才会被发现。
  4. 慢业务交给工作线程时,要自己管理 offset:只提交已经完成的连续位置,并在队列积压时用 pause() 暂停拉取。

五、减少不必要的重平衡 ​

5.1 静态成员身份 ​

滚动发布时,每个实例重启一次,默认会触发两次重平衡:下线一次、上线一次。设置 group.instance.id 后,成员在会话超时之内回来,组协调器会把原来的分区还给它,不触发重平衡。

实测:C1 关闭,3 秒后用同一个 group.instance.id 重启。

C2 经历的撤销与分配C1 重启后
经典协议,动态成员被撤销 2 次、分配 2 次约 3 秒后重新分到分区
经典协议,静态成员没有任何变化约 50ms 后拿回原来的分区
新协议,动态成员被撤销 1 次、分配 1 次(先接手 C1 的分区,再交回)约 5 秒后才拿回分区
新协议,静态成员没有任何变化约 50ms 后拿回原来的分区

代价是:静态成员离开期间,它的分区没有人消费,直到它回来或者会话超时。所以 session.timeout.ms(新协议下是服务端的 group.consumer.session.timeout.ms)要大于一次正常重启的时间,又不能大到真正宕机时要很久才转移分区。另外,group.instance.id 在组内必须唯一,通常用 Pod 名或主机名,重复会导致成员互相顶替。

5.2 其他做法 ​

  • 保持 poll() 的节奏:慢处理是被移出组的首要原因。
  • 优雅停机:先停止拉取、处理完手头的批次、提交,再关闭消费者。
  • 健康检查不要依赖消费进度:否则积压时容器会被重启,重启又触发重平衡,形成恶性循环。

六、消费者是不是越多越快 ​

同一个消费组里,一个分区同一时刻只分配给一个成员,所以有效并行度的上限是分区数。消费者比分区多时,多出来的实例空闲,只在别人故障时接手。

增加分区可以提高上限,但会改变 key 到分区的映射,破坏基于 key 的顺序,也会增加元数据和恢复的开销。需要更高并行度又不想加分区时,可以在消费者内部按 key 分发给多个工作线程,同一个 key 始终进同一个线程,保持单 key 有序。

七、监控什么 ​

指标看什么
每个分区的 lag总 lag 会掩盖单个分区卡住的情况
重平衡次数与耗时发布之外的重平衡通常意味着处理超时或实例抖动
单批处理耗时与 poll 间隔接近 max.poll.interval.ms 时要预警
提交失败次数CommitFailedException 意味着刚刚发生过重复消费
分区分配是否均匀分配倾斜会让个别实例成为瓶颈
发布期间的吞吐曲线周期性归零说明每次重启都在触发全组重平衡

八、常见误区 ​

  • 「协作式分配比 eager 停顿更短」:只在有慢成员时成立。所有成员都健康时,实测 eager 全组只停了约 140ms,而协作式被移动的分区要停 3 秒左右。
  • 「升级到 Kafka 4 就自动用上了新协议」:服务端默认启用,但客户端要显式设置 group.protocol=consumer。
  • 「用了新协议还能配置 session.timeout.ms」:这些配置改由服务端控制,客户端设置了会在创建消费者时抛出 ConfigException。迁移时要先从配置里删掉。
  • 「重平衡会导致重复消费」:重复来自没提交的进度,幂等消费才是根本。
  • 「处理慢就把 max.poll.interval.ms 调大」:会让真正卡死的成员更晚被发现,应该先缩小单批处理量。

小结 ​

重平衡的代价由两部分组成:撤销了多少分区,以及要等多久才能重新分配。eager 撤销得多,但只要所有成员都健康,恢复得很快;它的真正风险是一个慢成员能拖住全组。协作式分配和新协议只移动必要的分区,新协议还去掉了全组同步,更适合消费逻辑耗时不稳定的系统。在此之上,用静态成员身份避免发布抖动,用幂等消费兜住重复,才是完整的方案。


配套实验

参考资料

文章以 CC BY-NC-SA 4.0 授权 · 代码片段以 MIT 授权