Kafka 消费者组与 Rebalance

mq 📚 learning kafka-rebalance · kafka · consumer-group · rebalance · cooperative-sticky · consumer-assignment

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 是生产首选
结论 Cmax.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 继续消费
Cooperative 的关键:成员先声明"想保留哪些 Partition",Leader 计算最小变动方案,只让变化的 Partition 让出。

三种分配算法

算法特点推荐度
RangeAssignor(默认)按 range 平均分配老版本默认,可能热点
RoundRobinAssignor轮询分配均衡但每次变化大
StickyAssignor尽量保持原有分配减少变更
CooperativeStickyAssignor增量式,最优✅ 生产推荐

关键配置参数

参数默认值说明
session.timeout.ms10000心跳超时,判定成员失效
heartbeat.interval.ms3000心跳间隔(< session.timeout/3)
max.poll.interval.ms300000两次 poll 最大间隔,超过判定卡死
max.poll.records500单次 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 开始),必须做幂等

关联题目

  • ⚠️ 《什么是 Kafka 的重平衡机制?》— Round 2 Q3, ⭐(把 Rebalance 理解成消息迁移)
  • 关联知识

    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 张卡,点击翻面