KafkaNotes

第 15 章:Rebalance 内幕:分区分配与再平衡

zjc 于 2026-01-15 发布

这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 再平衡(Rebalance)是消费组弹性的来源,也是生产事故的高发区:一次大规模 rebalance 可能让整个组停止消费数十秒。本章讲清楚触发条件、协议流程、分配策略与如何避免 rebalance 风暴。

15.1 什么是 Rebalance

消费组的成员或分区发生变化时,协调器把分区在消费者之间重新分配,这个过程叫 rebalance。

触发条件只有三类:

  1. 成员变化:新消费者加入、旧消费者退出/崩溃;
  2. 订阅变化:消费者订阅的 topic 集合改变(含正则匹配到新 topic);
  3. 分区数变化: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

行为:

这是 K8s/滚动升级场景的标配,能把发布对消费的影响降到最低。注意 group.instance.id 必须保证唯一且与实例绑定(如 Pod 序号),否则会互相 fence。

15.6 Rebalance 风暴:症状、原因与治理

典型症状

消费 lag 突增、日志出现 Attempt to heartbeat failed / session timeout / max poll intervalrebalancing 频繁、消费延迟周期性抖动。

常见原因

  1. 处理太慢:两次 poll 间隔超过 max.poll.interval.ms,被踢出组;
  2. GC 停顿:长时间 STW,心跳/poll 都停;
  3. 网络抖动:心跳超时;
  4. 频繁发布:无静态成员 + 部署波次太密;
  5. 消费者订阅不一致:有的订阅 topic A,有的订阅 A+B,不断触发 rebalance。

治理清单

手段 配置/做法
降单次处理量 max.poll.records 调小(如 100)
放宽处理时限 max.poll.interval.ms 调大(如 600000)
心跳更宽容 session.timeout.ms=60000heartbeat.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 指标:

把“1 分钟内 rebalance 次数”做成告警,是发现风暴的第一道防线。

15.8 动手实验

  1. 建一个 6 分区 topic,启动 2 个消费者(控制台或代码),观察分区分配;
  2. 再启动第 3 个消费者,观察 rebalance 与新分配;
  3. kill 一个消费者,观察分区转移与短暂停顿;
  4. 在业务处理里 Thread.sleep(310000),复现 max.poll.interval.ms 超时导致的踢组;
  5. 切换为 CooperativeStickyAssignor,重复第 2 步,对比停顿范围。

本章小结

思考题

  1. 为什么 max.poll.records 调小反而能缓解 rebalance?
  2. Eager 策略下新增一个消费者,为什么没发生迁移的分区也暂停消费了?
  3. 两个消费者订阅的 topic 不一致会怎样?什么时候这反而是合理的?