Rebalance 机制
Rebalance(再均衡)是消费组重新分配分区的过程。它是消费组弹性的来源,也是无数线上事故的元凶——Rebalance 期间消费全停、完成后大概率重复消费。本章讲清它何时发生、怎么发生、如何让它少发生。
1. 什么时候触发
三类触发条件:
- 组成员变化(最常见)
- 新消费者加入(扩容、滚动发布)
- 消费者主动退出(close/优雅停机)
- 消费者被判死:心跳超时(session.timeout.ms)或 poll 间隔超时(max.poll.interval.ms)
- 订阅的 topic 变化:正则订阅匹配到新 topic
- 分区数变化:topic 扩分区
其中"被判死"往往不是真死,而是 GC 停顿过长、处理太慢——这类假死引发的 Rebalance 是重点治理对象。
2. Eager 协议:全停全分
经典(Eager)协议的流程简单粗暴:
成员变化
|
v
所有消费者 revoke 自己的全部分区 <- 全组停止消费! (stop the world)
|
v
重新 JoinGroup -> 组长计算新方案 -> SyncGroup 下发
|
v
所有消费者按新方案重新领分区, 从 committed offset 恢复消费问题显而易见:哪怕只有 1 个消费者加入,所有人都要放下手里的全部分区,重新走一遍入组流程。大消费组(几十个实例、上千分区)一次 Rebalance 可能停顿几十秒。
3. Cooperative 协议:增量式再均衡
Cooperative(协作式)协议(2.4+ 引入)把"一次全停"改成"多轮微调",只移动需要移动的分区:
场景: C1 有 p0,p1,p2 C2 有 p3,p4,p5 新成员 C3 加入
第一轮 Rebalance:
计算目标方案: C1[p0,p1] C2[p3,p4] C3[p2,p5]
C1 只 revoke p2, C2 只 revoke p5 <- 其余分区不停!
C1 继续消费 p0,p1; C2 继续消费 p3,p4
第二轮 Rebalance:
p2, p5 分配给 C3
完成. 全程只有 2 个分区短暂停顿启用方式是选用 CooperativeStickyAssignor(见下节)。新项目建议直接用它。
4. 四种分配策略
partition.assignment.strategy 决定组长怎么算分配方案:
| 分配器 | 算法 | 特点 |
|---|---|---|
| RangeAssignor(默认之一) | 每个 topic 独立按范围切分 | 多 topic 时前面的消费者总是多拿,负载倾斜 |
| RoundRobinAssignor | 所有 topic 分区排一起轮流发牌 | 均匀,但 Rebalance 后分配大变 |
| StickyAssignor | 均匀 + 尽量保留原有分配 | 减少分区移动,仍是 Eager 协议 |
| CooperativeStickyAssignor | Sticky + Cooperative 协议 | 推荐:增量再均衡,移动最少 |
Range 倾斜示意(2 个 topic 各 3 分区,2 个消费者):
RangeAssignor: RoundRobin:
C1: t1-p0 t1-p1, t2-p0 t2-p1 C1: t1-p0 t1-p2 t2-p1
C2: t1-p2, t2-p2 C2: t1-p1 t2-p0 t2-p2
C1 拿 4 个, C2 拿 2 个 (倾斜) 3 : 3 (均匀)配置:
partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor从 Eager 分配器(如默认 Range)切到 CooperativeSticky,不能一把全量替换。安全路径:第一次滚动发布把策略列表配成两个(Cooperative 在前、原策略在后),第二次滚动发布再去掉旧策略。跳过这个步骤混跑两种协议,会抛 IllegalStateException。
5. 静态成员:group.instance.id
默认情况下消费者重启会拿到新的成员 ID,"退出 + 加入"= 两次 Rebalance。静态成员给消费者一个固定身份:
group.instance.id=inventory-pod-0 # 每个实例唯一且重启不变
session.timeout.ms=120000 # 配合调大, 覆盖重启耗时效果:
- 消费者短暂重启(低于 session.timeout.ms)后以原身份回归,直接拿回原来的分区,不触发 Rebalance。
- 特别适合 Kubernetes StatefulSet:pod 名天然唯一稳定,滚动重启不再引发 Rebalance 风暴。
代价:实例真挂了要等满 session.timeout 才会把分区转移给别人,故障转移变慢——超时值要在"重启免 Rebalance"和"故障接管速度"之间权衡。
6. Rebalance 期间发生了什么坏事
- 消费停顿:Eager 协议下全组停止消费,积压瞬间上涨。
- 重复消费:revoke 时若有已处理未提交的消息,新接手的消费者会重放。
- 连锁反应:Rebalance 使 poll 变慢 → 更多消费者超时 → 再触发 Rebalance → 风暴。
在 ConsumerRebalanceListener.onPartitionsRevoked 里提交 offset 能把重复消费窗口降到最小:
consumer.subscribe(List.of("orders"), new ConsumerRebalanceListener() {
public void onPartitionsRevoked(Collection<TopicPartition> parts) {
consumer.commitSync(currentOffsets); // 交出分区前把进度存好
}
public void onPartitionsAssigned(Collection<TopicPartition> parts) {
log.info("assigned: {}", parts);
}
});7. 减少 Rebalance 的检查清单
| 措施 | 针对的问题 |
|---|---|
| 用 CooperativeStickyAssignor | 全停变增量 |
| 配 group.instance.id(静态成员) | 滚动重启不触发 Rebalance |
| max.poll.records 调小 / 处理提速 | 避免 poll 超时假死 |
| session.timeout.ms=45s、heartbeat=15s | 容忍短暂 GC/网络抖动 |
| 优雅停机(close 而非 kill -9) | 快速离组,缩短停顿 |
| JVM 调优减少长 GC | 避免心跳超时假死 |
| 监控 Rebalance 次数指标 | 及早发现风暴苗头 |
三个信号:消费者日志频繁出现 "(Re-)joining group"、LAG 呈锯齿状(停-追-停-追)、JMX 指标 kafka.consumer 下 rebalance-rate-per-hour 异常。看到这些先查 max.poll.interval 超时(处理慢),再查 GC 与网络。
小结
- Rebalance 由成员变化、订阅变化、分区数变化触发;假死(GC、处理慢)是最常见诱因
- Eager 协议全停全分,Cooperative 协议只移动必要分区——新项目直接用 CooperativeStickyAssignor
- Range 会倾斜,RoundRobin 均匀但漂移大,Sticky 系列兼顾均匀与稳定
- 静态成员 group.instance.id 让重启不触发 Rebalance,代价是故障接管变慢
- 在 onPartitionsRevoked 里提交 offset,把重复消费窗口压到最小
- 用默认分配器启动 3 个消费者,杀掉一个观察日志里的 Rebalance 全过程;换成 CooperativeStickyAssignor 重做,对比其余消费者是否继续消费。
- 给消费者配上 group.instance.id 并把 session.timeout 调到 2 分钟,快速重启一个实例,验证没有触发 Rebalance。
- 订阅两个各 3 分区的 topic、启动 2 个消费者,分别用 Range 与 RoundRobin 观察分配结果的倾斜差异。