这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 再平衡(Rebalance)是消费组弹性的来源,也是生产事故的高发区:一次大规模 rebalance 可能让整个组停止消费数十秒。本章讲清楚触发条件、协议流程、分配策略与如何避免 rebalance 风暴。
15.1 什么是 Rebalance
消费组的成员或分区发生变化时,协调器把分区在消费者之间重新分配,这个过程叫 rebalance。
触发条件只有三类:
- 成员变化:新消费者加入、旧消费者退出/崩溃;
- 订阅变化:消费者订阅的 topic 集合改变(含正则匹配到新 topic);
- 分区数变化:topic 增加分区。
注意:Broker 宕机本身不触发消费组 rebalance(分区 Leader 切换对客户端透明),除非它同时是协调器且导致组会话超时。
15.2 协议流程
消费者A 消费者B 消费者C Coordinator(Broker)
│ │ │ │
│--- JoinGroup ------------------->│ 所有成员报名
│ │ │ │ 选出组 leader(这里是 A)
│<-- JoinGroup 响应(你是leader) ----│
│<------ 你是普通成员 --------------│
│ │
│ A 本地运行分配策略:
│ A: p0,p1 B: p2,p3 C: p4,p5
│ │
│--- SyncGroup(分配方案) --------->│
│<------ SyncGroup(A: p0,p1) ------│
│ <------ (B: p2,p3) -------│
│ <---- (C: p4,p5)-│
│ │
│--- Heartbeat / OffsetCommit --->│ 进入稳定态
再次强调:分配方案由组 leader 消费者计算(第 12 章),协调器只是组织者。generation(代数)随每次 rebalance 递增,用于隔离旧成员的过期请求。
15.3 Eager vs Cooperative:两种再平衡协议
Eager(全部回收再分配)
传统模式。触发 rebalance 时,所有消费者撤销全部分区,进入 JoinGroup,等全组达成一致后再重新分配。
代价:即使只是新增一个消费者,也会造成全组短暂停止消费。
Cooperative / Incremental(增量式)
只撤销“需要移动的分区”,其他分区继续消费:
Eager: A[p0,p1] B[p2,p3] -> 全部撤销 -> A[p0] B[p1,p2] C[p3]
Cooperative: A[p0,p1] B[p2,p3] -> A继续p0, B继续p2 -> 只迁移 p1,p3
2.4+ 版本支持。新系统建议直接使用 CooperativeStickyAssignor,把 rebalance 的停顿范围降到最小。
15.4 四种分配策略
配置项:partition.assignment.strategy。
RangeAssignor(默认之一)
按 topic 逐个计算:分区排序、消费者按名称排序,再连续切块。
topicX 6分区, topicY 3分区, 2个消费者
topicX: C1=[0,1,2] C2=[3,4,5]
topicY: C1=[0,1] C2=[2]
-> C1 共 5 个,C2 共 4 个
多 topic 时前面的消费者容易分到更多分区,可能不均。
RoundRobinAssignor
把所有 topic 的分区放在一起轮询分配,通常更均匀。但要求组内订阅一致才语义明确。
StickyAssignor
尽量均匀,且再平衡时尽量保留原分配(粘性),减少分区迁移。
CooperativeStickyAssignor
在 Sticky 基础上支持协作式增量再平衡,是当前推荐的默认选择:
partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor
15.5 静态成员:免打扰的重启
滚动发布时,每个实例重启都会触发 rebalance。配置静态成员后,实例以固定身份加入组:
group.instance.id=pod-1
session.timeout.ms=60000
行为:
- 重启时间小于
session.timeout.ms:协调器认为成员只是短暂离线,不触发 rebalance,分区保留;新实例以同一身份回来,拿着原分区继续消费; - 超时未归:才按普通成员失效处理。
这是 K8s/滚动升级场景的标配,能把发布对消费的影响降到最低。注意 group.instance.id 必须保证唯一且与实例绑定(如 Pod 序号),否则会互相 fence。
15.6 Rebalance 风暴:症状、原因与治理
典型症状
消费 lag 突增、日志出现 Attempt to heartbeat failed / session timeout / max poll interval、rebalancing 频繁、消费延迟周期性抖动。
常见原因
- 处理太慢:两次 poll 间隔超过
max.poll.interval.ms,被踢出组; - GC 停顿:长时间 STW,心跳/poll 都停;
- 网络抖动:心跳超时;
- 频繁发布:无静态成员 + 部署波次太密;
- 消费者订阅不一致:有的订阅 topic A,有的订阅 A+B,不断触发 rebalance。
治理清单
| 手段 | 配置/做法 |
|---|---|
| 降单次处理量 | max.poll.records 调小(如 100) |
| 放宽处理时限 | max.poll.interval.ms 调大(如 600000) |
| 心跳更宽容 | session.timeout.ms=60000,heartbeat.interval.ms=15000 |
| 滚动发布免扰 | group.instance.id + 保持 session 内回归 |
| 协作式再平衡 | CooperativeStickyAssignor |
| 统一订阅 | 同组所有消费者订阅完全相同 |
| 排查 GC | 缩短停顿、调整堆、看 GC 日志 |
15.7 观察 Rebalance
客户端日志是最直接的证据,关注:
(Re)join group ... as member
Successfully synced group
Attempt to heartbeat failed because group is rebalancing
Member ... sending LeaveGroup request
JMX 指标:
kafka.consumer:type=consumer-coordinator-metrics:sync-rate、assigned-partitions;- 协调器端:
kafka.server:type=group-coordinator-metrics。
把“1 分钟内 rebalance 次数”做成告警,是发现风暴的第一道防线。
15.8 动手实验
- 建一个 6 分区 topic,启动 2 个消费者(控制台或代码),观察分区分配;
- 再启动第 3 个消费者,观察 rebalance 与新分配;
- kill 一个消费者,观察分区转移与短暂停顿;
- 在业务处理里
Thread.sleep(310000),复现max.poll.interval.ms超时导致的踢组; - 切换为
CooperativeStickyAssignor,重复第 2 步,对比停顿范围。
本章小结
- Rebalance 由成员、订阅、分区数三类变化触发,Broker 宕机不直接触发;
- 分配方案由组 leader 消费者计算,generation 隔离旧请求;
- Cooperative 增量式再平衡只迁移必要分区,是新系统默认选择;
- 静态成员让滚动重启不打扰消费组,是容器环境的必备配置;
- rebalance 风暴的根源多为处理超时、GC 停顿与订阅不一致,治理要对症下药。
思考题
- 为什么
max.poll.records调小反而能缓解 rebalance? - Eager 策略下新增一个消费者,为什么没发生迁移的分区也暂停消费了?
- 两个消费者订阅的 topic 不一致会怎样?什么时候这反而是合理的?