上一篇【第24篇】消息传递保证语义深度解析——at-most-once/at-least-once/exactly-once 下一篇【第26篇】ConsumerNetworkClient源码解析——消费者的"网络大脑"
摘要
Kafka消费者组的Rebalance(再均衡)是整个消费者体系中最精妙也最容易出问题的一环。当一个消费者加入或离开Group、订阅的Topic分区数发生变化、或者消费者心跳超时时,Kafka会自动触发Rebalance——把分区重新分配给存活的消费者,仿佛在进行一场"重新洗牌"。
本文从Rebalance的触发条件切入,深入剖析JoinGroup和SyncGroup两道协议请求的完整交互流程,用ASCII图告诉你每个消费者在Rebalance过程中在干什么。然后逐一拆解三种分区分配策略:RangeAssignor的"按Topic划分"、RoundRobinAssignor的"轮询公平分配"、StickyAssignor的"粘性保底"。最后分析Rebalance的性能影响——“停止消费"期间到底发生了什么,以及如何缩短这个"空窗期”。读完本文,你将能从源码层面理解Rebalance的每一个细节。
一、Rebalance是什么——"重新洗牌"的必要性
在Kafka的消费者组模型中,一个Partition在同一时刻只能被Group中的一个Consumer消费。当消费者组中消费者数量变化、或订阅的Topic分区数变化时,就必须重新分配"谁负责哪个分区"——这就是Rebalance。
用一个牌局比喻:五个人打牌,突然有人上厕所去了(Consumer下线),或者又来了个新朋友(Consumer上线),牌桌上的牌(Partition)就必须重新分配,这就是Rebalance。
【Rebalance触发前后对比】
Rebalance前(3 Consumer,3 Partition):
Consumer A → Partition 0
Consumer B → Partition 1
Consumer C → Partition 2
完美的一对一关系
Rebalance后(Consumer B下线,2 Consumer,3 Partition):
Consumer A → Partition 0, Partition 1
Consumer C → Partition 2
Consumer A需要多负责一个分区
触发条件
Rebalance有下面几种触发场景:
| Consumer加入 | 新的Consumer实例加入Consumer Group | 可能导致重新分配给新Consumer |
| Consumer退出 | Consumer主动调用close()或进程崩溃 | 该Consumer负责的分区需分配给其他Consumer |
| 心跳超时 | session.timeout.ms内未收到Consumer心跳 | GroupCoordinator认为该Consumer"假死" |
| 元数据变更 | 订阅的Topic分区数增加 | 新分区需要分配 |
| Consumer取消订阅 | Consumer主动取消订阅某些Topic | subscription集合收缩 |
二、Rebalance协议演进——从ZooKeeper到两阶段提交
Kafka的Rebalance设计经历了三个版本的演进,每个版本都在解决上一版的问题。
版本一:基于ZooKeeper Watcher(Kafka 0.8)
【版本一架构】
ZooKeeper
┌─────────────┐
│ /consumers/ │
│ group_id/ │
│ ├── ids/ │◄── Consumer 1 注册临时节点
│ ├── owners/│◄── Consumer 2 注册临时节点
│ └── offsets│◄── Consumer 3 注册临时节点
└─────────────┘
↑ Watcher
┌────────────┼────────────┐
│ │ │
Consumer 1 Consumer 2 Consumer 3
问题:
1. 羊群效应(Herd Effect):一个节点变化通知所有Consumer
2. 脑裂(Split Brain):不同Consumer看到ZooKeeper状态不一致
这种方式下,每个Consumer都在ZooKeeper上注册Watcher,任何人上线/下线都会触发所有人的Watcher回调——这就是著名的"羊群效应"。想象一下:一个百人群里,有人@了所有人,每个人都收到一条通知,这就是羊群效应。
版本二:GroupCoordinator集中管理(Kafka 0.8.2)
【版本二架构】
Consumer 1 ──Heartbeat──┐
Consumer 2 ──Heartbeat──┤
Consumer 3 ──Heartbeat──┤
▼
┌────────────────┐
│ GroupCoordinator│ ← 集中管理Consumer Group
│ (某个Broker) │ ← 分区分配在服务端完成
└────────────────┘
改进:
– 羊群效应解决了 ✓
– 脑裂解决了 ✓
新问题:
– 分区分配策略写死在服务端,无法灵活定制 ✗
版本三:客户端分区分配(Kafka 0.9+,当前版本)
【版本三架构——两阶段协议】
阶段一:Join Group
Consumer ──JoinGroupRequest──► GroupCoordinator
◄──JoinGroupResponse── (选出Leader,指定分区策略)
阶段二:Synchronizing Group State
Leader Consumer:
① 根据策略计算分区分配结果
② ──SyncGroupRequest(含分配结果)──► GroupCoordinator
Follower Consumer:
③ ──SyncGroupRequest(空结果)──► GroupCoordinator
GroupCoordinator:
④ 收集所有SyncGroupRequest
⑤ ◄──SyncGroupResponse(含各自的分区)── Consumer
核心变化:分区分配计算工作从服务端移到客户端(Leader Consumer)
优点:自定义PartitionAssignor只需更新客户端,无需改服务端
三、JoinGroup与SyncGroup完整流程——ASCII图解
下面用ASCII图展示Rebalance的完整时序:
【Rebalance完整交互序列图】
Consumer A Consumer B Consumer C GroupCoordinator
(Member1) (Member2) (Member3=Leader)
│ │ │ │
│ ① 发现需要Rebalance (rejoinNeeded=true) │
│ │ │ │
├─JoinGroupRequest─────────────────────────────────────────►│
│ members: [], protocol: {RangeAssignor,RoundRobin} │
│ │ │ │
│ ├─JoinGroupRequest───────────────────────►│
│ │ members: [], protocol: {…} │
│ │ │ │
│ │ ├─JoinGroupRequest────►│
│ │ │ members: [], │
│ │ │ protocol: {…} │
│ │ │ │
│ │ │ ② GroupCoordinator收集完所有 JoinGroupRequest
│ │ │ ③ 选出Leader=Consumer C
│ │ │ ④ 选取partition.assignment.strategy
│ │ │ │
│◄─────────────────────────────────────────JoinGroupResponse│
│ leaderId: Consumer C, members: [Consumer C] │
│ │ │ │
│ ◄─────────────────────────JoinGroupResponse│
│ │ leaderId: Consumer C, members: [Consumer C]│
│ │ │ │
│ │ ◄───────────────JoinGroupResponse│
│ │ │ leaderId: Consumer C, │
│ │ │ members: [A, B, C], │
│ │ │ partitions: […] │
│ │ │ │
│ │ │ ⑤ Leader执行分区分配计算 │
│ │ │ rangeAssignor.assign() │
│ │ │ │
│ │ ├─SyncGroupRequest────────►│
│ │ │ assignments: [{ │
│ │ │ Consumer A → [P0,P1], │
│ │ │ Consumer B → [P2], │
│ │ │ Consumer C → [P3] │
│ │ │ }] │
│ │ │ │
├─SyncGroupRequest────────────────────────────────────────►│
│ assignments: [] (Follower不传分配结果) │
│ │ │ │
│ ├─SyncGroupRequest──────────────────────►│
│ │ assignments: [] │
│ │ │ │
│ │ │ ⑥ GroupCoordinator收到Leader的分配结果
│ │ │ ⑦ 为每个Member生成各自的分配信息
│ │ │ │
│◄─────────────────────────────────────────SyncGroupResponse│
│ assignment: [P0, P1] │
│ │ │ │
│ ◄─────────────────────────SyncGroupResponse│
│ │ assignment: [P2] │
│ │ │ │
│ │ ◄───────────────SyncGroupResponse│
│ │ │ assignment: [P3] │
│ │ │ │
│ ⑧ 状态 → STABLE │ ⑧ 状态 → STABLE │ ⑧ 状态 → STABLE │
│ │ │ │
关键源码:ConsumerCoordinator发送JoinGroupRequest
// ConsumerCoordinator.java (Kafka 0.10.x源码)
private RequestFuture<ByteBuffer> sendJoinGroupRequest() {
// 先将offset提交(防止Rebalance期间消息重复消费)
if (coordinatorUnknown())
lookupCoordinator(); // 找不到GroupCoordinator,先查找
// 构建JoinGroup请求
JoinGroupRequest request = new JoinGroupRequest(
groupId, // Consumer Group ID
this.sessionTimeoutMs, // session超时时间
this.generation.memberId, // memberId
protocolType(), // "consumer"
metadata() // 订阅的Topic+支持的分配策略
);
// 发送请求
return client.send(coordinator, ApiKeys.JOIN_GROUP, request)
.compose(new JoinGroupResponseHandler());
}
四、三种分区分配策略——"谁拿哪副牌"的学问
4.1 RangeAssignor——按Topic范围划分
RangeAssignor是Kafka默认的分配策略,它的核心思路是:按Topic维度,每个Topic的分区按范围均分给订阅了它的Consumer。
【RangeAssignor分配示例】
假设场景:
Topic-A: Partition 0,1,2,3,4,5 (6个分区)
Topic-B: Partition 0,1,2 (3个分区)
Consumer Group: C1, C2, C3
Topic-A分配计算:
6个分区 / 3个消费者 = 每个消费者2个分区
C1 → Partition 0,1
C2 → Partition 2,3
C3 → Partition 4,5
Topic-B分配计算:
3个分区 / 3个消费者 = 每个消费者1个分区
C1 → Partition 0
C2 → Partition 1
C3 → Partition 2
最终分配结果:
C1 → [TopicA-P0, TopicA-P1, TopicB-P0] (3个分区)
C2 → [TopicA-P2, TopicA-P3, TopicB-P1] (3个分区)
C3 → [TopicA-P4, TopicA-P5, TopicB-P2] (3个分区)
看起来挺均匀的!
RangeAssignor的"倾斜"问题:
【RangeAssignor不公平场景】
假设场景:
Topic-A: Partition 0,1,2 (3个分区)
Topic-B: Partition 0,1,2,3,4,5,6,7 (8个分区)
Consumer Group: C1, C2
Topic-A分配:
3 / 2 = 1.5 → C1得2个, C2得1个
C1 → [A-P0, A-P1]
C2 → [A-P2]
Topic-B分配:
8 / 2 = 4
C1 → [B-P0, B-P1, B-P2, B-P3]
C2 → [B-P4, B-P5, B-P6, B-P7]
最终:
C1 → 2 + 4 = 6个分区
C2 → 1 + 4 = 5个分区
当Topic数量多时,C1累积领先,消费负载不均!
RangeAssignor核心源码:
// RangeAssignor.java
public Map<String, List<TopicPartition>> assign(
Map<String, Integer> partitionsPerTopic, // key=topic, value=分区数
Map<String, List<String>> subscriptions) { // key=consumerId, value=订阅的topic列表
Map<String, List<TopicPartition>> assignment = new HashMap<>();
for (Map.Entry<String, List<String>> memberEntry : subscriptions.entrySet()) {
String memberId = memberEntry.getKey();
assignment.put(memberId, new ArrayList<>());
}
// 按Topic维度分配
for (Map.Entry<String, Integer> topicEntry : partitionsPerTopic.entrySet()) {
String topic = topicEntry.getKey();
int numPartitions = topicEntry.getValue();
// 找出订阅了这个Topic的Consumer
List<String> consumersForTopic = new ArrayList<>();
for (Map.Entry<String, List<String>> memberEntry : subscriptions.entrySet()) {
if (memberEntry.getValue().contains(topic))
consumersForTopic.add(memberEntry.getKey());
}
// Range分配核心逻辑
int numConsumers = consumersForTopic.size();
int partitionsPerConsumer = numPartitions / numConsumers;
int consumersWithExtraPartition = numPartitions % numConsumers;
// 前面的消费者多分一个分区
for (int i = 0; i < numConsumers; i++) {
int start = i * partitionsPerConsumer + Math.min(i, consumersWithExtraPartition);
int length = partitionsPerConsumer + (i < consumersWithExtraPartition ? 1 : 0);
String consumer = consumersForTopic.get(i);
for (int p = start; p < start + length; p++) {
assignment.get(consumer).add(new TopicPartition(topic, p));
}
}
}
return assignment;
}
4.2 RoundRobinAssignor——轮询公平分配
RoundRobinAssignor把所有Topic的所有Partition排成一个队列,然后按轮询方式逐个分配给Consumer,就像发扑克牌一样。
【RoundRobinAssignor分配示例】
同一场景:
Topic-A: Partition 0,1,2
Topic-B: Partition 0,1,2,3,4,5,6,7
Consumer Group: C1, C2
所有分区排成队列:
[A-P0, A-P1, A-P2, B-P0, B-P1, B-P2, B-P3, B-P4, B-P5, B-P6, B-P7]
轮询分配(像发扑克牌):
第1轮: C1←A-P0, C2←A-P1
第2轮: C1←A-P2, C2←B-P0
第3轮: C1←B-P1, C2←B-P2
第4轮: C1←B-P3, C2←B-P4
第5轮: C1←B-P5, C2←B-P6
第6轮: C1←B-P7, (没了)
最终结果:
C1 → [A-P0, A-P2, B-P1, B-P3, B-P5, B-P7] (6个分区)
C2 → [A-P1, B-P0, B-P2, B-P4, B-P6] (5个分区)
RoundRobinAssignor的局限:同一Consumer Group内所有Consumer的订阅Topic集合必须完全一致,否则轮询时可能出现分配错误。
4.3 StickyAssignor——粘性保底策略
StickyAssignor是Kafka 0.11引入的策略,核心目标是:在Rebalance时尽可能保持原有分配不变,减少状态迁移。
【StickyAssignor vs RangeAssignor 对比】
初始状态(3 Consumer, 3 Partition):
C1 → [P0]
C2 → [P1]
C3 → [P2]
场景:C2下线,触发Rebalance
RangeAssignor的结果:
C1 → [P0, P1] ← C1原本消费P0, 现在多了P1
C3 → [P2] ← C3不变
变化:C1需要接管P1,导致状态迁移
StickyAssignor的结果:
C1 → [P0] ← 完全不变!
C3 → [P1, P2] ← C3接管P1
优势:
– C1完全不受影响
– P2保持由C3消费,不需要迁移
– 减少了ConsumerRebalanceListener的onPartitionsRevoked触发次数
三种策略对比总结:
| 分配粒度 | 按Topic分别独立分配 | 所有Topic的分区统一轮询 | 最优化的全局分配 |
| 分配均匀性 | 中等(可能倾斜) | 较均匀 | 较均匀 |
| 状态保持 | 差(每次都重新算) | 差(每次都重新算) | 优(尽量保持原分配) |
| 订阅要求 | 消费者可订阅不同Topic | 同一Group必须订阅相同Topic集合 | 消费者可订阅不同Topic |
| 适用场景 | 简单场景,Topic不多 | 统一订阅场景 | 有状态消费,需要减少Rebalance影响 |
| 引入版本 | Kafka 0.9 | Kafka 0.9 | Kafka 0.11 |
五、Rebalance对消费性能的影响——"停止消费"的真相
Rebalance期间,整个Consumer Group会暂停消息消费,这被称为"Stop-The-World"效应。
【Rebalance期间的时间轴】
时间 ──────────────────────────────────────────────────────────►
正常消费期 Rebalance期 恢复正常消费期
[拉消息 处理 提交] [停止消费 协调 分配] [从新offset开始]
▲ ▲ ▲
│ │ │
不断poll() poll()被阻塞 完成SyncGroup后
Consumer状态 恢复消费
→ PREPARING_REBALANCE
→ COMPLETING_REBALANCE
→ STABLE
影响Rebalance耗时的因素:
| session.timeout.ms | 检测Consumer"假死"的时间窗口 | 适当增大(如30s→60s),减少误判 |
| heartbeat.interval.ms | 心跳间隔 | 设为session.timeout.ms的1/3 |
| max.poll.interval.ms | poll()之间的最大间隔 | 增大以避免Consumer被踢出 |
| 分区数量 | 分区越多分配计算越慢 | 合理规划分区数 |
| Consumer数量 | 越多JoinGroup等待越长 | 不要超过分区数 |
缩短Rebalance "空窗期"的5个建议:
六、源码中的状态流转
消费者在Rebalance过程中的状态变化由ConsumerCoordinator管理:
// 消费者端Rebalance状态流转
状态:UNJOINED ──► PREPARING_REBALANCE ──► COMPLETING_REBALANCE ──► STABLE
UNJOINED:
– 初始状态或Consumer离开Group后
– 还未发送JoinGroupRequest
PREPARING_REBALANCE:
– 已发送JoinGroupRequest,等待响应
– needsJoinPrepare=true时,先做onJoinPrepare钩子
COMPLETING_REBALANCE:
– 已收到JoinGroupResponse
– Leader正在计算分区分配
– Follower等待Leader完成计算
– 所有Consumer发送SyncGroupRequest
STABLE:
– 已收到SyncGroupResponse
– 分配结果应用完毕
– 正常消费消息
核心代码:ConsumerCoordinator.poll()中的协调逻辑
// AbstractCoordinator.poll() 核心流程(简化)
public void poll(long now) {
// 1. 检查是否需要触发Rebalance
if (rejoinNeeded) {
ensureActiveGroup(); // 触发JoinGroup
}
// 2. 检查心跳超时
if (needRejoin()) {
rejoinNeeded = true;
return;
}
// 3. 发送心跳
pollHeartbeat(now);
// 4. 自动提交offset
maybeAutoCommitOffsetsAsync(now);
}
本篇小结
本文系统剖析了Kafka Consumer Group Rebalance的完整设计:
- 触发条件:Consumer上下线、心跳超时、Topic分区数变化都会触发Rebalance
- 协议演进:从ZooKeeper Watcher(羊群效应+脑裂)→ GroupCoordinator集中管理(分区策略不灵活)→ 两阶段提交(JoinGroup + SyncGroup,策略下放到客户端)的三代演进
- 三种分配策略:RangeAssignor按Topic范围分配(默认,可能倾斜)、RoundRobinAssignor轮询分配(要求订阅一致)、StickyAssignor粘性分配(减少状态迁移,推荐使用)
- 性能影响:Rebalance期间整个Group停止消费,通过合理配置session.timeout.ms和选择正确的分区策略可以缩短空窗期
理解了Rebalance的机制,下一篇我们将深入ConsumerNetworkClient源码,看看消费者是如何通过它进行网络通信的!
上一篇【第24篇】消息传递保证语义深度解析——at-most-once/at-least-once/exactly-once 下一篇【第26篇】ConsumerNetworkClient源码解析——消费者的"网络大脑"







