KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA
04 · 消费组与再平衡 — keel 龙骨
消费组卡住的那几秒,是 Kafka 排障里最常被误判的地方。有人看到「消费停了 10 秒」就以为是下游慢,实际是重平衡;有人看到 lag 涨了就去加消费者,实际是分区不够。这一章把位点提交、心跳、再平衡拆开,并且把停顿时间量出来。
消费组卡住的那几秒,是 Kafka 排障里最常被误判的地方。有人看到「消费停了 10 秒」就以为是下游慢,实际是重平衡;有人看到 lag 涨了就去加消费者,实际是分区不够。这一章把位点提交、心跳、再平衡拆开,并且把停顿时间量出来。
现场
一个消费组有 3 个分区要消费。先起 1 个消费者,再起第 2 个,然后杀掉第 2 个。整个过程用带回调的客户端记录分派变化,同时一个持续生产者以 60 条/秒的速度往里写,这样就能量出「没有数据可读」的窗口到底有多长。
全链路:一次成员变更引发的重平衡
flowchart TD
subgraph M1["成员 m1 进程"]
direction TB
M1J["① JoinGroup"]
M1A["② ASSIGN p0,p1,p2"]
M1C["③ 独占 3 个分区正常消费"]
M1R["④ 收到 REVOKE(全部收回)"]
M1A2["⑤ 重新 ASSIGN p2"]
M1H["⑥ 心跳继续"]
M1A3["⑦ REVOKE p2 -> ASSIGN p0,p1,p2"]
end
subgraph M2["成员 m2 进程"]
direction TB
M2J["①' JoinGroup(新成员加入)"]
M2A["②' ASSIGN p0,p1"]
M2K["⑧ 进程被 kill"]
end
subgraph CO["组协调者(__consumer_offsets 分区 leader)"]
direction TB
CO1["② 发起一轮 rebalance<br/>range 策略切分"]
CO2["⑦ 等 session.timeout.ms 未见 m2 心跳后把它踢出"]
end
LOG["orders-02 分区 0 / 1 / 2 的日志"]
M1J --> CO1
M2J --> CO1
CO1 --> M1A
CO1 --> M2A
M1A --> M1C
LOG --> M1C
M1C -. "REVOKE->ASSIGN 间隔 32ms,成员内消息间隔峰值 279ms" .-> M1R
M1R --> M1A2
M2K -. "kill 到重新分派约 11s,这 11s 里 p0/p1 无人消费" .-> CO2
CO2 --> M1A3
M1H -. "m2 死后 m1 仍在心跳,但不拥有 p0/p1" .-> CO2
style M1J fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style M1A fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style M1C fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style M1R fill:#fff4e6,stroke:#e8590c,color:#1a1a2e
style M1A2 fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style M1H fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style M1A3 fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style M2J fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style M2A fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style M2K fill:#ffe3e3,stroke:#c92a2a,color:#1a1a2e
style CO1 fill:#fff9db,stroke:#f08c00,color:#1a1a2e
style CO2 fill:#fff9db,stroke:#f08c00,color:#1a1a2e
style LOG fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
图里有两个不同的「停顿」,混淆它们就会误判:
- 进程内可感的停顿:m1 自己的消息到达间隔峰值,实测 279 ms(第一次重平衡)到 349 ms(成员死亡后那次)。
- 死亡成员所辖分区的空窗:m2 占据 p0/p1,它被 kill 之后,这两块分区根本没人消费,直到 m1 重新分派拿到它们——约 11 秒。
精确定义
消费组(consumer group):一组通过相同 group.id 协作消费的消费者。组内每个分区最多被一个成员消费,所以成员数超过分区数时多出来的成员闲置。
分派策略(partition.assignment.strategy):
| 策略 | 做法 | 代价 |
|---|---|---|
range |
对每个 topic,按分区序均分给成员 | 成员数不能整除分区数时,前面的成员多拿一个;多个 topic 时容易倾斜 |
roundrobin |
所有 topic 的分区混起来轮询分配 | 需要所有成员订阅同一组 topic 才均衡 |
cooperative-sticky |
增量式:只挪动必须挪的分区,其余保留 | 通信轮次更多(可能要两轮),但重平衡代价小 |
本机实测用的是 range,两成员时 m1 拿到 p2、m2 拿到 p0/p1,就是按分区序号切的结果。
心跳与会话超时:heartbeat.interval.ms(默认 3000)是成员向协调者报活的间隔;session.timeout.ms(默认 45000,本机客户端设了 10000)是协调者判定成员死亡、踢出组的时限;max.poll.interval.ms(默认 300000)是「两次 poll() 之间允许的最长时间」,超了客户端会主动离开组。前两个管「进程死没死」,第三个管「进程活着但是不是在卡」。
位点提交的三种时机:
- 自动提交(
enable.auto.commit=true):后台线程每隔auto.commit.interval.ms提交一次「已经返回给上层」的位点。它不等你处理完,所以进程崩溃会丢一批。 - 同步提交(
commitSync):处理完再提交,阻塞直到得到结果。语义最清楚,代价是提交往返延迟进了消费路径。 - 异步提交(
commitAsync):提交不阻塞,失败靠回调;重平衡时要用同步提交兜一次,否则可能提交到已经被撤销的分派上(CommitFailedException)。
auto.offset.reset 决定「位点不存在时从哪读」:earliest(最早)/ latest(最新)/ none(报错)。它不是「重读」开关——重读要用 --reset-offsets 显式改位点。
证据一:lag 的构成
orders-02 用消费组 g1 消费 120 条后停下(--max-messages 120),groups --describe:
GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG
g1 orders-02 1 14 99 85
g1 orders-02 0 0 95 95
g1 orders-02 2 106 106 0
逐列读:LOG-END-OFFSET 是分区末端(95/99/106),CURRENT-OFFSET 是这组读到的位置,LAG = 末端 − 位点。三个分区的 lag(95+85+0=180)加上已读的 120,正好等于总条数 300。
这个例子还暴露了一件事:lag 在分区之间是不均衡的。这次消费几乎全从分区 2(106 条)读走了,分区 0 一条没动(lag 95)。原因是 --max-messages 120 只限制总数,不保证均摊到各分区。运维看 lag 必须看分区级,只看总数会得到「还差 180 条」这个掩盖了「某个分区完全没动」的假象。
消费完之后 lag 归零:
g1 orders-02 1 99 99 0
g1 orders-02 0 95 95 0
g1 orders-02 2 106 106 0
证据二:位点的重置语义
--reset-offsets --to-earliest --execute 把位点倒回最早:
GROUP TOPIC PARTITION NEW-OFFSET
g1 orders-02 1 0
g1 orders-02 0 0
g1 orders-02 2 0
(随后 describe)CURRENT-OFFSET 全变 0,LAG 重新变成 95/99/106
--to-latest --execute 推到最新,位点回到 95/99/106,lag 归零。这两条命令把第 00 章的结论落地了:位点是一个可以被任意改写的外部状态。它不在消息上,所以想重读就往前调,想跳过就往后调。
证据三:重平衡现场
成员 m1 单独在线时(3 个分区全给它):
g2 orders-02 0 506 536 30 rdkafka-51c6...-51a002 /0:0:0:0:0:0:0:1 rdkafka
g2 orders-02 1 485 523 38 rdkafka-51c6...-51a002 /0:0:0:0:0:0:0:1 rdkafka
g2 orders-02 2 499 539 40 rdkafka-51c6...-51a002 /0:0:0:0:0:0:0:1 rdkafka
CONSUMER-ID 一栏三行相同,说明三个分区都是同一个成员在消费。
启动第二个成员 m2 之后:
g2 orders-02 0 860 887 27 rdkafka-10ba...-b303c65 /127.0.0.1 rdkafka
g2 orders-02 1 825 856 31 rdkafka-10ba...-b303c65 /127.0.0.1 rdkafka
g2 orders-02 2 777 866 89 rdkafka-51c6...-51a002 /0:0:0:0:0:0:0:1 rdkafka
分区 0/1 换成了 m2(CONSUMER-ID 变了),分区 2 还是 m1。这就是 range 策略的结果。
客户端的回调日志把时序摊开了:
[15:05:25.684] m1 REVOKE [('orders-02', 0), ('orders-02', 1), ('orders-02', 2)]
[15:05:25.716] m1 ASSIGN [('orders-02', 2)]
REVOKE 到 ASSIGN 间隔 32 ms。整个过程中 m1 的消息到达间隔峰值从 28 ms 跳到 279 ms——这 279 ms 就是「用户在重平衡期间感知到的消费停顿」。它不是 3 秒也不是 30 秒,而是一次分派交接的开销。
再杀掉 m2:
[15:05:42.398] m1 heartbeat ok, msgs=1174, max_gap=279ms
(kill m2 发生在 15:05:42)
[15:05:53.278] m1 REVOKE [('orders-02', 2)]
[15:05:53.311] m1 ASSIGN [('orders-02', 0), ('orders-02', 1), ('orders-02', 2)]
[15:05:57.449] m1 heartbeat ok, msgs=2123, max_gap=349ms
kill 发生在 15:05:42,重新分派在 15:05:53——中间约 11 秒。这 11 秒里,m2 持有的 p0/p1 没有任何消费者,而 m1 只在消费 p2。停顿 = 会话超时检测(本机 10 s)+ 一轮重平衡。m1 自己的消息间隔峰值也升到 349 ms,但那远小于 11 秒——只盯着「进程内消息间隔」会严重低估这次故障的影响,因为受影响的是 m2 那两个分区。
证据四:重平衡的固定开销
ConsumerPerformance 单消费者读 20 万条,结果行里有一列 rebalance.time.ms:
data.consumed.in.nMsg nMsg.sec rebalance.time.ms fetch.time.ms
200180 57128.99 3268 236
3.27 秒花在加入组和拿分派上,真正 fetch 数据只用了 236 ms。消费者刚启动时这段「什么都不干只协商」的时间就是重平衡成本。消费者数量多、频繁上下线时,这段成本会反复发生。
主动破坏:让消费者卡住但不死
把消费者的处理逻辑里加一段比 max.poll.interval.ms 更长的 sleep(本机可以设小一点,比如 10 秒),观察会发生什么:协调者不会因为「进程还活着」就放过它——超过 max.poll.interval.ms 没 poll,客户端会主动发 LeaveGroup,触发一轮重平衡,把分区交给别人。这个破坏方式能让你区分「进程崩了」和「进程在卡」,两者触发的重平衡路径不一样。
生产边界
- 教学替身 vs 真实依赖:本机两个消费者是两个本地 Python 进程,组协调者是本机 broker。真实环境里成员分布在不同机器,网络延迟会放大
session.timeout.ms的检测时间;cooperative-sticky在这种场景下比range明显更容易不中断消费,但本机只测了range。 - 要盯的指标:分区级
records-lag-max(不要只看总 lag)、rebalance-rate/rebalance-latency、消费者的poll-interval与commit-latency、JoinGroup阶段的耗时。消费端「处理一条消息的耗时」也要埋点——max.poll.interval.ms超限的根因通常在这里。 - 失败策略:处理慢导致的踢出会形成「处理慢 → 被踢 → 重平衡 → 从头再来 → 更慢」的雪崩。兜底做法是「先提交再处理」之外的第三种选择:把重活丢给下游异步做,
poll循环只负责取数和入队。重平衡时的REVOKE回调里要同步提交一次位点,否则已处理的部分会重复。
动手
- 起 4 个消费者消费一个 3 分区的 topic,
--describe观察第 4 个成员的ASSIGNMENT是否为空,解释原因。 - 把
session.timeout.ms从 10000 调到 30000,重复第 2 个成员被 kill 的实验,量出空窗从约 11 秒变成多少。 - 断言题:在
range策略下,3 个成员消费 4 个分区,谁拿 2 个?把成员数改成 2,再回答一次。
自测
- 「进程内消息间隔峰值 279 ms」和「死亡成员分区空窗 11 秒」为什么差这么多?哪一个才对应业务方观察到的「数据延迟」?
- lag 总量为 0,能不能说明消费健康?本机那个 p0 完全没消费的例子说明了什么?
auto.offset.reset=earliest和--reset-offsets --to-earliest有什么区别?各自在什么场景用?- 自动提交为什么可能丢消息?同步提交和异步提交的重平衡行为差在哪?
cooperative-sticky为什么能减少重平衡停顿?它为此付出了什么代价?
↓ 下一步:05 章 · 端到端可靠性