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

图里有两个不同的「停顿」,混淆它们就会误判:

精确定义

消费组(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() 之间允许的最长时间」,超了客户端会主动离开组。前两个管「进程死没死」,第三个管「进程活着但是不是在卡」。

位点提交的三种时机:

  1. 自动提交(enable.auto.commit=true):后台线程每隔 auto.commit.interval.ms 提交一次「已经返回给上层」的位点。它不等你处理完,所以进程崩溃会丢一批。
  2. 同步提交(commitSync):处理完再提交,阻塞直到得到结果。语义最清楚,代价是提交往返延迟进了消费路径。
  3. 异步提交(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,触发一轮重平衡,把分区交给别人。这个破坏方式能让你区分「进程崩了」和「进程在卡」,两者触发的重平衡路径不一样。


生产边界

动手

  1. 起 4 个消费者消费一个 3 分区的 topic,--describe 观察第 4 个成员的 ASSIGNMENT 是否为空,解释原因。
  2. 把 session.timeout.ms 从 10000 调到 30000,重复第 2 个成员被 kill 的实验,量出空窗从约 11 秒变成多少。
  3. 断言题:在 range 策略下,3 个成员消费 4 个分区,谁拿 2 个?把成员数改成 2,再回答一次。

自测

  1. 「进程内消息间隔峰值 279 ms」和「死亡成员分区空窗 11 秒」为什么差这么多?哪一个才对应业务方观察到的「数据延迟」?
  2. lag 总量为 0,能不能说明消费健康?本机那个 p0 完全没消费的例子说明了什么?
  3. auto.offset.reset=earliest 和 --reset-offsets --to-earliest 有什么区别?各自在什么场景用?
  4. 自动提交为什么可能丢消息?同步提交和异步提交的重平衡行为差在哪?
  5. cooperative-sticky 为什么能减少重平衡停顿?它为此付出了什么代价?

↓ 下一步:05 章 · 端到端可靠性

进入 keel 阅读