Rebalance 是 Kafka 消费者组最复杂、也最容易引发生产事故的机制。本文从源码层面完整走读 Rebalance 的协议流程(JoinGroup → SyncGroup → Heartbeat),深入解析 Eager 和 Cooperative 两种 Rebalance 协议的差异,揭示 Stop-the-World 问题的根源——为什么一次 Rebalance 能让消费暂停数十秒。包含 Consumer 端状态机、Coordinator 端处理逻辑、Rebalance 触发条件全景图,以及生产环境减少 Rebalance 影响的 6 条最佳实践。
一、Rebalance 是什么,为什么它这么痛
1.1 一句话定义
Rebalance(再均衡) 是 Kafka 消费者组在成员变化时,重新分配分区归属的过程。
触发 Rebalance 的事件:
1.2 为什么 Rebalance 是痛点
Rebalance 的影响时间线:

总暂停时间:25 秒 ← 对实时业务来说就是一次事故
Eager Rebalance(全量暂停) 在大集群中可能暂停 30-60 秒,这是 Kafka 被诟病最多的设计。
二、Rebalance 的协议全景
2.1 三大核心协议
Rebalance 通过三个核心 API 协议完成:

2.2 协议详解
① JoinGroup:加入消费者组
// JoinGroupRequest 结构
{
"groupId": "my-consumer-group",
"memberId": "client-uuid-xxx", // 消费者 ID
"sessionTimeoutMs": 30000, // 会话超时
"rebalanceTimeoutMs": 300000, // Rebalance 超时(等成员加入的最长时间)
"protocolType": "consumer", // 协议类型
"protocols": [ // 支持的分配策略
{
"name": "cooperative-sticky", // 策略名称
"metadata": <userData> // 自定义元数据(StickyAssignor 用来传上次分配)
}
]
}
Coordinator 收到 JoinGroup 后的处理逻辑:
// GroupCoordinator.scala(简化)
def handleJoinGroup(groupId, memberId, protocols, …): Unit = {
val group = groupManager.getGroup(groupId) match {
case None =>
// 组不存在 → 创建新组
val newGroup = new DelayedHeartbeatGroup(…)
groupManager.addGroup(newGroup)
newGroup
case Some(existing) => existing
}
// 检查组成员变化
if (isNewMember || memberLeft || topicChanged) {
// 触发 Rebalance
group.prepareRebalance() // 状态 → PreparingRebalance
}
// 等待所有成员加入
// 第一个加入的成员成为 Leader
if (group.allMembersJoined) {
// Leader 执行分配
val assignment = assignor.assign(
group.partitionsPerTopic,
group.subscriptions // 所有成员的订阅信息
)
// 发送 JoinGroup 响应给所有成员
// Leader 收到所有成员的信息
// Follower 只收到自己的分配结果
group.sendJoinGroupResponse(assignment)
}
}
② SyncGroup:确认分配方案
// SyncGroupRequest 结构(Leader 发送分配方案)
{
"groupId": "my-consumer-group",
"memberId": "client-uuid-xxx",
"generationId": 5, // 第几代 Rebalance
"groupAssignment": { // Leader 发送完整的分配方案
"client-uuid-1": [topic-a-0, topic-a-3],
"client-uuid-2": [topic-a-1, topic-a-4],
"client-uuid-3": [topic-a-2, topic-a-5],
}
}
// SyncGroupResponse(每个消费者收到自己的分配)
{
"memberAssignment": [topic-a-0, topic-a-3], // 我的分区
"errorCode": 0
}
③ Heartbeat:维持心跳
// 心跳请求
{
"groupId": "my-consumer-group",
"generationId": 5, // 当前 Rebalance 代次
"memberId": "client-uuid-xxx"
}
// Coordinator 的心跳处理
def handleHeartbeat(groupId, memberId, generationId): Unit = {
val group = groupManager.getGroup(groupId)
// 检查 generationId 是否匹配
if (group.generationId != generationId) {
// 代次不匹配 → 消费者需要重新 JoinGroup
return RebalanceInProgress
}
// 检查会话是否过期
if (group.isExpired(memberId)) {
return UnknownMember
}
// 更新心跳时间
group.updateHeartbeat(memberId)
// 检查是否需要触发 Rebalance
if (group.needsRebalance) {
return RebalanceInProgress // 通知消费者重新 JoinGroup
}
return NoError
}
三、Consumer 端状态机
3.1 完整状态流转
┌──────────────────────────────────────────────────────┐
│ Consumer Rebalance 状态机 │
│ │
│ ┌──────────┐ │
│ │ UNINIT │ ← 初始状态 │
│ └────┬─────┘ │
│ │ subscribe() │
│ ▼ │
│ ┌──────────┐ poll() ┌──────────────┐ │
│ │ STABLE │ ←─────────────│ REBALANCING │ │
│ │ (正常消费)│ │ (Rebalance中) │ │
│ └────┬─────┘ └──────┬───────┘ │
│ │ │ │
│ │ 需要Rebalance │ │
│ │ (成员变化/Topic变化) │ │
│ └──────────────────────────────┘ │
│ │ │
│ ┌─────────▼──────────┐ │
│ │ onPartitionsRevoked│ │
│ │ (提交offset/清理) │ │
│ └─────────┬──────────┘ │
│ │ │
│ ┌─────────▼──────────┐ │
│ │ JoinGroup │ │
│ │ (等待Coordinator) │ │
│ └─────────┬──────────┘ │
│ │ │
│ ┌─────────▼──────────┐ │
│ │ SyncGroup │ │
│ │ (收到新分配) │ │
│ └─────────┬──────────┘ │
│ │ │
│ ┌─────────▼──────────┐ │
│ │ onPartitionsAssigned│ │
│ │ (恢复offset/重建) │ │
│ └─────────┬──────────┘ │
│ │ │
│ ┌─────────▼──────────┐ │
│ │ STABLE │ │
│ │ (恢复正常消费) │ │
│ └────────────────────┘ │
└──────────────────────────────────────────────────────┘
3.2 poll() 方法中的 Rebalance 逻辑
// KafkaConsumer.poll() 的简化逻辑(KafkaConsumer.java)
public ConsumerRecords<K, V> poll(Duration timeout) {
// 1. 检查是否需要 Rebalance
if (coordinator.needsRebalance()) {
// 触发 Rebalance
coordinator.poll(time.milliseconds());
// 注意:这里可能阻塞很长时间!
}
// 2. 正常拉取数据
Map<TopicPartition, List<ConsumerRecord<K, V>>> records =
fetcher.fetchedRecords();
// 3. 检查心跳(在 poll 循环中维护心跳)
coordinator.maybeHeartbeat();
return new ConsumerRecords<>(records);
}
// ConsumerCoordinator.poll() 的 Rebalance 逻辑
public void poll(long now) {
// 1. 发送心跳
maybeHeartbeat();
// 2. 检查是否需要 JoinGroup
if (state == MemberState.PREPARING_REBALANCE) {
// 发送 JoinGroup 请求(阻塞等待响应)
sendJoinGroupRequest();
// 这里会阻塞,直到 Coordinator 返回响应
// 阻塞时间 = rebalanceTimeoutMs(默认 5 分钟!)
}
// 3. 检查是否需要 SyncGroup
if (state == MemberState.COMPLETING_REBALANCE) {
sendSyncGroupRequest();
}
// 4. 如果 Rebalance 完成,执行回调
if (state == MemberState.STABLE) {
// onPartitionsAssigned 回调在这里执行
invokePartitionsAssigned(assignment);
}
}
3.3 RebalanceListener 回调时序
// 消费者注册 RebalanceListener
consumer.subscribe(topics, new ConsumerRebalanceListener() {
@Override
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
// 分区被撤销前调用
// 关键:在这里提交 offset,否则可能丢消息!
consumer.commitSync();
log.info("分区被撤销: {}", partitions);
// 执行清理操作(如关闭文件句柄、释放资源)
}
@Override
public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
// 分区被分配后调用
log.info("分区被分配: {}", partitions);
// 恢复 offset(如果用 auto.offset.reset 可能跳过消息)
// 初始化状态(如 Flink 状态恢复)
}
});
Eager vs Cooperative 回调时序对比:
Eager Rebalance (传统):
时间线:
t0: 所有消费者收到 onPartitionsRevoked(所有分区)
→ 所有分区暂停消费 ← 全员 Stop-the-World!
t1: JoinGroup + SyncGroup
t2: 所有消费者收到 onPartitionsAssigned(新分配的分区)
→ 恢复消费
特点:先全部撤销,再全部分配
问题:即使某个消费者的分区没有变化,也被撤销又重新分配
Cooperative Rebalance (增量):
时间线:
t0: 只有需要迁移的分区的消费者收到 onPartitionsRevoked(部分分区)
→ 只有被撤销的分区暂停 ← 大部分分区继续消费!
t1: JoinGroup + SyncGroup (第一轮)
t2: 新消费者收到 onPartitionsAssigned(被撤销的分区)
→ 恢复消费
特点:只撤销需要迁移的分区,其他分区不受影响
优势:Stop-the-World 范围最小化
四、Coordinator 端处理逻辑
4.1 Group 状态机

4.2 GroupCoordinator 核心处理逻辑
// GroupCoordinator.scala(核心逻辑简化)
class GroupCoordinator {
def handleJoinGroup(
groupId: String,
memberId: String,
protocols: List[(String, ByteBuffer)],
sessionTimeoutMs: Int,
rebalanceTimeoutMs: Int
): JoinGroupResponse = {
val group = groupManager.getGroup(groupId) match {
case Some(g) => g
case None =>
// 新组 → 创建
val g = new GroupMetadata(groupId)
groupManager.addGroup(g)
g
}
group.inLock {
group.currentState match {
case Dead =>
// 组不存在 → 返回错误
JoinGroupResponse(UNKNOWN_GROUP_ID)
case Empty | Stable =>
// 正常状态 → 检查是否需要 Rebalance
val member = group.getOrCreateMember(memberId, protocols)
if (group.hasNewMember || group.topicChanged) {
// 新成员加入或 Topic 变化 → 触发 Rebalance
group.transitionTo(PreparingRebalance)
prepareRebalance(group)
} else {
// 无变化 → 直接返回当前分配
JoinGroupResponse(SUCCESS, group.generationId,
group.leaderId, group.assignment)
}
case PreparingRebalance =>
// 正在等待成员加入
val member = group.addMember(memberId, protocols)
if (group.allMembersJoined(rebalanceTimeoutMs)) {
// 所有成员已加入 → 选出 Leader,执行分配
group.transitionTo(CompletingRebalance)
val leader = group.leader
// Leader 执行 Assignor.assign()
val assignment = performAssignment(group, leader)
// 返回分配结果
JoinGroupResponse(SUCCESS, group.generationId,
group.leaderId, assignment)
} else {
// 还在等待其他成员 → 延迟响应
// 消费者端会阻塞在 JoinGroup 调用上
delayJoinGroupResponse(group, member)
}
case CompletingRebalance =>
// 之前正在分配 → 新成员来了,需要重新 Rebalance
group.transitionTo(PreparingRebalance)
prepareRebalance(group)
delayJoinGroupResponse(group, group.getMember(memberId))
}
}
}
def handleSyncGroup(
groupId: String,
memberId: String,
generationId: Int,
groupAssignment: Map[String, Assignment]
): SyncGroupResponse = {
val group = groupManager.getGroup(groupId).get
group.inLock {
if (group.generationId != generationId) {
// 代次不匹配 → 需要重新 JoinGroup
return SyncGroupResponse(REBALANCE_IN_PROGRESS)
}
if (memberId == group.leaderId) {
// Leader 发送了完整的分配方案
group.storeAssignment(groupAssignment)
// 通知所有成员
group.allMembers.foreach { member =>
val assignment = groupAssignment.get(member.memberId)
member.completeSync(assignment)
}
group.transitionTo(Stable)
}
// 返回该成员的分配
SyncGroupResponse(SUCCESS, group.assignment.get(memberId))
}
}
def handleHeartbeat(
groupId: String,
memberId: String,
generationId: Int
): HeartbeatResponse = {
val group = groupManager.getGroup(groupId) match {
case Some(g) => g
case None => return HeartbeatResponse(UNKNOWN_GROUP_ID)
}
group.inLock {
group.currentState match {
case Dead => HeartbeatResponse(UNKNOWN_GROUP_ID)
case Empty => HeartbeatResponse(UNKNOWN_MEMBER_ID)
case PreparingRebalance =>
// 正在 Rebalance → 通知消费者重新 JoinGroup
HeartbeatResponse(REBALANCE_IN_PROGRESS)
case CompletingRebalance =>
// 分配方案还没确认
HeartbeatResponse(REBALANCE_IN_PROGRESS)
case Stable =>
// 检查会话超时
if (group.isSessionExpired(memberId)) {
HeartbeatResponse(UNKNOWN_MEMBER_ID)
} else {
// 检查是否需要触发新的 Rebalance
if (group.needsRebalance) {
group.transitionTo(PreparingRebalance)
HeartbeatResponse(REBALANCE_IN_PROGRESS)
} else {
group.updateHeartbeat(memberId)
HeartbeatResponse(SUCCESS)
}
}
}
}
}
}
五、Rebalance 触发条件全景
5.1 触发条件分类

5.2 每种触发的处理逻辑
| 新消费者加入 | Coordinator | JoinGroup 请求检测到新 memberId | 正常行为,不可避免 |
| 消费者主动退出 | Consumer | close() 时发送 LeaveGroup | 正常行为 |
| 订阅 Topic 变化 | Consumer | subscribe() 变化 → 下次 poll 触发 | 正常行为 |
| 心跳超时 | Coordinator | session.timeout.ms 内无心跳 | ✅ 调大 session.timeout |
| poll 间隔超时 | Coordinator | max.poll.interval.ms 内无 poll | ✅ 调大或调小 max.poll.records |
| 分区数变化 | Coordinator | 分区数增加 → 检测到新分区 | 不可避免 |
| GC 停顿 | Coordinator | 间接导致心跳超时 | ✅ 优化 JVM |
| 网络抖动 | Coordinator | 心跳包丢失 | ✅ 增加心跳频率 |
| Coordinator 切换 | Broker | Group Coordinator 所在 Broker 宕机 | ✅ KRaft 模式缓解 |
六、生产环境减少 Rebalance 影响的 6 条实践
实践 1:切换到 CooperativeStickyAssignor
# 从 Eager Rebalance 切换到 Cooperative Rebalance
# 分两步滚动升级
# Step 1: 同时配置两种策略(过渡期)
partition.assignment.strategy=\\
org.apache.kafka.clients.consumer.CooperativeStickyAssignor,\\
org.apache.kafka.clients.consumer.RangeAssignor
# Step 2: 全部消费者升级后,移除 RangeAssignor
partition.assignment.strategy=\\
org.apache.kafka.clients.consumer.CooperativeStickyAssignor
实践 2:合理设置超时参数
# 会话超时:GC 停顿 30 秒内不会被误判宕机
session.timeout.ms=30000
heartbeat.interval.ms=10000 # session 的 1/3
# poll 间隔:确保 > 单批处理时间
# 经验值:单批处理时间 × 3
max.poll.interval.ms=600000 # 10 分钟
# 单批拉取量:确保在 poll interval 内能处理完
# max.poll.records × 单条处理时间 < max.poll.interval.ms
max.poll.records=100
实践 3:异步处理 + 仅 poll 心跳
// 问题:如果消息处理慢,会超过 max.poll.interval → 被踢出组
// 错误做法(同步处理)
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
// 如果处理这批数据超过 5 分钟 → Rebalance!
for (ConsumerRecord<String, String> record : records) {
processMessage(record); // 慢操作
}
}
// 正确做法(异步处理 + 心跳维护)
ExecutorService executor = Executors.newFixedThreadPool(4);
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
// 异步提交处理
Future<?> future = executor.submit(() -> {
for (ConsumerRecord<String, String> record : records) {
processMessage(record);
}
});
// 主线程不阻塞 → 心跳正常 → 不会触发 Rebalance
// 但需要确保处理完成后才提交 offset
}
实践 4:使用 Static Membership(Kafka 2.3+)
# Static Membership:消费者固定 memberId
# 重启后不需要 Rebalance(在 session.timeout 内)
group.instance.id=consumer-1
# 配合更大的 session timeout
session.timeout.ms=300000 # 5 分钟
Static Membership 的效果:
普通模式:
消费者重启 → LeaveGroup → Rebalance → 其他消费者暂停
Static Membership:
消费者重启 → Coordinator 认为只是临时离线
→ 在 session.timeout 内重启回来 → 不触发 Rebalance
→ 其他消费者完全无感知
适合:消费者需要频繁重启的场景(如部署更新)
实践 5:监控 Rebalance 频率
// 通过 JMX 监控 Rebalance 次数
// kafka.consumer:type=coordinator-metrics,name=rebalance-rate-per-sec
// 或者在代码中监听
consumer.subscribe(topics, new ConsumerRebalanceListener() {
private final AtomicInteger rebalanceCount = new AtomicInteger(0);
@Override
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
int count = rebalanceCount.incrementAndGet();
log.warn("Rebalance #{} 触发,分区被撤销: {}", count, partitions);
// 上报到监控系统
metricsReporter.increment("kafka.rebalance.count");
metricsReporter.gauge("kafka.rebalance.revoked.partitions",
partitions.size());
}
@Override
public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
log.info("Rebalance 完成,分区被分配: {}", partitions);
metricsReporter.gauge("kafka.rebalance.assigned.partitions",
partitions.size());
}
});
告警阈值:
正常: < 3 次/天
警告: 3-10 次/天
严重: > 10 次/天 → 检查消费者配置和 GC 日志
实践 6:避免在 onPartitionsRevoked 中做耗时操作
// 错误做法:在回调中做耗时操作
@Override
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
// ❌ 这些操作太慢,延长 Rebalance 时间
flushAllBuffers(); // 可能要几秒
closeAllFileHandles(); // 可能要几秒
commitSync(); // 可能要几秒
// 总计可能 10-20 秒 → 其他消费者都在等!
}
// 正确做法:只做必要的操作
@Override
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
// ✅ 只提交 offset(异步提交也行)
consumer.commitSync(); // 这一步必须做
// 其他清理操作放到后台线程异步执行
cleanupExecutor.submit(() -> {
flushAllBuffers();
closeAllFileHandles();
});
}
七、Rebalance 问题排查清单
现象:消费暂停,日志出现 "Attempt to heartbeat failed since group is rebalancing"
排查步骤:
1. 检查 Rebalance 触发原因
→ 查看 Consumer 日志中 Rebalance 前的最后一条日志
→ 是正常部署(消费者重启)还是异常(心跳超时)?
2. 检查 max.poll.interval.ms
→ 日志中是否有 "max.poll.interval.ms expired"?
→ 单批处理时间是否超过 max.poll.interval.ms?
3. 检查 GC 日志
→ 是否有 > 10 秒的 GC 停顿?
→ Full GC 会导致心跳超时
4. 检查网络
→ 消费者到 Broker 的网络延迟是否正常?
→ 心跳是否被网络抖动丢失?
5. 检查 Coordinator
→ Group Coordinator 所在的 Broker 是否宕机?
→ kafka-consumer-groups.sh –describe –group <group> –state
6. 检查消费者数量
→ 是否有大量消费者同时重启?(如 K8s 滚动更新)
→ 考虑使用 Static Membership
八、Rebalance 核心知识速查
协议流程: JoinGroup → SyncGroup → Heartbeat 两种模式:
Stop-the-World 根源:
- Eager 模式下,所有分区在 JoinGroup/SyncGroup 期间暂停
- onPartitionsRevoked 中的耗时操作延长暂停时间
生产环境最佳实践:
1. CooperativeStickyAssignor(增量 Rebalance) 2. 合理设置 session.timeout / max.poll.interval 3. 异步处理消息,主线程只做 poll + heartbeat 4. Static Membership(频繁重启场景) 5. 监控 Rebalance 频率(< 3 次/天为正常) 6. onPartitionsRevoked 只做 commit offset
专栏导航:
-
AI 推理优化系列—vLLM PagedAttention 解析:显存利用率从 40% 提升到 90% 的秘密
-
Kafka 深度解剖 ①:StickyAssignor 分区分配策略
Kafka 深度解剖,覆盖分区分配、Rebalance、Offset 提交、Exactly-Once 语义、Producer 幂等与事务。




