TL;DR(30 秒扫完)
- 本质:Consumer Group 内部对 Partition 的订阅关系重新协商;消息固定不变
- 三类触发:组成员变动 / 订阅关系变化 / 分区数变化
- 两种模式:Eager(旧,STW) vs Cooperative(0.11+,增量式,生产推荐)
- 三种算法:RangeAssignor(默认)/ RoundRobinAssignor / CooperativeStickyAssignor
- 心跳参数:
session.timeout.ms(默认 10s)+max.poll.interval.ms(默认 5min)+heartbeat.interval.ms - 副作用:Rebalance 期间整个组暂停消费,可能重复消费(新消费者从上次 commit offset 开始)
关键结论
结论 ARebalance 的主体是 Consumer Group,不是分区也不是消息
结论 BEager 模式 STW 严重,CooperativeStickyAssignor 是生产首选
结论 C
max.poll.interval.ms 超时会误伤——业务处理慢被判定"卡死",触发 Rebalance结论 DRebalance 期间可能重复消费,业务必须做幂等
完整讲解(费曼四步)
STEP 1 · 概念
Rebalance 是 Consumer Group 内部对 Partition 订阅关系的重新协商。Partition 归属不变(消息固定不变),变的是"哪个消费者负责消费哪个 Partition"。STEP 2 · 大白话
餐厅服务员比喻:- 每个分区 = 一张桌子
- 消费者 = 服务员
- Consumer Group = 一个班的员工
- Rebalance = 排班调整,重新分配每张桌子给哪个服务员——桌子(消息)不动,动的是服务员(消费者)
STEP 3 · 底层
三类触发条件
| 触发条件 | 场景 |
|---|---|
| 组成员变动 | 新增消费者、消费者退出(正常 shutdown 或崩溃)、心跳超时 |
| 订阅关系变化 | subscribe() 的 topic 列表变了 |
| 分区数变化 | Topic 分区被扩容 |
心跳机制
每个 Consumer → 周期性发心跳 → Group Coordinator
↓
超过 session.timeout.ms(默认 10s)未收到
↓
Coordinator 判定成员失效
↓
触发 Rebalance
Rebalance 流程(JoinGroup 协议)
1. JoinGroup:所有成员加入 → Leader 收到
2. Leader 收集所有成员信息 → 计算分配方案
3. SyncGroup:Leader 广播方案 → 各成员按方案订阅 Partition
Group Leader:一般由第一个 join 的消费者担任。
两种 Rebalance 模式
| 模式 | 时机 | 影响 |
|---|---|---|
| Eager(旧) | 所有 Partition 全部释放再重新分配 | STW 时间长,整个组暂停 |
| Cooperative(0.11+) | 只调整必要变化的 Partition | 未变的 Partition 继续消费 |
三种分配算法
| 算法 | 特点 | 推荐度 |
|---|---|---|
| RangeAssignor(默认) | 按 range 平均分配 | 老版本默认,可能热点 |
| RoundRobinAssignor | 轮询分配 | 均衡但每次变化大 |
| StickyAssignor | 尽量保持原有分配 | 减少变更 |
| CooperativeStickyAssignor | 增量式,最优 | ✅ 生产推荐 |
关键配置参数
| 参数 | 默认值 | 说明 |
|---|---|---|
session.timeout.ms | 10000 | 心跳超时,判定成员失效 |
heartbeat.interval.ms | 3000 | 心跳间隔(< session.timeout/3) |
max.poll.interval.ms | 300000 | 两次 poll 最大间隔,超过判定卡死 |
max.poll.records | 500 | 单次 poll 最大记录数 |
max.poll.interval.ms 的坑:- 消费者处理上一批消息太慢 → 超过间隔 → 被认为"卡死" → 触发 Rebalance
- 解决:调大该值或减小
max.poll.records
Rebalance 期间的副作用
- 重复消费:原消费者 Offset 未 commit,新消费者从上次 commit offset 重新消费
- 短暂停顿:Eager 模式下整个组暂停;Cooperative 下未变的分区继续
- 必须幂等:消费端要做去重(DB 唯一索引、Redis 记录、去重表)
STEP 4 · 简化
一句话总结:Rebalance 是 Consumer Group 内 Partition 订阅关系重新协商,Partition 消息不动,动的是消费者分配。 记忆口诀:- 主体:Consumer Group(不是分区,不是消息)
- 三触发:成员变、订阅变、分区变
- 两模式:Eager(STW)vs Cooperative(增量)
- 参数坑:
max.poll.interval.ms超时会误伤 - 副作用:Rebalance 期间可能重复消费
常见误区
认为 Rebalance 是"消息在分区之间迁移"
Partition 消息固定不变,变的是消费者分配
认为 Rebalance 是分区层的行为
Rebalance 的主体是 Consumer Group
不知道 Cooperative 增量式算法
0.11+ 有 CooperativeStickyAssignor,避免 STW
调小
max.poll.interval.ms 解决堆积可能反而更频繁触发 Rebalance,应该看根因
认为 Rebalance 不会影响消费
Eager 模式下整个组暂停消费
延伸追问
max.poll.interval.ms 超时也会触发 Rebalance,为什么?表示消费者处理上一批消息太慢,被认为"卡死",其他成员接手;调大该值或减小
max.poll.records。Cooperative Rebalance 怎么做到"增量式"的?
成员先声明"想保留哪些 Partition",Leader 计算最小变动方案,只让变化的 Partition 让出。
Rebalance 期间消费者会重复消费吗?
会。因为原消费者 Offset 还没 commit,新消费者从上次 commit 的 Offset 开始消费。
如何减少 Rebalance 频率?
用 CooperativeStickyAssignor;调大
session.timeout.ms 和 max.poll.interval.ms;减小 max.poll.records;避免频繁增删消费者。Cooperative 和 Sticky 什么区别?
Cooperative 是"增量式重平衡",Sticky 是"保持原有分配不变";CooperativeStickyAssignor 是两者结合的最优实现。
速查表
主体: Consumer Group(不是分区,不是消息)
触发: 组成员变动 / 订阅变化 / 分区数变化
模式: Eager(STW)vs Cooperative(增量式,生产推荐)
算法: Range / RoundRobin / Sticky / CooperativeSticky
心跳: session.timeout=10s / heartbeat.interval=3s / max.poll.interval=5min
坑: max.poll.interval.ms 超时会误伤;Rebalance 期间可能重复消费
Anki 候选卡片
Q: Kafka Rebalance 的主体是什么?
A: Consumer Group,不是分区也不是消息
Q: Kafka Rebalance 的三类触发条件?
A: 组成员变动 / 订阅关系变化 / 分区数变化
Q: Kafka 0.11+ 推荐的重平衡算法?
A: CooperativeStickyAssignor(增量式)
Q:
max.poll.interval.ms 超时会导致什么?A: 消费者被判定卡死,触发 Rebalance
Q: Rebalance 期间的副作用?
A: 可能重复消费(新消费者从上次 commit offset 开始),必须做幂等
关联题目
关联知识
Rebalance 是 Consumer Group 内 Partition 订阅关系重新协商,不是消息迁移
餐厅排班:桌子(分区)不动,动的是服务员(消费者)的分配
✦ 记 忆 口 诀 ✦
主体是 Consumer Group;三触发、两模式、Cooperative 生产推荐
关键可视化
Rebalance 流程(JoinGroup 协议)
flowchart TD A[所有成员 JoinGroup] --> B[Leader 收到] B --> C[Leader 收集成员信息] C --> D[Leader 计算分配方案] D --> E[SyncGroup 广播方案] E --> F[各成员按方案订阅 Partition] F --> G[开始消费]
Eager vs Cooperative 对比
flowchart LR
subgraph Eager 旧模式
A1[所有 Partition 全部释放] --> B1[重新分配]
B1 --> C1[STW 整个组暂停]
end
subgraph Cooperative 0.11 增量式
A2[成员声明想保留哪些 Partition] --> B2[Leader 计算最小变动]
B2 --> C2[只让变化的 Partition 让出]
end心跳与失效判定
flowchart TD
A[Consumer 周期性发心跳] --> B[Group Coordinator]
B --> C{超过 session.timeout.ms 默认 10s 未收到}
C -->|是| D[判定成员失效]
D --> E[触发 Rebalance]
C -->|否| F[继续监控]知识关系
⬆️ 前置(Prerequisite)
kafka-message-structure🔄 延伸(Extends)
暂无⚡ 对比(Contrast)
暂无
🎯 概念
📏 规则
⚠️ 误区
🔍 追问
✨ 口诀
共 0 张卡,点击翻面