KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA

02 · 副本、ISR 与 leader 选举 — keel 龙骨

acks=all 是 Kafka 里被误解最多的一行配置。很多人把它读成「绝不丢」,然后在把一个 broker 停掉之后发现生产端写不进去,报错还不是他预期的那个。这一章把 ISR 这个中间概念拆开,并且用真实报错把「什么情况报什么错」钉死。

acks=all 是 Kafka 里被误解最多的一行配置。很多人把它读成「绝不丢」,然后在把一个 broker 停掉之后发现生产端写不进去,报错还不是他预期的那个。这一章把 ISR 这个中间概念拆开,并且用真实报错把「什么情况报什么错」钉死。

现场

集群 A 是 3 个节点(a1/a2/a3),都是 combined 模式(broker 和 controller 同进程)。关键配置:

default.replication.factor=3
min.insync.replicas=2
unclean.leader.election.enable=false

一个 3 副本的 topic,生产者用 acks=all。现在停掉 broker,看生产端到底会拿到什么。我原本的预期(也是很多文档给的结论)是「停掉两个 broker 就会报 NOT_ENOUGH_REPLICAS」。实测结果不是。这一章的重点就是把「预期 vs 实测」的差讲清楚。

全链路:acks=all 的写路径与两条失败分支

flowchart TD
  P["① 生产者 acks=all<br/>orders-03 / orders-04"] --> META["② 从元数据拿到<br/>分区 leader 与 ISR 列表"]
  META --> L["③ 写到 leader 的 .log<br/>(a1/a2/a3 之一)"]
  L --> F1["④ follower a2 fetch 复制"]
  L --> F2["④ follower a3 fetch 复制"]
  F1 -. "落后超过 replica.lag.time.max.ms" .-> SHRINK["⑤ leader 判定 follower 掉队<br/>发起 ISR 收缩"]
  F2 -. "追上则保留在 ISR" .-> KEEP["⑤ 留在 ISR"]
  L --> CHK{"⑥ ISR 大小 >= min.insync.replicas ?"}
  CHK -- "否" --> E1["NotEnoughReplicasError (错误码 19)<br/>leader 直接拒绝,不写"]
  CHK -- "是" --> WAIT["⑦ 等 ISR 里每个成员都 ack"]
  WAIT --> OK["⑧ 返回成功,HW 推进"]
  WAIT -. "某个 ISR 成员已死且 ISR 变不了<br/>(controller 仲裁丢失)" .-> E2["RequestTimedOutError<br/>客户端一直等不到 ack"]
  SHRINK --> CTRL{"⑨ ISR 变更要写进<br/>controller 元数据日志"}
  CTRL -- "controller 有多数派" --> APPLY["⑩ ISR 变更生效"]
  CTRL -- "controller 少数派(3 个 voter 停了 2 个)" --> STUCK["⑩ 元数据冻结,ISR 不变<br/>→ 走 E2 虚线"]
  style P fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
  style META fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
  style L fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
  style F1 fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
  style F2 fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
  style SHRINK fill:#fff4e6,stroke:#e8590c,color:#1a1a2e
  style KEEP fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
  style CHK fill:#fff9db,stroke:#f08c00,color:#1a1a2e
  style CTRL fill:#fff9db,stroke:#f08c00,color:#1a1a2e
  style E1 fill:#ffe3e3,stroke:#c92a2a,color:#1a1a2e
  style E2 fill:#ffe3e3,stroke:#c92a2a,color:#1a1a2e
  style WAIT fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
  style OK fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
  style APPLY fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
  style STUCK fill:#ffe3e3,stroke:#c92a2a,color:#1a1a2e

两条失败分支的区别就是整章的题眼:E1(ISR 不够)是 leader 主动拒绝,客户端立刻拿到明确错误;E2(ISR 变不了)是客户端干等超时。它们看起来都是「写不进去」,但根因和处置完全不同。

精确定义

ISR(in-sync replicas) 是「跟上了 leader 的副本集合」,包含 leader 自己。它不是「所有活着的副本」——一个活着但落后太多的 follower 会被踢出 ISR。判定标准是 follower 是否在 replica.lag.time.max.ms(默认 30000 ms)内跟上了 leader 的 LEO。

谁维护 ISR:leader 负责判断 follower 掉队并发起 ISR 收缩,但这个变更要作为元数据写进 controller 的元数据日志才算生效。所以 ISR 的维护是 leader 判断 + controller 提交的两段式——记住这一点,后面失败分支就是它导致的。

ISR 增长:follower 重启后追上 leader 的进度,leader 把它加回 ISR,同样要经 controller 提交。

HW(high watermark) = 所有 ISR 成员都已复制的最高 offset。消费者只能读到 HW 之前的消息(read_uncommitted 下也是),这是「已提交」的边界。

min.insync.replicas 是「写成功要求 ISR 至少有几个成员」。它与 acks=all 配对使用:acks=all 说「等所有 ISR ack」,min.insync.replicas 说「ISR 少于这个数就别写」。单独设一个没有意义。

unclean leader election 指「从不在 ISR 里的副本里选 leader」。允许它会牺牲一致性换可用性。

证据一:ISR 收缩与 leader 切换

三副本齐全时,orders-03(min.insync.replicas=2):

Partition: 0   Leader: 1   Replicas: 1,2,3   Isr: 1,2,3
Partition: 1   Leader: 2   Replicas: 2,3,1   Isr: 2,3,1
Partition: 2   Leader: 3   Replicas: 3,1,2   Isr: 3,1,2

停掉 a3 之后:

Partition: 0   Leader: 1   Replicas: 1,2,3   Isr: 1,2
Partition: 1   Leader: 2   Replicas: 2,3,1   Isr: 2,1
Partition: 2   Leader: 1   Replicas: 3,1,2   Isr: 1,2

两个变化:a3 从三个 ISR 里消失;分区 2 的 leader 从 3 变成了 1——因为 leader a3 没了,controller 从 ISR({3,1,2})里选了个新 leader,按副本顺序选了 1。

时间上,从 kill 到 ISR 收缩完成约 10 秒(broker.session.timeout.ms 是 9 秒级)。这个数字值得记:broker 挂了之后,ISR 不会立刻缩,要等 controller 把它 fence 掉。这段窗口里 ISR 还是旧的,acks=all 会等一个已经死掉的副本。

证据二:ISR=2 时 acks=all 正常

同一个时刻,用 acks=all 往 orders-03 发 5 条:

[isr2] 首条成功: partition=1 offset=2
[isr2] 汇总: 成功 5 / 失败 0

ISR 是 2,min.insync.replicas 是 2,2 ≥ 2,所以写成功。这说明 acks=all 不是「等所有副本」,是「等 ISR 里所有副本」;ISR 缩到 2 之后,它就只等这 2 个。你以为的「3 副本强一致」在掉了一个副本后就悄悄降级成了「2 副本一致」。

证据三:停两个 broker ≠ NOT_ENOUGH_REPLICAS

把 a2、a3 都停掉,只留 a1。等了一分钟,orders-03 的 ISR 一直是:

Partition: 0   Leader: 1   Isr: 1,2
Partition: 1   Leader: 2   Isr: 2,1
Partition: 2   Leader: 1   Isr: 1,2

ISR 没有收缩到 1,分区 1 的 leader 还是已经死掉的 2。为什么?因为这是 3 个 combined 节点,controller 也是这 3 个。KRaft 的 controller 需要多数派(2/3)才能提交任何元数据变更。停了 2 个,controller 只剩 a1 一个 voter,仲裁丢失,元数据被冻结——ISR 变更提交不了,leader 也改不了。

这时查 controller 状态:

TimeoutException: Call(callName=describeMetadataQuorum ...) timed out ...
Caused by: DisconnectException: ... due to node 1 being disconnected

再用 acks=all 发一条:

[onlyA1] 首条失败: RequestTimedOutError: [Error 7] RequestTimedOutError: Request timed out after 10000 ms
[onlyA1] 失败类型: RequestTimedOutError

拿到的不是 NOT_ENOUGH_REPLICAS,是超时。 因为 leader(a1)眼里的 ISR 仍然是 {1,2}(2 个 ≥ min.isr 2),它接受写入,然后等 a2 的 ack——a2 永远不会来,也没有任何机制能把 ISR 缩到 {1},所以客户端一直等到自己的 request.timeout.ms 超时。

这条是本课最反直觉的一条:「用 3 个 combined 节点,停 2 个来演示 NOT_ENOUGH_REPLICAS」这个方案本身是错的,它先破坏的是 controller 仲裁,不是 ISR。

证据四:正确地复现 NOT_ENOUGH_REPLICAS

要让 leader 真正拒绝,得保住 controller 仲裁(只停 1 个)同时让 ISR < min.isr。做法是把 min.insync.replicas 抬到 3:

./kafka-cli.sh topics --create --topic orders-04 --partitions 3 --replication-factor 3 --config min.insync.replicas=3

然后停掉 a3(voter 剩 a1/a2,仲裁还在),ISR 在约 10 秒内从 3 缩到 2:

[15:15:25] p0 leader=3 isr=3,1,2   p1 leader=1 isr=1,2,3   p2 leader=2 isr=2,3,1
[15:15:35] p0 leader=1 isr=1,2     p1 leader=1 isr=1,2     p2 leader=2 isr=2,1

再发一条 acks=all:

[isr2min3] 首条失败: NotEnoughReplicasError: [Error 19] NotEnoughReplicasError: None
[isr2min3] 失败类型: NotEnoughReplicasError

Error 19 就是 NOT_ENOUGH_REPLICAS。注意报错正文是 None——broker 侧的 NotEnoughReplicasException 没带消息体,所以别指望从错误文案里看出 ISR 数字,要么开 debug 日志,要么直接查 topic describe。

同一个时刻对 orders-03(min.insync.replicas=2)发同样的 acks=all:

[minisr2ok] 首条成功: partition=0 offset=6
[minisr2ok] 汇总: 成功 2 / 失败 0

同集群、同时刻、同为 ISR=2,min.isr=3 报错、min.isr=2 成功。这一组对照把「ISR 大小 vs min.insync.replicas」的关系钉死了。

顺带看配置的优先级链(configs --describe):

min.insync.replicas=3 sensitive=false synonyms={
  DYNAMIC_TOPIC_CONFIG:min.insync.replicas=3,
  DYNAMIC_DEFAULT_BROKER_CONFIG:min.insync.replicas=2,
  STATIC_BROKER_CONFIG:min.insync.replicas=2,
  DEFAULT_CONFIG:min.insync.replicas=1}

topic 级覆盖 broker 级,broker 级覆盖默认值。生产上常见的事故是「topic 建的时候用了默认值」,事后想改要记得改 topic 级配置,不是改 broker 配置文件就完事。

证据五:ISR 回填与 leader 回到 preferred

把 a3 重启,观察 orders-04 的 ISR:

[15:17:40] p0 L1 isr=1,2       p1 L1 isr=1,2       p2 L2 isr=2,1
[15:17:48] p0 L1 isr=1,2,3     p1 L1 isr=1,2,3     p2 L2 isr=2,1,3
[15:18:26] p0 L3 isr=1,2,3     p1 L1 isr=1,2,3     p2 L2 isr=2,1,3

两段时间:a3 起来后约 8 秒,它追上进度回到 ISR;再过约 45 秒,分区 0 的 leader 从 1 变回 3——3 是它的 preferred leader(replicas 列表里的第一个)。这说明 leader 不会自动回到最优位置,要等 controller 做一次 leader 再平衡。

证据六:unclean 选举的约束

手动触发 unclean 选举:

./kafka-cli.sh leader --election-type UNCLEAN --topic orders-05 --partition 0
-> Valid replica already elected for partitions orders-05-0

当已经有一个合法(在 ISR 里)的 leader 时,unclean 选举是 no-op——它只在「没有任何 ISR 副本可用」时才会去选 ISR 外的副本。这和 unclean.leader.election.enable=false 的语义一致:控制器不会把 leader 交给不在 ISR 里的副本。

对照 preferred 选举(请求把一个离线副本选成 leader):

./kafka-cli.sh leader --election-type PREFERRED --topic orders-05 --partition 0
-> PreferredLeaderNotAvailableException: The preferred leader was not available.

preferred 副本离线时直接报不可用,不会硬选。

证据七:offline 分区长什么样

把 orders-06 的副本缩到只剩 a3(单副本),再停 a3,这个分区就彻底没有可用副本了:

Partition: 0   Leader: none   Replicas: 3   Isr: (空)   Elr: 3   LastKnownElr: 3

Leader: none、ISR 为空,--unavailable-partitions 也能列出它。往它写:

[offline] 首条失败: KafkaTimeoutError: Expiring 1 record(s) ... 20000 ms has passed since batch creation

注意这里客户端看到的还是超时。LEADER_NOT_AVAILABLE 是一个可重试错误码,客户端拿到之后会去刷新元数据、等待,对应用代码来说最终表现为「批过期超时」。所以「为什么 Kafka 写不进去」这个现象,可能对应三种完全不同的根因:ISR 不足、ISR 变不了、没有 leader。只看客户端报错文案是分不清的,要去查 topics --describe 和 controller 状态。

acks=all 保证什么、不保证什么

保证的:写入被当时 ISR 里的所有副本接收(不是所有副本,不是所有存活副本),且 ISR ≥ min.insync.replicas。如果这个条件不满足,写入被拒绝,不会静默降级。

不保证的:

主动破坏:把 min.insync.replicas 调高

在 orders-03 上动态把 min.insync.replicas 改成 3,然后停一个 broker,看生产端是否立刻开始报 NotEnoughReplicasError。这个练习的价值在于:你会在不改代码、不改 topic 的情况下,让一个正在跑的生产者从「正常」变成「全部失败」——这正是配置事故的真实形态。


生产边界

动手

  1. 在 orders-04(min.isr=3)停一个 broker,断言生产端报 NotEnoughReplicasError;把 min.isr 改回 2,断言同一条命令变成成功。
  2. 用 MetadataQuorumCommand ... describe --status 在停 2 个 broker 前后各跑一次,观察 TimeoutException 出现。
  3. 停一个 broker,用 --unavailable-partitions 和 --under-replicated-partitions 分别列出分区,说出两者的区别。

自测

  1. 为什么「停 2 个 broker」拿到的是超时而不是 NOT_ENOUGH_REPLICAS?把 controller 仲裁这一环补进你的解释。
  2. acks=all 在 ISR 从 3 缩到 2 之后,语义发生了什么变化?「3 副本强一致」还成立吗?
  3. ISR 收缩为什么要经 controller 提交?如果 leader 能自己决定 ISR,会有什么风险?
  4. unclean.leader.election.enable=false 让分区在什么情况下宁可不可用?如果设成 true,可用性和一致性分别换来/付出什么?
  5. Leader: none 和 Isr: 1,2(小于副本数)这两种状态,运维处置上有什么不同?

↓ 下一步:03 章 · 生产者的确认与顺序

进入 keel 阅读