阶段二 · 架构原理与调优

分区再平衡深入:Eager,Cooperative 与 Sticky

前置回顾

  • 第 04 篇讲了 rebalance 的概念:触发条件、Eager 全员暂停的外部表现、两种协议的对比表
  • 第 10 篇讲了 Controller 管集群级元数据;本篇主角是另一套机制,Group Coordinator 管消费者组的分区分配
  • 本篇往协议内部拆:分配算法在哪里执行,Sticky 怎么最小化迁移,频繁 rebalance 怎么排查

一句话总结

Rebalance 的分配算法不在 Broker 上,而在消费者客户端执行:Coordinator 只负责组织成员、传递结果.
Eager 要求全员先放弃所有分区再重分,Cooperative 只迁移变动的分区;Sticky assignor 在均衡前提下最小化迁移.三者组合是现代消费端的稳态.

从一次上线事故说起

一个 Go 消费者服务,6 个 Pod,订阅 12 分区的 Topic,每周跟着 K8s 滚动发布上线一次.

发布过程:新 Pod 启动(组成员 +1,触发 rebalance),老 Pod 退出(成员 -1,再触发一次).一轮发布下来,组内震荡十几次.

现象:

  • 发布窗口内消费停停走走,端到端处理延迟从 200ms 飙到 30s
  • lag 曲线呈锯齿:停一段、追一轮、再停
  • 发布结束后还要抖十分钟才回到稳态

要解这个问题,先得看清每次 rebalance 里究竟停了什么.这要从协议本身讲起.


协议舞步:JoinGroup 与 SyncGroup

先纠正一个常见误解:分区分配不是 Broker 算的.Coordinator(第 04 篇讲过它的角色)只负责组织成员,真正的分配算法在一个消费者客户端上执行.

一次 rebalance 的完整通信:

各成员 ──JoinGroup(订阅列表 + 支持的 assignor 列表)──► Coordinator
│
等待全员到齐(或 rebalance 超时),指定第一个加入者为 leader member
│
leader member ◄──── 全部成员的订阅元数据 ───────────┘
│
│ 在本地执行 assignor 算法, 算出每人拿哪些分区
▼
leader member ──SyncGroup(分配方案)──► Coordinator
│
所有成员 ◄──── SyncGroup 响应(各自的分区)────────────┘

三个关键细节:

细节 说明
leader member 是谁 Coordinator 指定第一个加入的成员;与第 09/10 篇的 Broker Leader 无关
算法在哪执行 leader member 的客户端进程里,对 franz-go 就是 balancer 那段代码
何时生效 成员收到 SyncGroup 响应后才切换分配,之前按旧分配继续消费

Coordinator = 会议主持人,负责点名叫齐人,自己不排座位.
leader member = 轮值记录员,按规则算出座位表,交给主持人宣布.

分配逻辑放客户端是历史选择:Broker 专注数据面,分配策略留给各语言客户端自己发挥.代价是行为依赖客户端实现,这正是新一代协议 KIP-848 要把分配收回 Broker 侧的动机,末尾 QA 展开.


Eager:全员停摆的老协议

第 04 篇讲过 Eager 的外部表现(所有人一起停),现在看协议内部为什么必然如此.

分配生效的前提是:每个分区同一时刻最多属于一个成员.为此 Eager 协议要求成员在发 JoinGroup 之前,放弃手上的全部分区:

Eager rebalance 时间线(C3 加入,C1/C2 正在消费):

C1 [P0,P1,P2] C2 [P3,P4,P5] C3 启动
│ │ │
▼ ▼ ▼
┌──────── 全员 revoke 全部 3 个分区 ─────────────────────┐
│ C1 停消费, C2 停消费, 等 C3 加入 │ ← stop-the-world
└───────────────────────────────────────────────────────┘
│ │ │
▼ ▼ ▼
JoinGroup ×3 ───► leader 算分配 ───► SyncGroup 广播结果
│ │ │
▼ ▼ ▼
C1 [P0,P1] C2 [P3,P4] C3 [P2,P5]
恢复消费 恢复消费 开始消费

为什么不能一边消费一边等?因为 leader 计算分配的依据是"当前谁持有什么",这份快照必须静止.否则一个分区可能在旧 owner 还没放手时就被分给新 owner,两个成员同时消费同一分区:offset 提交互相覆盖,消费进度直接损坏.

所以 Eager 的停顿范围是全部分区,哪怕只变动其中两块.


Cooperative:只迁移变动的分区

KIP-429 换了个思路:不要求一次清场,而是分两轮,只搬要动的:

Cooperative rebalance(C3 加入, 同样 6 分区):

Round 1: leader 算出"最小迁移"方案
C1 交出 P2, C2 交出 P5, C3 暂时拿到空(要等 P2/P5 被释放)
│
├── C1 revoke P2, P0/P1 继续消费
└── C2 revoke P5, P3/P4 继续消费
│
▼ 成员重新 JoinGroup(第二轮)

Round 2: P2, P5 已空出
把 P2, P5 分给 C3 ──► SyncGroup 广播 ──► 完成

同一场景下,C1 和 C2 只在自己"迁出的那一块"上短暂中断,其余分区不停.

两个必须知道的代价:

代价 说明
通常需要两轮 follow-up rebalance 是常态,总耗时未必比 Eager 短,优化的是停顿面
协议协商 组内成员的 assignor 支持列表必须有交集,升级需要两次滚动发布

升级路径的坑值得展开.存量集群用 Eager 的 RangeAssignor,想切到 CooperativeSticky:

错误做法: 直接发布"只支持 CooperativeSticky"的新版本
→ 新老成员的支持列表无交集 → 组组建失败(INCONSISTENT_GROUP_PROTOCOL)

正确做法(两次滚动):
第一次滚动: 发布 assignor 列表 = [CooperativeSticky, Range] 的版本
协商仍命中 Range, 行为不变
第二次滚动: 全员就位后, 发布只含 CooperativeSticky 的版本 → 切换生效

franz-go 默认就是 CooperativeStickyBalancer(第 04/05 篇提过),没有存量包袱;从 sarama 或老版本 Java 客户端迁移时才需要走两步流程.


Assignor:分配算法细节

协议决定"什么时候重分",assignor 决定"怎么分".内置策略的差异集中在两点:均衡性、迁移量.

Range:按 topic 逐段切

对每个 Topic 单独切分:

2 个成员, 2 个 Topic, 各 3 分区:

Topic A: C1 → [A0, A1], C2 → [A2]
Topic B: C1 → [B0, B1], C2 → [B2]

结果: C1 拿 4 个, C2 拿 2 个 → 倾斜

单 Topic 时问题不大;多 Topic 时排在前面的成员处处多拿,倾斜累积.

RoundRobin:全局轮转

不区分 Topic,所有分区排队轮流发:

C1 → [A0, A2, B1]
C2 → [A1, B0, B2]
均衡: 3 / 3

均衡性满分,但成员数一变余数全挪位,迁移量大;成员订阅列表不一致时还可能分到无效分区.

Sticky:均衡优先,迁移最小

Sticky 有两个目标,按优先级排队:

  1. 均衡:成员间分区数尽量持平
  2. 粘性:尽量保留成员已持有的分区,只挪必须挪的

Kafka 的实现是启发式:

1. 上一轮分配中仍然有效的部分, 原样保留
2. 未分配的分区, 发给当前负载最轻的成员
3. 仍不均衡时, 移动少量分区补齐(此时才打破粘性)

效果对比,C3 退出、6 分区 3 成员:

Eager + Range/RoundRobin:   C1/C2 的分配全部重算, 大面积迁移
Eager + Sticky: C3 的 [P4,P5] 直接分给 C1/C2, 其余原地不动

Sticky 分配 = 重新排座位:尽量让每个人坐回原位,只挪必须动的少数人.
对比 RoundRobin 的"全员重新抽签",迁移成本天差地别.

CooperativeSticky

现代组合:Sticky 的算法 + Cooperative 的协议.既最小化迁移量,又最小化迁移时的停顿面.这是 franz-go 的默认值,也是本系列一直使用的配置.


Go 实践:三个回调

rebalance 在 franz-go 里暴露为回调,实战代码要处理好三个时机:

client, _ := kgo.NewClient(
kgo.SeedBrokers("localhost:9092"),
kgo.ConsumerGroup("order-processor"),
kgo.ConsumeTopics("order-events"),

// CooperativeSticky(默认值,显式写出便于理解)
kgo.Balancers(kgo.CooperativeStickyBalancer()),

// 分区到手: 初始化分区级状态(本地缓存、定时器等)
kgo.OnPartitionsAssigned(func(ctx context.Context, c *kgo.Client, assigned map[string][]int32) {
for topic, partitions := range assigned {
slog.Info("partitions assigned", "topic", topic, "partitions", partitions)
}
}),

// 分区即将被收回: 最后的机会提交 offset
// 不提交 → 下次拿到这些分区会从旧位置重复消费(第 06 篇的重复窗口)
kgo.OnPartitionsRevoked(func(ctx context.Context, c *kgo.Client, revoked map[string][]int32) {
if err := c.CommitUncommittedOffsets(ctx); err != nil {
slog.Error("commit on revoke failed", "err", err)
}
}),

// 分区意外丢失(心跳超时被踢): 可能已被 fencing, 提交不可靠
kgo.OnPartitionsLost(func(ctx context.Context, c *kgo.Client, lost map[string][]int32) {
slog.Warn("partitions lost", "partitions", lost)
}),
)
回调 时机 典型动作
OnPartitionsAssigned 新分配生效后 初始化状态、打日志
OnPartitionsRevoked 正常移交前 提交 offset、flush 状态
OnPartitionsLost 异常丢失(超时被踢) 只清理本地状态

再叠上第 04 篇的 BlockRebalanceOnPoll,处理中的那批消息在 rebalance 期间不会被夺走:

for {
fetches := client.PollFetches(ctx)
processInParallel(fetches) // 期间 rebalance 被阻塞
client.AllowRebalance() // 处理完毕, 放行
client.CommitUncommittedOffsets(ctx)
}

注意副作用:BlockRebalance 的窗口等于处理耗时,窗口越长,rebalance 收尾越慢(follow-up 轮次被拖住).重活交给 worker pool,主循环尽快回到 Poll.


治频繁 rebalance:Static Membership

回到开头的发布事故.Pod 更替是计划内事件,但协议只看到"成员离开又加入",照常全场震荡.

KIP-345 的 static membership 给出正解:给实例一个固定身份:

kgo.InstanceID("order-consumer-0"),   // StatefulSet 中常用 Pod 名
场景 动态成员(默认) 静态成员(instance id)
实例秒级重启 走完整 rebalance 两趟 保留原分配,零 rebalance
实例消亡(缩容) 立即触发重分配 等 session.timeout 超时才转移
同 ID 双实例并存 不会出现 新实例把旧实例 fencing(发布场景正是要这个)

K8s 上配合 StatefulSet,用有序 Pod 名做 instance id,滚动发布就不再扰动整个组.两个注意点:

  • 缩容时分区要等 session.timeout.ms(默认 45s)才转移,这段时间这些分区无人消费
  • instance id 必须全局唯一且稳定,别用随机数或 IP

频繁 rebalance 排查清单

高频 rebalance 基本逃不出三类原因:

原因 症状 对策
心跳超时(网络抖动、GC 停顿) 成员莫名离组 调大 session.timeout.ms;排查网络与 GC
处理超时(两次 Poll 间隔过长) 分区被夺走,日志出现 lost 控制单批消息量;处理下沉到 worker pool
计划内重启 有规律的成对 rebalance static membership
# 组状态: Stable / PreparingRebalance / CompletingRebalance
kafka-consumer-groups.sh --describe --state --group order-processor \
--bootstrap-server localhost:19092

# 成员与各自持有的分区
kafka-consumer-groups.sh --describe --members --verbose --group order-processor \
--bootstrap-server localhost:19092

第 07 篇的 rebalance-rate-per-hour 指标是核心监控:稳态应趋近 0,持续大于 0 就按上表逐项排查.


动手验证

用第 07 篇的集群,起两个消费者实例(上面的代码加实例日志),做对照实验:

# 第三个终端: 实时观察组状态与成员分配
kafka-consumer-groups.sh --describe --members --verbose \
--group order-processor --bootstrap-server localhost:19092

对照点:

  1. CooperativeSticky 与 Range 分别怎么划分这两个消费者的分区
  2. 杀掉一个消费者,看 revoke 回调收到的分区集合(Cooperative 只含迁移的)
  3. 处理逻辑故意 sleep 超过 max.poll.interval.ms,看分区被夺走的日志
  4. 配置 InstanceID 后重启进程,确认不再出现 rebalance

生产注意事项

别把重处理写进 Poll 循环

主循环的处理时间直接吃满 max.poll.interval(默认 5 分钟).一批 500 条、每条 1s,单批 8 分钟就超线,分区被判"丢失".
正确姿势:Poll → 投递 worker pool → 尽快回来 Poll 下一批.

rebalance 期间的提交失败是常态

迁移中的分区提交 offset 可能收到 REBALANCE_IN_PROGRESS 或 NOT_COORDINATOR.
franz-go 内部会重试;回调里的提交失败要记录留痕,它对应着一次可能的重复消费,别静默吞掉.

稳态组合拳

CooperativeStickyBalancer + 静态成员 + BlockRebalanceOnPoll + 手动提交:
分别解决迁移量、发布扰动、处理安全、重复窗口四个问题,是本系列的推荐基线.

扩容分区 = 主动触发 rebalance

增加分区数必然触发一轮 rebalance,还会打乱第 03 篇讲的 Key 路由(顺序保证被打破).
分区规划放在创建时做足,运行期尽量不动.


Cooperative 就没有停顿了吗?

不是.迁移中的分区仍会短暂中断,只是范围从"全部分区"缩小到"变动的分区".
而且 Cooperative 常常需要 follow-up 第二轮才能到位,总耗时不一定比 Eager 短.它优化的是停顿面,不是总时长.

为什么分配算法在客户端算,而不是 Coordinator?

历史设计:Broker 保持数据面简单,分配策略留给客户端生态.
代价是不同客户端、不同版本的分配行为可能不一致.KIP-848 新协议(3.7 预览、4.0 正式)把分配收回服务端,组内成员不再跑分配算法,rebalance 也更快;Go 客户端生态在跟进,现阶段主流用法仍是经典协议.

同一个 group 里能混用 assignor 吗?

能混用"支持列表",但最终生效的只有一个:所有成员支持列表的交集中,按客户端配置顺序取第一个.
交集为空则组无法组建(报 INCONSISTENT_GROUP_PROTOCOL).升级协议时按上文"两次滚动"走,先加支持,再切默认.

Sticky 一定是最优分配吗?

不保证.它优先均衡、其次最小迁移,是启发式算法,复杂订阅场景下只能算近似解.
对绝大多数场景"足够好",不用纠结理论最优.

rebalance 时 Producer 会受影响吗?

不受影响.Producer 只往 Broker 写,不感知消费者组的分配变化.
看到的现象是消费暂停导致 lag 短期上翘,恢复后自动消化.

静态成员是万能的吗?

不是.它解决"计划内重启",不解决"处理超时"和"网络抖动".
而且实例永久下线(缩容)时,分区要等 session 超时才转移;想缩短这段时间,用优雅关停触发 LeaveGroup(第 05 篇),别 kill -9.


快速回顾

  • 协议舞步:JoinGroup 组织成员,SyncGroup 广播结果,分配算法跑在 leader member 的客户端里
  • Eager:JoinGroup 前全员放弃全部分区,停摆面是全部;胜在协议简单
  • Cooperative:两轮迁移,只 revoke 变动的分区,停顿面最小;代价是 follow-up 与升级兼容性
  • Assignor:Range 多 Topic 倾斜,RoundRobin 均衡但迁移大,Sticky 均衡 + 最小迁移,CooperativeSticky 是默认组合
  • 回调三件套:assign / revoke / lost,revoke 是提交 offset 的最后窗口
  • Static membership:固定 instance id,把计划内重启的 rebalance 降为零
  • 排查:三查,心跳、处理耗时、发布事件;盯 rebalance-rate-per-hour

动手练习

  1. 对照协议:2 个消费者订阅 6 分区 Topic,分别用 CooperativeSticky 和 Range,打印 assigned 分区,验证均衡差异.
  2. 观察 revoke 范围:打印 revoke 回调收到的分区集合,对比 Range(Eager)和 CooperativeSticky 两种 balancer 的差异.
  3. 制造处理超时:把处理逻辑 sleep 到超过 max.poll.interval.ms,观察分区被收回、lag 抖动、恢复的全过程.
  4. 静态成员实验:给消费者配置 InstanceID,杀进程重启,对比配置前后触发 rebalance 的次数.
  5. 看状态机:rebalance 期间用 --describe --state 反复查看,记录 PreparingRebalance 到 Stable 的完整状态迁移.
  6. 模拟发布:起 3 个消费者,逐个重启一轮,统计这轮"发布"造成的 rebalance 次数;配上静态成员再跑一遍对比.

下一篇进入生产调优:批量、压缩、linger.ms、buffer.memory 的联动关系.第 03 篇给过参数速查,下篇讲这些旋钮怎么配合榨出吞吐,又该在什么场景主动放弃吞吐换延迟.