消费者组与Rebalance#
一句话答案#
Consumer Group 内消费者分摊 Partition,Rebalance 在成员变化时重新分配,期间消费暂停需优化避免频繁触发。
核心要点
Rebalance 定义:
当 Consumer Group 中的消费者数量或订阅的 Topic Partition 数量发生变化时,Kafka 会重新分配 Partition 到消费者的映射关系,这个过程叫 Rebalance。
触发条件(三种):
- 消费者加入/离开 Group:新消费者上线、消费者宕机、消费者主动调用
close() - 消费者心跳超时:消费者未在
session.timeout.ms(3.0 起默认 45s,之前 10s)内发送心跳,被 Coordinator 踢出;两次poll()间隔超过max.poll.interval.ms(默认 5min)也会主动离组 - Topic Partition 数量变化:新增 Partition
谁来算分配方案?(高频纠错点)
常见错误是说「Coordinator 计算分区分配」。在经典协议(classic,4.x 客户端仍默认)下,分区分配方案是由 Consumer Group Leader 在客户端计算的,Coordinator(Broker 端)只负责选出 Leader、转发方案、维护成员状态。这样设计的好处是分配逻辑下沉到客户端,加新分配策略不用改 Broker。
注意版本差异:Kafka 4.0 起 GA 的新协议 KIP-848(消费者设
group.protocol=consumer)反过来了——分配改由 Broker 端 Group Coordinator 计算,没有 JoinGroup/SyncGroup,也没有 Group Leader,成员通过心跳(ConsumerGroupHeartbeat)拿到自己的目标分配,按成员增量调整,不再有全组同步屏障。
两阶段协议:JoinGroup + SyncGroup
阶段一 JoinGroup(选 Leader、收集订阅信息)
1. 所有消费者向 Coordinator 发 JoinGroup 请求,上报订阅的 Topic
2. Coordinator 选第一个加入的消费者为 Group Leader
3. Coordinator 把"全部成员列表 + 订阅信息"只回给 Leader
(其他成员收到空响应,知道自己不是 Leader)
阶段二 SyncGroup(下发方案、成员领取)
4. Leader 在本地用分配策略算出"成员→分区"方案,
通过 SyncGroup 请求把方案上交给 Coordinator
5. 其他成员也发 SyncGroup(不带方案),等待领取
6. Coordinator 把方案按成员拆分,分别回给每个消费者
7. 各消费者拿到自己的分区,开始消费plaintextgeneration/epoch 防脑裂:
- 每完成一轮 Rebalance,Coordinator 把 generation id(代次)+1。所有请求都要带当前 generation。
- 若某个「掉队」的老成员(如 GC 卡顿后恢复)拿着过期的 generation 发请求,Coordinator 直接拒绝(
ILLEGAL_GENERATION),逼它重新走 JoinGroup 加入新一代。 - 作用:防止网络分区/假死场景下,新老两批分配方案同时生效导致同一分区被两个消费者消费(脑裂 / 重复消费)。
分配策略:
| 策略 | 特点 |
|---|---|
| RangeAssignor(默认首选) | 按 Topic 维度均匀分配,可能导致部分消费者分到更多 Partition;3.0 起默认配置是 [RangeAssignor, CooperativeStickyAssignor],实际用 Range |
| RoundRobinAssignor | 跨 Topic 轮询分配,更均匀 |
| StickyAssignor | 尽量保持原有分配不变,减少不必要的迁移 |
| CooperativeStickyAssignor | 增量式 Rebalance,不需要 Stop-The-World |
CooperativeSticky 增量再平衡(避免 Stop-The-World):
- Eager 协议每次 Rebalance 都让所有人先撤销全部分区再重新分配,期间全组停止消费(STW)。即使大多数分区最终还分给原主人,也白白经历了一次「全放下→重新拿」。
- Cooperative(协作式)协议把 Rebalance 拆成两轮:
- 第一轮:成员只「上报」自己当前持有的分区,Leader 算出新方案,但只通知那些需要被迁走的分区让其撤销(revoke),不需要变动的分区继续正常消费、不中断。
- 第二轮:被释放出来的分区在第二次 Rebalance 中分配给新的 owner,触发一次后续 Rebalance 完成补领。
- 效果:只有真正发生迁移的少数分区会短暂停顿,绝大多数分区不停消费,避免了全组 STW,大集群扩缩容时抖动显著降低。
Rebalance 的问题与优化:
- 问题:Rebalance 期间整个 Group 停止消费(Eager 模式),影响可用性
- 优化方案:
- 使用
CooperativeStickyAssignor(Kafka 2.4+),实现增量 Rebalance,只迁移需要变更的 Partition - 合理设置
session.timeout.ms和heartbeat.interval.ms,避免误判消费者下线 - 合理设置
max.poll.interval.ms,避免消费逻辑太慢导致被踢出 - 静态成员(Kafka 2.3+,配置
group.instance.id):实例重启后在 session 超时内回来,不触发 Rebalance,适合滚动发布 - Kafka 4.0+ 可切到 KIP-848 新协议(
group.protocol=consumer),Broker 端增量分配,彻底去掉全组 STW
- 使用
面试回答(2分钟版)
Consumer Group 是 Kafka 实现消费负载均衡的核心机制,同一个 Group 内的消费者分摊消费 Topic 的 Partition,每个 Partition 同一时刻只被 Group 内一个消费者消费。当 Group 内消费者数量变化、消费者心跳超时被踢出、或 Topic Partition 数变化时,会触发 Rebalance 重新分配 Partition 到消费者的映射。传统的 Eager 协议 Rebalance 是 Stop-The-World 的——所有消费者先撤销当前分配再重新分配,期间整个 Group 停止消费,影响可用性。分配策略有四种:RangeAssignor 按 Topic 维度分配是默认策略,RoundRobinAssignor 跨 Topic 轮询更均匀,StickyAssignor 尽量保持原有分配减少迁移,CooperativeStickyAssignor 是 Kafka 2.4 引入的增量式 Rebalance 只迁移需要变更的 Partition 不用全停。生产中优化 Rebalance 的关键是合理设置 session.timeout.ms 和 max.poll.interval.ms 避免误判下线,推荐使用 CooperativeStickyAssignor 减少 Rebalance 对消费的影响。Kafka 4.0 起新消费者组协议 KIP-848 已经 GA,分配改由 Broker 端计算、按成员增量调整,不再有全组停顿,消费者配 group.protocol=consumer 即可启用。
追问与易错
追问方向:
- “Rebalance 有什么问题?怎么减少影响?”→ Eager 协议下 Rebalance 期间所有消费者停止消费(Stop The World),延迟高;可通过增大 session.timeout.ms、合理设置 max.poll.interval.ms、静态成员、避免频繁上下线来减少触发频率,并用 CooperativeSticky 或 4.0+ 的 KIP-848 新协议减少停顿
- “消费者数大于分区数会怎样?”→ 普通消费者组里多余的消费者处于空闲状态,不会被分配到任何分区;所以消费者数不要超过分区数,否则浪费资源。例外是 4.2 起生产可用的共享组(Share Group,Queues for Kafka),多个消费者可以按条分摊同一分区,但不保序
- “Kafka 有哪些分区分配策略?”→ Range(按范围均分,可能不均匀)、RoundRobin(轮询,较均匀)、Sticky(尽量保持原分配减少迁移)、CooperativeSticky(增量 Rebalance 不停消费)
易错点:
- ❌ Rebalance 只在消费者宕机时发生——消费者加入或离开、订阅的 Topic 分区数变化、两次 poll 间隔超过
max.poll.interval.ms都会触发;滚动发布时每个实例重启都会触发一次(静态成员可避免) - ❌ 处理慢被踢出组,就调大
session.timeout.ms——处理慢由max.poll.interval.ms判定,和心跳无关;应该调小max.poll.records或调大max.poll.interval.ms - ❌ 用了 CooperativeSticky 就完全没有停顿——增量协作只让需要迁移的分区暂停,其余分区继续消费;被迁移的分区仍要等两轮 Rebalance 完成才恢复