Skip to content

第 9 章 副本机制与 ISR:Kafka 的「合同多人会签」

目标读者:知道 Kafka 有「副本」这件事,但说不清 ISR 到底是什么、HW 和 LEO 谁推动谁、acks=all 是不是真的不丢消息、unclean.leader.election.enable 该开还是该关。

学完你会:能画出一张分区在三 Broker 之间的复制时序图,能解释为什么 min.insync.replicas + acks=all 才是「真不丢」的姿势,能用一句话讲清 Leader Epoch 解决了什么历史 bug,遇到 ISR 抖动 / Leader 频繁切换的告警知道从哪几个配置入手。


0. 导读:副本是 Kafka 的「灵魂」,不是装饰

很多同学第一次看 Kafka 文档,对副本的理解停在:

「副本就是把数据多存几份呗,跟 MySQL 主从复制差不多。」

这是一个非常表面的认识。更准确的说法是:

Kafka 的副本机制 = Leader 写日志 + Follower 主动拉取 + 一个动态的「在追副本名单」(ISR)+ 一个对消费者可见的水位线(HW)+ 一份用于恢复一致性的版本号(Leader Epoch)

副本机制是一切高可用的基础,它直接决定了:

  • 消息是否会丢acks / min.insync.replicas / unclean.leader.election.enable 三件套)
  • 消息何时对消费者可见(HW 推进的时机)
  • Broker 宕机后集群多快恢复(ISR 切换 vs Unclean Election vs 选不出 Leader)
  • 写入吞吐和延迟(同步副本越多越慢,但越安全)

本章将用「公司多份合同会签」「值日生交接班」「快递员名单」这些类比,把抽象的概念变成你能在脑子里跑得动的画面。


1. 三种角色:Leader、Follower、Observer

1.1 Leader(主副本)

每个 分区(Partition)有且仅有一个 Leader。Leader 是这个分区上唯一处理读写请求的副本:

  • 所有 Producer 的 produce 请求只能打到 Leader;
  • 所有 Consumer 的 fetch 请求默认也只能从 Leader 读(Kafka 2.4 之后允许通过 Follower Fetching 从最近副本读,但只是延迟优化,本章先按「读写都走 Leader」讲)。

生活类比:Leader 像一个项目的「合同签字主管」——所有需要进文件柜的合同,必须先经过他签字盖章,他才把合同的最新页码(Offset)告诉所有人。

1.2 Follower(从副本)

同一个分区在其他 Broker 上的副本就是 Follower。Follower 不响应客户端读写,它只做一件事:

不停地向 Leader 发 Fetch 请求,把 Leader 已经写入的日志条目回来追平。

注意「拉(pull)」这个词:Follower 主动拉,Leader 不主动推。这是 Kafka 复制模型的关键设计——Follower 行为和普通 Consumer 一模一样,只不过它读的 Topic 是「自己分区的日志」,并且把 fetch 到的内容写到自己的本地 log 文件里。

生活类比:Follower 像合同副本管理员,每隔一段时间就跑去主管办公室抄一遍最新的合同进去。抄得快的算「跟得上」,抄得慢的会被踢出名单。

1.3 Observer(KRaft / 企业版概念)

⚠️ 澄清:标准开源 Kafka 的「数据副本」只有 Leader 和 Follower 两种。Observer 是 Confluent 商业版的多区域副本特性(Multi-Region Clusters,MRC),用来在跨 IDC 场景里放一份「不参与 ISR、不影响吞吐、只用于灾备读取」的副本。

同时,KRaft 的 __cluster_metadata Topic 里也有「Observer」概念——指的是还没成为 Voter 的 Controller(不投票,只同步),这是 Raft 协议层面的 Observer,跟数据副本是两码事。

本章后续提到的副本机制,默认指开源 Kafka 的「数据副本」(Leader / Follower)。Observer 的两种含义一并知道即可,不会在 ISR 计算里出现。

📌 与其他 MQ 的对比

  • RabbitMQ Mirrored Queue:早期采用主从镜像,每条消息要写到所有镜像才 ack,写慢且选主复杂;新版 Quorum Queue 走 Raft,每条写需要 Quorum 多数确认。
  • RocketMQ:早期 Master/Slave 异步或同步刷盘 + 异步或同步复制,Slave 只做读分担和灾备;4.5 之后引入 DLedger(Raft)支持自动选主。
  • Pulsar:存储与计算分离,分区数据在 BookKeeper 的多个 Bookie 上保存 N 份,用 Quorum 写入(Ensemble Size / Write Quorum / Ack Quorum)。
  • Kafka:「Leader + Follower 拉取 + ISR」是自成一派的方案,不要求 Quorum 多数,只要求 ISR 多数——这是 Kafka 在写入吞吐上能领先一截的关键。

2. ISR:动态的「在追副本」名单

2.1 什么是 ISR

ISR(In-Sync Replicas)= 当前与 Leader 进度同步的副本集合

注意这个集合是动态的

  • Follower 跟得上时在 ISR 里;
  • Follower 卡住了,超过一定时间就被剔出 ISR
  • 卡住的 Follower 重新追上后,再被加回 ISR

ISR 集合存放在 Controller 的内存里(KRaft 模式则在 __cluster_metadata 日志里),并通过元数据请求广播给所有 Broker。

生活类比:ISR 像一个「可信快递员名单」。今天小张请假了一上午,快递追不上派送进度,调度员就把他从今日可信名单里划掉;等他下午上班并把今日所有件都送完了,才允许他重新上名单。

2.2 Leader 是怎么判断 Follower 「跟得上」的?

这是 Kafka 历史上一个演进了好几个版本的细节。

早期版本(0.9 之前):用 replica.lag.max.messages 判断——Follower 落后 Leader 超过 N 条消息就踢出去。

🚨 缺陷:N 条这个阈值在突发流量下毫无意义。Producer 一秒钟刚批量写了 5000 条,所有 Follower 瞬间都 lag 超过 5000,整个 ISR 直接清空,集群可用性反而下降。

0.9 之后:改成 replica.lag.time.max.ms(默认 30000 ms = 30 秒)。判断标准变为:

Follower 多长时间没有 fetch 到 Leader 的最新 LEO,超过这个时间窗就剔出 ISR。

「fetch 到最新 LEO」的细节是:每次 Follower 发 Fetch 请求时带上自己的 fetch offset,如果 Leader 看到 Follower 这次请求的 offset 已经追到了 Leader 的 LEO,就更新 Follower 的「最近一次跟上时间」为当前时间。一旦超过 replica.lag.time.max.ms 都没有「追到 LEO」过一次,就踢出 ISR。

为什么时间比条数好?

  • 时间天然抗突发流量:Producer 一波猛写,所有 Follower 同时 lag,但只要它们很快就拉到了最新 LEO,时间窗内是「跟上的」,就不踢;
  • 时间能识别真正的故障:网络挂掉的 Follower 30 秒都没成功 fetch,肯定该踢了。

2.3 OSR:被踢出的副本去哪了

ISR 之外还有个名字叫 OSR(Out-of-Sync Replicas):

所有副本(AR, Assigned Replicas)
  = ISR(在追的)  +  OSR(掉队的)

AR 是「这个分区分配在哪些 Broker 上」的静态集合,由你创建 Topic 时指定,扩缩容前不会变;ISR / OSR 是 AR 在运行时的动态划分。

被踢出去的副本(OSR)依然在拉数据,只是不计入 ISR 投票。一旦它追上了 Leader 的 LEO,会自动重新加入 ISR。


3. HW 与 LEO:消费者究竟能看到哪条消息

这是最容易被新手忽视、却是面试必问的概念。先看定义:

3.1 LEO(Log End Offset)

每个副本(Leader 和每个 Follower)都有自己的 LEO,意思是:

这个副本本地日志中下一条待写入的位置(也就是当前最大 offset + 1)。

注意是「下一条待写入」,不是「最后一条已写入」。这个细节会在后面的 HW 推进里关键。

3.2 HW(High Watermark,高水位)

HW 是分区级别的概念,由 Leader 计算并维护:

HW = ISR 中所有副本的 LEO 的最小值

它的含义是:

Offset < HW 的消息,已经被 ISR 中所有副本都同步了(被认为「永久持久化」),消费者可以读到; Offset ≥ HW 的消息,还没被所有 ISR 副本同步,消费者读不到

HW 是 Kafka 给消费者的「可见窗口」。消费者 fetch 时,Leader 只会返回 [fetchOffset, HW) 之间的数据。

3.3 ASCII 图示:HW 是怎么推进的

假设一个分区有三副本(B1=Leader,B2/B3=Follower),都在 ISR 里。初始所有 LEO=0,HW=0。

步骤 ①:Producer 写入两条消息到 Leader

Leader B1   : [m0][m1]              LEO=2   HW=0
Follower B2 :                       LEO=0
Follower B3 :                       LEO=0
分区 HW = min(2, 0, 0) = 0          ← 消费者还看不到 m0, m1

步骤 ②:B2 发起 Fetch,拿到 m0, m1,写入本地

Leader B1   : [m0][m1]              LEO=2   HW=0
Follower B2 : [m0][m1]              LEO=2
Follower B3 :                       LEO=0
分区 HW = min(2, 2, 0) = 0          ← B3 还没追上,HW 仍是 0

步骤 ③:B3 也发起 Fetch,追上

Leader B1   : [m0][m1]              LEO=2   HW=0
Follower B2 : [m0][m1]              LEO=2
Follower B3 : [m0][m1]              LEO=2
此时 Leader 还**不知道** B3 的 LEO=2

注意:Leader 是通过下一次 Fetch 请求里带的 fetch.offset 才知道 Follower 当前的 LEO。所以 HW 推进永远滞后于真实写入 一个 RTT。

步骤 ④:B2 / B3 发出下一次 Fetch(fetch.offset=2)

Leader 收到所有 Follower 的 fetch.offset=2,更新各自 LEO=2
重新计算 HW = min(2, 2, 2) = 2      ← HW 推进到 2

随后 Leader 在 Fetch Response 里把新 HW=2 带回 Follower
Follower 自己也把本地 HW 更新成 2

此时消费者再 fetch 就能读到 m0, m1 了 ✓

3.4 关键结论

  1. HW 一定 ≤ LEO(不可能比已经写入的还高)。
  2. HW 推进至少需要一轮 Follower 的 Fetch(Leader 必须看到 Follower 的最新 fetch.offset)。
  3. HW 由 Leader 单点计算,再通过 FetchResponse 同步给 Follower。Follower 拿到最新 HW 之后才能更新自己本地 HW。
  4. 消费者的可见性取决于 HW,不是 LEO——这就是为什么你 acks=1 的消息可能还没被 Follower 拉到,但 consumer.poll() 在某些版本就能读到(因为 Leader 本地写入后会立刻把 LEO 往前推,但 HW 不动)。

4. Leader Epoch:合同的「版本号」

4.1 历史问题:HW 截断(HW Truncation)的 bug

Kafka 0.11 之前,副本切主后会发生一个非常隐蔽的 数据不一致 bug。我们来还原一下。

场景:分区 P 有副本 B1(Leader)、B2(Follower)。acks=all,HW=2。某一时刻:

Leader B1   : [m0][m1][m2]   LEO=3   HW=2  ← 已经写入 m2,但 HW 还没来得及推进
Follower B2 : [m0][m1][m2]   LEO=3   HW=2  ← 跟 Leader 同步过了

注意 HW=2 但 LEO=3:m2 已经写入双方本地,只是 HW 还没来得及推进到 3。

事故 1:B1 突然掉电。Controller 把 B2 选成新 Leader。

事故 2:B1 复电起来,发现自己是 Follower。

旧版本的恢复逻辑是:

Follower 启动时,会截断本地日志到自己当前的 HW(认为 HW 之上的数据可能是脏数据,要重新从新 Leader 拉)。

于是 B1 把自己的日志截断到 HW=2(删掉 m2),变成:

Old Leader B1(现在是 Follower): [m0][m1]        LEO=2   HW=2
New Leader B2                  : [m0][m1][m2]   LEO=3   HW=2

事故 3:B2 也掉电了。集群只剩 B1。

如果开启 unclean.leader.election.enable=true,Controller 会让 B1 当 Leader:

Leader B1   : [m0][m1]   LEO=2   HW=2   ← m2 永远丢了 ❌

事故还有更狠的版本:B2 没死,但 m2 之后又有新消息 m2'。B1 复电后截断,再去拉 m2',导致两边的 m2 完全不同——也就是「同一个 offset 上有两条不同的消息」,对消费者而言是幽灵消息

根因:Follower 启动时只看 HW 决定是否截断,没办法识别「这条消息究竟是不是被旧 Leader 写过」

4.2 Leader Epoch 的解决方案

Kafka 0.11 引入 Leader Epoch——每次发生 Leader 切换时,Controller 都会把这个分区的 epoch +1,并写进所有副本的 leader-epoch-checkpoint 文件。

Epoch 0: B1 当 Leader([m0, m1, m2] 都在 epoch 0 下写入)
Epoch 1: B2 当 Leader(在 epoch 1 下写入新消息)
Epoch 2: B1 重新当 Leader
...

每个 epoch 在每个副本上都对应一个起始 offset(这个 epoch 的第一条消息所在的 offset)。

leader-epoch-checkpoint 文件示意:
0  0       ← epoch 0 从 offset 0 开始
1  3       ← epoch 1 从 offset 3 开始
2  10      ← epoch 2 从 offset 10 开始

新的恢复流程

Follower 启动后,先向新 Leader 发 OffsetForLeaderEpoch 请求

「我本地最大的 epoch 是 X,X 这个 epoch 在你那里的 end offset 是多少?」

新 Leader 回答:「我这里 epoch X 的 end offset 是 Y。」

Follower 比较:

  • 如果自己的 LEO ≤ Y:保留全部数据,从 LEO 开始正常 fetch;
  • 如果自己的 LEO > Y:截断到 Y,然后从 Y 开始正常 fetch。

带回上面的事故

B1 复电时本地:[m0(epoch=0)][m1(epoch=0)][m2(epoch=0)]   LEO=3
新 Leader B2 此时:B2 接管后没写新数据,epoch=1 的 end offset = 3

B1 询问 B2:「epoch=0 的 end offset 是多少?」B2 答:「3」
B1 LEO=3 ≤ 3,不截断,保留 m2 ✓

m2 不会再丢。

生活类比:Leader Epoch 像合同的版本号。每换一任主管就启用新版本号;副本管理员复职后,先去问主管「v0 版本最后一页是多少」,再决定自己手上的合同要不要撕掉一页。

4.3 不再依赖 HW 做截断,是 Kafka 副本协议的一次根本性升级

Leader Epoch 让恢复逻辑从「相信 HW」变成「相信 epoch 链」。这是一次非常关键的升级,所有 1.0+ 版本默认都开启,你不用配置。但理解这个故事能让你在排查「同 offset 不同消息」「ISR 切换后丢数据」时,知道首先去看 leader-epoch-checkpoint 文件。


5. unclean.leader.election.enable 的代价

5.1 什么是 Unclean Leader Election

当 ISR 中所有副本都不可用(全部宕机或全部掉队)时,Controller 面临两个选择:

  • (默认)unclean.leader.election.enable=false:等,等到 ISR 中至少一个副本回来。期间这个分区完全不可写也不可读(消费者读到 HW 就停了)。
  • unclean.leader.election.enable=true:从 OSR(不在 ISR 的副本)里随便挑一个当 Leader。它可能比原 Leader 落后很多消息——但起码集群可用了。

5.2 代价权衡

场景false 的代价true 的代价
业务对可用性敏感(如告警、监控)全 ISR 宕机时分区不可用牺牲已 ack 的部分消息换可用
业务对数据敏感(如订单、账务)等待恢复,期间停写但不丢已 ack 的消息可能丢失

金句unclean.leader.election.enable=true 是「为了可用性,允许 Kafka 丢已经 ack 的消息」。

acks=all + unclean.leader.election.enable=true 是一对自相矛盾的配置——你 ack 时承诺数据安全,转头又允许从掉队副本里选主丢数据。生产中强烈建议

  • 数据敏感(金融、计费、订单):unclean.leader.election.enable=false + min.insync.replicas=2,宁可分区短暂不可用,不丢消息;
  • 可用性敏感(监控、日志):可以 true,但要清楚自己在「拿可用性换数据」。

Kafka 自从 1.0 起,默认值就是 false


6. min.insync.replicas + acks=all 才是「真不丢」

6.1 单独的 acks 远远不够

acks 是 Producer 端的参数:

acks含义风险
0发出去就当成功,不等 Broker 回复网络丢包就丢
1Leader 写入本地就回 ackLeader 写入但 Follower 还没拉,Leader 死了就丢
all(或 -1ISR 中所有副本都写入后 才回 ack「ISR 中所有」如果只剩 Leader 一个,相当于 acks=1

划重点:acks=all 的语义是「ISR 里所有副本都确认了」。如果你的 ISR 已经因为 Follower 全挂掉而只剩 Leader 一个,那么 acks=all 等价于 acks=1——Leader 一个人写入就 ack 了。这种场景下数据依然可能丢(这个唯一的 Leader 死了就什么都没了)。

6.2 min.insync.replicas 才是「兜底」

min.insync.replicas(简称 min.isr)是 Topic 或 Broker 级别的配置:

Producer 在 acks=all 时,当前 ISR 副本数 必须 ≥ min.insync.replicas 才允许写入;否则 Broker 直接返回 NotEnoughReplicasException

也就是说,min.isr=2 的语义是:「我至少要 2 个副本同时确认才允许你写入;如果 ISR 缩到只剩 1 个,宁可拒绝写入也不让你产生丢数据风险。」

6.3 推荐组合(生产最佳实践)

副本数 replication.factor=3min.insync.replicas=2,Producer 端 acks=all

  • 任何时刻最多挂 1 个 Broker,集群仍可写入;
  • 挂 2 个 Broker 时分区拒绝写入(保护数据),但不会丢已 ack 的消息;
  • 配合 unclean.leader.election.enable=false,是 Kafka 可用性 ↔ 一致性 平衡的最甜点。
副本数 = 3, min.isr = 2, acks = all
┌─────────┬──────────────────┬──────────────────────────────┐
│ 健康副本│ Producer 写入    │ 说明                         │
├─────────┼──────────────────┼──────────────────────────────┤
│ 3 / 3   │ ✅ 写入并 ack    │ 全员同步,最安全             │
│ 2 / 3   │ ✅ 写入并 ack    │ 1 个副本掉队,照常工作       │
│ 1 / 3   │ ❌ 拒绝写入      │ ISR < 2,触发 min.isr 保护   │
│ 0 / 3   │ ❌ 拒绝写入      │ 等 ISR 恢复或 unclean 选主    │
└─────────┴──────────────────┴──────────────────────────────┘

📌 常见误区:「我建了 3 副本就一定不丢消息」——错。如果 acks=1min.insync.replicas=1,仍然可能丢。三件套必须同时满足:副本数 ≥ 3、min.insync.replicas ≥ 2acks=all


7. 故障场景剧本

下面用三个剧本演练真实集群的故障行为。每个剧本前先约定环境:

集群:3 Broker(B1, B2, B3)
Topic:learn.09.orders,3 分区,每分区 3 副本,min.insync.replicas=2
Producer:acks=allenable.idempotence=true
看分区 P0 的副本布局:Leader=B1, Followers=[B2, B3]

7.1 剧本 ①:Follower 掉队后恢复

T0:稳态。ISR=[B1, B2, B3],Producer 持续写入 100 msg/s。

T1:B3 所在机器磁盘 IO 突然飙高(比如同机有大数据任务),Follower fetcher 线程跟不上:

Leader B1 LEO 在每秒推进 100,B3 的 LEO 涨幅只有 30

T2:30 秒过去(replica.lag.time.max.ms=30000),B3 始终没追上 Leader 的最新 LEO。

Controller 收到 B1 的「ISR 收缩」请求 → 把 B3 从 ISR 移除
ISR 变为 [B1, B2]

T3:Producer 继续写入。因为 ISR=[B1, B2] 还有 2 个,min.insync.replicas=2 满足,写入照常成功

Leader 计算 HW = min(LEO_B1, LEO_B2) = ...
B3 不参与 HW 计算,HW 推进不会被它拖慢

T4:B3 那台机器磁盘 IO 恢复正常,fetcher 加速追赶。

T5:B3 在某次 Fetch 中追到了 Leader 的最新 LEO。

Leader 检测到「B3 的 fetch offset == 我的 LEO」
→ 向 Controller 发起「ISR 扩张」请求
→ Controller 把 B3 加回 ISR
ISR 变回 [B1, B2, B3]

告警建议:监控 kafka.server:type=ReplicaManager,name=UnderReplicatedPartitions,长时间 > 0 就要排查;瞬时抖动可以容忍。

7.2 剧本 ②:Leader 宕机

T0:稳态。ISR=[B1, B2, B3],Leader=B1。

T1:B1 进程崩溃(kill -9 或机器宕机)。

T2:Controller 通过元数据通道(KRaft 下是 __cluster_metadata 心跳)发现 B1 失联。

Controller 把 B1 从所有相关分区的 ISR 中剔除
对 P0:从 ISR 中剩下的副本里选一个新 Leader
        通常按 AR 列表的「Preferred Replica」顺序
        这里 AR=[B1, B2, B3],B1 没了 → B2 当新 Leader

T3:Controller 把新元数据广播给所有 Broker。

新 ISR = [B2, B3]
新 Leader = B2,Epoch +1

T4:客户端缓存的元数据失效。Producer 下一次写 P0 时收到 NOT_LEADER_FOR_PARTITION 错误,自动刷新元数据,重连到 B2。

默认情况下整个切换过程在几秒内完成(取决于 controller 处理速度 + 心跳超时)
KRaft 模式下通常 < 1s(基于 Raft 直接拿到事件)
ZK 模式下可能需要 5~30s(依赖 ZK Watcher 触发)

T5:B1 复电起来,发现 epoch 已经 +1,自己变成 Follower,向 B2 发起 OffsetForLeaderEpoch,拿到截断点(如果有),然后正常 fetch。最终 B1 重新加入 ISR。

Preferred Leader Election:注意 B2 现在是 Leader。原本 AR 第一个是 B1(Preferred Leader)。如果你想让 Leader 回到「均衡」状态(避免 Leader 都堆在某些 Broker 上),可以触发:

bash
kafka-leader-election.sh \
  --bootstrap-server localhost:9092 \
  --election-type preferred \
  --topic learn.09.orders \
  --partition 0

或在 Broker 配置开启 auto.leader.rebalance.enable=true,Controller 会定期自动把 Leader 切回 Preferred。

7.3 剧本 ③:网络分区(脑裂)

T0:稳态。三 Broker 互联。

T1:B1 和 B2 之间的网络断开,但 Controller(假设在 B3)能联系到它俩。

B1(仍是 P0 的 Leader)发现 B2 fetch 卡住
→ 30s 后向 Controller 申请把 B2 踢出 ISR
ISR 变成 [B1, B3]

Kafka 不会真的「双 Leader 脑裂」,因为:

  1. Controller 是单点(KRaft 下是 Active Controller,ZK 下争 /controller 节点),所有 Leader/ISR 变更必须经过它。
  2. Broker 必须与 Controller 通信才能确认自己是 Leader;如果 Broker 与 Controller 失联(fenced),它会主动 step down,不再接受写入。
  3. 即使发生「老 Leader 还以为自己是 Leader」的瞬间,新 Leader 也会带着更高的 Leader Epoch;老 Leader 收到任何带新 epoch 的请求会立刻 step down,老 Leader 写入的孤立消息被新副本通过 epoch 截断。

真正的脑裂风险Controller 层面——ZK 时代如果出现 ZK 抖动,可能两个 Broker 都自认为是 Controller。Kafka 用 controller_epoch 来打标,所有 Broker 只听 epoch 大的那个 Controller,老的 Controller 发的指令会被忽略。

KRaft 时代这个问题被根本上消除:Controller 选举走 Raft,最多只有一个 Leader Quorum 成员,其它都自动降级为 Follower。


8. Replica Fetcher 线程模型

8.1 线程模型概览

每个 Broker 上都有若干 ReplicaFetcherThread 负责从其它 Broker 拉数据。每个 Follower → 每个 Source Broker 之间共享一个 fetcher 线程。

线程数由 num.replica.fetchers 控制(默认 1):

Broker B2 上的 fetcher 线程:
  fetcher-thread-0 → 从 B1 拉
  fetcher-thread-0 → 从 B3 拉   ← 同一个线程串行拉两个 source

如果你的集群规模大(几十个 Broker)、单 Broker 上有几百个 Follower 分区,num.replica.fetchers=1 会让单线程被打满。建议生产num.replica.fetchers=4~8,这样每个 Source Broker 由不同线程负责,能并行拉取。

8.2 Fetcher 的工作循环

loop:
  把所有「在跟这个 Source Broker 同步的分区」打包成一个 Fetch 请求
    包含每个分区的 (topic, partition, fetch_offset, max_bytes)
  发送到 Source Broker
  收到响应:每个分区的 [records...] + 新的 HW
  for each partition in response:
    把 records 追加到本地 log
    更新本地 LEO
    更新本地 HW = min(本地 LEO, 响应里的 HW)
  等待 replica.fetch.wait.max.ms 或者 fetch.min.bytes 满足,再发下一轮

replica.fetch.max.bytes(默认 1MB)控制每个分区每次最多拉多少。如果你有大消息(单条 > 1MB),这个值要调大,否则 Follower 永远拉不到那条消息,被 ISR 踢出去。

8.3 三类阈值要联动配置

Producer 端Broker 端Fetcher 端
max.request.sizemessage.max.bytesreplica.fetch.max.bytes
默认 1MB默认 1MB默认 1MB

踩坑:曾经有同学只把 message.max.bytes 调大到 10MB,结果发了 5MB 的消息,Producer / Broker 写入都成功,但 Follower fetch 一直失败(fetcher 默认只拉 1MB),Follower 被 ISR 踢掉,然后 Producer 因为 min.isr 拒绝写入,整个分区不可用。调大消息上限三处都要改


9. 实操:用命令行观察 ISR

9.1 创建 Topic 并查看副本分布

bash
kafka-topics.sh --bootstrap-server localhost:9092 \
  --create --topic learn.09.orders \
  --partitions 3 --replication-factor 3 \
  --config min.insync.replicas=2

kafka-topics.sh --bootstrap-server localhost:9092 \
  --describe --topic learn.09.orders

输出(节选):

Topic: learn.09.orders  TopicId: ...  PartitionCount: 3  ReplicationFactor: 3
  Configs: min.insync.replicas=2
  Topic: learn.09.orders  Partition: 0  Leader: 1  Replicas: 1,2,3  Isr: 1,2,3
  Topic: learn.09.orders  Partition: 1  Leader: 2  Replicas: 2,3,1  Isr: 2,3,1
  Topic: learn.09.orders  Partition: 2  Leader: 3  Replicas: 3,1,2  Isr: 3,1,2

读法:

  • Replicas: 1,2,3 表示 AR(分配副本)依次是 B1, B2, B3,第一个就是 Preferred Leader
  • Isr: 1,2,3 表示当前 ISR;
  • Leader: 1 当前 Leader。

9.2 模拟一个 Broker 宕机后再看

bash
docker compose stop kafka1

kafka-topics.sh --bootstrap-server localhost:9094 \
  --describe --topic learn.09.orders
Topic: learn.09.orders
  Partition: 0  Leader: 2  Replicas: 1,2,3  Isr: 2,3   ← B1 被踢,B2 当 Leader
  Partition: 1  Leader: 2  Replicas: 2,3,1  Isr: 2,3
  Partition: 2  Leader: 3  Replicas: 3,1,2  Isr: 3,2

把 B1 起回来后稍等几秒,ISR 会自动恢复成 2,3,1 / 1,3,2 / 3,2,1,但 Leader 不会自动切回——除非开了 auto.leader.rebalance.enable=true 或手动触发 preferred election。

9.3 看 Under-Replicated Partitions

bash
kafka-topics.sh --bootstrap-server localhost:9094 \
  --describe --under-replicated-partitions

只要还有 ISR 不全的分区就会列出来。这是一个非常关键的运维指标,应当接入告警。

9.4 看 At-Min-Isr Partitions(已经踩到 min.isr 红线的分区)

bash
kafka-topics.sh --bootstrap-server localhost:9094 \
  --describe --at-min-isr-partitions

这些分区再挂一个副本就要拒绝写入了,必须立刻处理。

9.5 看 Unavailable Partitions(已经没有 Leader 的分区)

bash
kafka-topics.sh --bootstrap-server localhost:9094 \
  --describe --unavailable-partitions

如果有,分区已经没法读写,必须立刻把任意一个副本拉起来。


10. 配套代码与演示页面

本章配套:

  • 09_replication/code/observe_isr.py:用 AdminClient 实时打印每个分区的 Leader / ISR / Replicas(每秒刷新一次)。
  • 09_replication/code/kill_leader_demo.py:脚本化地停 Broker、观察 Leader 切换、再起回来,演示 ISR 收缩和 Leader 切换全过程。
  • 09_replication/code/min_isr_demo.py:演示在 ISR < min.insync.replicas 时,acks=all 的 Producer 收到 NotEnoughReplicas 异常。
  • 09_replication/demo.html:浏览器里可视化三 Broker 集群,可点击让任一 Broker 宕机/恢复,实时显示每个分区的 Leader/ISR/OSR、LEO 与 HW,模拟 Follower 卡顿被踢出 ISR、Unclean Leader Election 等场景。

11. 调优清单(速查)

配置项默认推荐作用
replication.factor1(建 Topic 时指定)3副本数
min.insync.replicas12(与 RF=3 配套)acks=all 写入门槛
unclean.leader.election.enablefalsefalse(金融场景必须)是否允许从 OSR 选 Leader
replica.lag.time.max.ms30000视集群网络/磁盘抖动调多久不追上就踢出 ISR
num.replica.fetchers14~8(多 Broker 时)Follower 拉取并发线程数
replica.fetch.max.bytes1048576message.max.bytes 联动每分区每次拉取上限
auto.leader.rebalance.enabletruetrue自动把 Leader 还给 Preferred
leader.imbalance.check.interval.seconds300300自动均衡的检查频率
leader.imbalance.per.broker.percentage1010触发均衡的不均衡百分比

Producer 端:

配置项推荐说明
acksall配合 min.isr 才有意义
enable.idempotencetrue重试不会产生重复
max.in.flight.requests.per.connection≤ 5(idempotence 限制)顺序保障
retriesInteger.MAX_VALUE配合 delivery.timeout.ms
delivery.timeout.ms120000总超时(包含所有重试)

12. 与其他 MQ 的副本机制对比

┌──────────────────┬───────────────────────────┬──────────────────────────────────┐
│ 系统              │ 复制模型                  │ 写入确认                         │
├──────────────────┼───────────────────────────┼──────────────────────────────────┤
│ Kafka             │ Leader-Follower(拉模式) │ acks + ISR + min.isr             │
│ RabbitMQ Quorum   │ Raft                       │ 多数派 commit                     │
│ RocketMQ DLedger  │ Raft                       │ 多数派 commit                     │
│ Pulsar            │ BookKeeper Quorum 写入    │ Ack Quorum 满足即返回             │
│ MySQL 半同步      │ Master 推送 binlog        │ 至少一个 Slave ack                │
│ TiDB / CockroachDB│ Raft(Region 级)          │ 多数派 commit                     │
└──────────────────┴───────────────────────────┴──────────────────────────────────┘

关键差异:

  • Kafka 是 「ISR 多数」而不是「副本总数多数」。3 副本 + min.isr=2 在 1 个副本掉线时仍可写,等价于 Quorum;但 Kafka 比 Raft 更灵活——Follower 慢了不会拖慢写入,只是被踢出 ISR。
  • Kafka 写入 latency 取决于「ISR 中最慢的副本」,但 Raft Quorum 系统的 latency 取决于「Quorum 中第 N/2+1 快的副本」。在副本数较多时 Raft 更稳定,副本数少时 Kafka 更快。

13. 小结:副本机制的「四不」原则

把本章压缩成 4 句话送给你:

  1. 不要把 acks=1 当成「至少一次」——Leader 还没把数据同步到 Follower 就死掉,消息就没了。
  2. 不要在 min.insync.replicas=1 的情况下声称「3 副本不丢消息」——一个副本的同步等于没同步。
  3. 不要轻易开 unclean.leader.election.enable=true——开了相当于「丢已 ack 消息换可用性」。
  4. 不要忽视 Leader Epoch——它是 Kafka 副本协议在 0.11 后真正可靠的根本,理解它能帮你排查很多「同 offset 不同消息」的诡异问题。

14. 面试高频题

Q1:acks=all 一定不丢消息吗?为什么?

考察点:对 ack 语义和 ISR 的理解。

标准答案

  • acks=all 的语义是「ISR 中所有副本都确认后才返回成功」,不是「所有副本」。
  • 如果 ISR 已经因为 Follower 全挂掉而只剩 Leader 一个acks=all 等价于 acks=1,Leader 一个人写入就 ack——这种情况下唯一的 Leader 一旦死掉,消息就丢了。
  • 真正的「不丢」需要三件套:replication.factor ≥ 3 + min.insync.replicas ≥ 2 + acks=all,并且 unclean.leader.election.enable=false

加分项:能补充说幂等 Producer + enable.idempotence=true 解决的是「重试导致的重复」,不是「丢消息」;事务 Producer 解决的是「跨分区原子写入」,也不是「不丢」。三者层次不同。

Q2:HW 和 LEO 有什么区别?HW 是怎么推进的?

考察点:副本同步底层机制。

标准答案

  • LEO(Log End Offset):每个副本的「下一条待写入位置」,每次写入或 fetch 之后都会推进。
  • HW(High Watermark):分区级别的水位线,由 Leader 计算 = min(ISR 中所有副本的 LEO),消费者只能读到 HW 之前的数据。
  • HW 推进流程
    1. Producer 写入 Leader,Leader 的 LEO 推进;
    2. Follower 发起 Fetch 请求(请求里带上自己当前的 fetch.offset),Leader 把数据返给 Follower,Follower 写入本地、LEO 推进;
    3. Follower 下一次 Fetch 时带上新的 fetch.offset,Leader 据此更新「Follower 的 LEO」记录;
    4. Leader 重新计算 HW = min(所有 ISR 副本的 LEO),HW 推进;
    5. Leader 在下一个 Fetch Response 中把新 HW 同步给 Follower。
  • 也就是说HW 推进永远比真正写入慢一个 RTT

加分项:能指出 Leader Epoch 替代了「靠 HW 截断」的旧逻辑,避免了历史上的数据不一致 bug。

Q3:Leader Epoch 是为了解决什么问题?

考察点:副本协议历史演进、对 0.11 之前 bug 的认知深度。

标准答案

  • 旧版本(0.11 前)Follower 启动时按本地 HW 截断日志,但 HW 推进有延迟,可能出现「一条消息已经写入双方本地,但 HW 还没推到那里」的状态。
  • 这种状态下如果 Leader 切换:旧 Leader 复活后会按 HW 截断丢消息;甚至可能出现「同 offset 不同消息」的幽灵数据。
  • Leader Epoch 给每次主切换分配一个递增 ID,每个副本本地保存 (epoch, start_offset) 链。Follower 启动时不再按 HW 截断,而是向新 Leader 请求 OffsetForLeaderEpoch,得到「我这个 epoch 在你那里的 end offset」,再决定是否截断。
  • 从此恢复逻辑由「相信 HW」变成「相信 epoch 链」,根本上避免了数据不一致。

加分项:能讲出 leader-epoch-checkpoint 文件存在每个分区目录下,并能描述它的格式(epoch start_offset 多行)。

Q4:什么时候会发生「ISR 收缩」?怎么排查 ISR 频繁抖动?

考察点:实际运维与排障能力。

标准答案

  • ISR 收缩的本质:某个 Follower 在 replica.lag.time.max.ms(默认 30s)内没有追到 Leader 的最新 LEO。
  • 常见原因:
    1. 网络抖动 / 跨机房延迟过大——fetch 请求来回慢;
    2. 磁盘 IO 打满——Follower 写本地 log 慢;
    3. GC 长时间停顿——fetcher 线程被卡住;
    4. num.replica.fetchers 太小——一个线程跨多个 Broker 串行 fetch,被慢的 Source 拖累;
    5. 大消息超过 replica.fetch.max.bytes——Follower 永远拉不动;
    6. Leader 被打满——Producer 写入压力过大,Leader 自顾不暇,fetch 请求处理慢。
  • 排查工具:JMX UnderReplicatedPartitionsserver.log 中的「Shrinking ISR」/「Expanding ISR」日志、kafka-topics.sh --describe --under-replicated-partitions

加分项:能补充「配合 kafka-broker-api-versions.sh 看 Broker 版本一致性、用 iostat / iftop 看 IO/网络」。

Q5:unclean.leader.election.enable=true 有什么风险?什么场景该开?

考察点:可用性与一致性的权衡。

标准答案

  • 风险:从 OSR(不在 ISR 的副本)里选 Leader 意味着这个新 Leader 可能比原 Leader 落后很多消息,已经 ack 给 Producer 的消息在新 Leader 上不存在,就永久丢失了
  • 场景
    • 可以开:日志聚合、监控埋点、IoT 设备状态——这类业务可用性 > 一致性,丢一些数据可以接受,不能容忍分区长时间不可用;
    • 不能开:订单、计费、账务、支付——这类业务一致性 > 可用性,宁可分区暂时不可写也不能丢消息。
  • 默认值:1.0+ 是 false。生产中绝大多数业务保持默认即可。

加分项:能补充「即使开了 true,配合 min.insync.replicas=2 也能尽量减少丢的概率——两个副本同时挂的概率比一个副本挂小得多」。

Q6:Kafka 为什么不像 Raft 一样要求「多数派写入」就 ack?

考察点:副本协议设计哲学。

标准答案

  • Raft 要求 N/2+1 副本写入才 commit,意味着 5 副本至少要 3 个 ack——这是一种「写多读少」的优化。
  • Kafka 选择 ISR 多数(实际是「ISR 中所有」),理由:
    1. 吞吐优先:Kafka 设计目标是高吞吐日志,要求「全员同步」可以让 Leader 快速推进 HW,配合批量 + 压缩压榨吞吐;
    2. 故障识别精细:用 replica.lag.time.max.ms 时间窗自动识别掉队,掉队的 Follower 不进 ISR,就不会拖累 commit;
    3. min.isr 控制底线:用户可调,本质上模拟了 Quorum 多数,但比 Quorum 灵活;
  • Trade-off:Kafka 写入 latency 取决于「ISR 中最慢副本」,Raft 取决于「Quorum 中第 N/2+1 快的副本」。副本多时 Raft 更稳定,副本少(3 副本)时 Kafka 更快、更灵活。

加分项:能补充「Kafka 的 KRaft 内部用的是真正的 Raft(仅用于元数据),数据副本仍然是 ISR 模型——这是两个独立协议」。

Q7:一个分区有 3 副本,Leader 写入后 Follower 拉取前 Leader 死了,会发生什么?

考察点:Producer ack 与 ISR 切换的边界。

标准答案:取决于 Producer 的 acks 配置:

  • acks=0:Producer 根本没等 ack,消息可能在 socket 里就丢了,更别说 Leader 死了。
  • acks=1:Leader 本地写入成功立刻回 ack;如果此时 Follower 还没拉取,Leader 死了,新 Leader 没有这条消息,消息丢失,但 Producer 已经认为成功。
  • acks=all:Leader 必须等 ISR 中所有副本都拉到这条消息才 ack;如果还没收到所有 ack 之前 Leader 死了:
    • Producer 收到错误(如 NotEnoughReplicas 或网络错误),会重试
    • 新 Leader 起来后没有这条消息,重试会再次写入;
    • 如果开了 enable.idempotence=true,Broker 通过 PID + 序列号去重,不会重复
    • 否则可能出现「同一条消息被重试多次写入」的重复。

加分项:能讲出「Leader 在写本地之前甚至还没写盘(PageCache 未刷盘),如果同时整机断电,已 ack 的消息也可能丢」——这是为什么生产环境推荐配合 flush.messages / flush.ms 或者干脆依赖副本数(多副本即多机的「逻辑刷盘」)。

Q8:Follower Fetching(KIP-392)是什么?跟读写分离一样吗?

考察点:对 Kafka 2.4+ 新特性的了解。

标准答案

  • KIP-392 允许 Consumer 通过 client.rack 配置从最近的副本(而不是必须 Leader)读取,主要解决跨 IDC 场景的高带宽费和延迟。
  • 不是真正的「读写分离」:写入仍然只能到 Leader;只是 Consumer 读取可以到 Follower。
  • 实现机制:Consumer 在 fetch 时携带自己的 rack ID,Broker 端根据 replica.selector.class(默认 RackAwareReplicaSelector)选一个同 rack 的副本返回。
  • 注意一致性:Follower 上的 HW 是滞后的,从 Follower 读会有更高的延迟(看到的最新消息可能比 Leader 慢一两个 fetch RTT)。需要业务能容忍。

加分项:能区分「Follower Fetching」和「跨 Region 异步复制」的不同(前者还是 ISR 内同步副本,只是读的位置不同;后者通常借助 MirrorMaker 2.0 在远端集群做异步复制)。


下一章我们将深入 Kafka 的「大脑」——Controller 与 KRaft:从 ZK 时代的争抢式选举,到 Raft 化的 __cluster_metadata 元数据流,再到 4.0 完全告别 ZK 的影响与迁移路径。

🎬 可视化演示

演示加载缓慢或样式异常?点此在新标签页打开 ↗

💻 示例代码

python
#!/usr/bin/env python3
"""
kill_leader_demo.py — 自动化「停 Broker → 观察 Leader 切换 → 起回来」演示

它会做这些事:
    1) 用 AdminClient 创建一个 3 副本 Topic(如果不存在);
    2) 起一个后台 Producer,持续以 acks=all 写入;
    3) 拉取分区 P0 的 Leader Broker;
    4) 用 docker compose stop kafkaX 把 Leader 那台 Broker 干掉;
    5) 实时打印 ISR / Leader 变化(直到 Leader 切换完成);
    6) sleep 几秒再 docker compose start 起回来;
    7) 等 ISR 重新扩张回去,结束。

前置:
    * 已经用 `docker compose up -d` 起好了 3 broker 集群;
    * 容器名为 kafka1 / kafka2 / kafka3,对外端口 9092 / 9094 / 9096;
    * 当前用户能执行 `docker compose`(在 docker-compose.yml 同目录)。

用法:
    python kill_leader_demo.py
    python kill_leader_demo.py --topic learn.09.killdemo --partition 0 \
        --compose-dir /data/workspace/learnNote/kafka

⚠️ 仅用于学习环境,不要在生产集群运行!
"""

from __future__ import annotations

import argparse
import os
import subprocess
import threading
import time
from datetime import datetime

from confluent_kafka import Producer
from confluent_kafka.admin import AdminClient, NewTopic

# 容器名 → 对外端口(与 09 章 docker-compose.yml 保持一致)
BROKER_ID_TO_CONTAINER = {1: "kafka1", 2: "kafka2", 3: "kafka3"}


def parse_args():
    p = argparse.ArgumentParser()
    p.add_argument("--bootstrap", default="localhost:9092,localhost:9094,localhost:9096")
    p.add_argument("--topic", default="learn.09.killdemo")
    p.add_argument("--partition", type=int, default=0)
    p.add_argument("--compose-dir", default="/data/workspace/learnNote/kafka",
                   help="docker-compose.yml 所在目录")
    p.add_argument("--down-secs", type=int, default=15,
                   help="Leader 容器停掉后等待多少秒再起回来")
    return p.parse_args()


def ensure_topic(admin: AdminClient, topic: str, partitions: int = 3, rf: int = 3):
    md = admin.list_topics(timeout=10)
    if topic in md.topics and md.topics[topic].error is None:
        print(f"[init] topic {topic} 已存在")
        return
    print(f"[init] 创建 topic {topic} (partitions={partitions}, rf={rf})")
    nt = NewTopic(topic, num_partitions=partitions, replication_factor=rf,
                  config={"min.insync.replicas": "2"})
    fs = admin.create_topics([nt])
    for t, f in fs.items():
        try:
            f.result()
            print(f"[init] topic {t} 创建成功")
        except Exception as e:
            print(f"[init] topic {t} 创建失败: {e}")
            raise


def get_partition_meta(admin: AdminClient, topic: str, pid: int):
    md = admin.list_topics(topic=topic, timeout=10)
    t = md.topics[topic]
    p = t.partitions[pid]
    return {"leader": p.leader, "replicas": list(p.replicas), "isr": list(p.isrs)}


def docker_compose(cmd: list[str], compose_dir: str):
    full = ["docker", "compose"] + cmd
    print(f"[shell] cd {compose_dir} && {' '.join(full)}")
    return subprocess.run(full, cwd=compose_dir, check=False,
                          stdout=subprocess.PIPE, stderr=subprocess.STDOUT)


class BackgroundProducer(threading.Thread):
    """后台持续 produce,统计成功 / 失败次数"""

    def __init__(self, bootstrap, topic, partition):
        super().__init__(daemon=True)
        self.topic = topic
        self.partition = partition
        self.stop_flag = threading.Event()
        self.ok = 0
        self.fail = 0
        self.last_err = None
        self.p = Producer({
            "bootstrap.servers": bootstrap,
            "acks": "all",
            "enable.idempotence": True,
            "delivery.timeout.ms": 30000,
            "retries": 1000000,
            "linger.ms": 5,
        })

    def _on_delivery(self, err, msg):
        if err is not None:
            self.fail += 1
            self.last_err = str(err)
        else:
            self.ok += 1

    def run(self):
        i = 0
        while not self.stop_flag.is_set():
            payload = f"msg-{i}-{datetime.now().isoformat()}".encode()
            try:
                self.p.produce(self.topic, value=payload, partition=self.partition,
                               on_delivery=self._on_delivery)
                self.p.poll(0)
            except BufferError:
                self.p.poll(0.1)
            i += 1
            time.sleep(0.05)
        self.p.flush(10)


def watch_until(admin, topic, pid, predicate, label, timeout=60):
    """周期查询元数据,直到 predicate(meta) 为 True 或超时。"""
    print(f"[watch] 等待:{label} (超时 {timeout}s)")
    t0 = time.time()
    last = None
    while time.time() - t0 < timeout:
        meta = get_partition_meta(admin, topic, pid)
        snap = (meta["leader"], tuple(meta["isr"]))
        if snap != last:
            print(f"[watch] {datetime.now().strftime('%H:%M:%S')} "
                  f"Leader={meta['leader']} ISR={meta['isr']} Replicas={meta['replicas']}")
            last = snap
        if predicate(meta):
            print(f"[watch] ✓ 满足条件:{label}")
            return meta
        time.sleep(1)
    raise TimeoutError(f"{label}{timeout}s 内没有发生")


def main():
    args = parse_args()
    admin = AdminClient({"bootstrap.servers": args.bootstrap})

    ensure_topic(admin, args.topic, partitions=3, rf=3)

    meta = get_partition_meta(admin, args.topic, args.partition)
    print(f"[init] {args.topic}-P{args.partition} "
          f"Leader={meta['leader']}  ISR={meta['isr']}  Replicas={meta['replicas']}")
    if len(meta["isr"]) < 3:
        print("[warn] ISR 不足 3,建议先等到 ISR=Replicas 再演示")

    leader_id = meta["leader"]
    target = BROKER_ID_TO_CONTAINER.get(leader_id)
    if target is None:
        raise SystemExit(f"未知 Leader broker_id={leader_id},请把它加入 BROKER_ID_TO_CONTAINER")

    bg = BackgroundProducer(args.bootstrap, args.topic, args.partition)
    bg.start()
    print(f"[producer] 后台开始 produce(acks=all, idempotent=true)")
    time.sleep(2)
    print(f"[producer] 已成功 {bg.ok} / 失败 {bg.fail}")

    print(f"\n=== STEP 1: 停掉 Leader 容器 {target} ===")
    docker_compose(["stop", target], args.compose_dir)

    new_meta = watch_until(admin, args.topic, args.partition,
                           lambda m: m["leader"] != leader_id and m["leader"] >= 0,
                           label=f"Leader 从 {leader_id} 切换到其它 Broker",
                           timeout=60)

    print(f"[result] 新 Leader = {new_meta['leader']}, ISR={new_meta['isr']}")
    print(f"[producer] 切换期间累计成功 {bg.ok} / 失败 {bg.fail}, 最后错误: {bg.last_err}")

    print(f"\n=== STEP 2: 等 {args.down_secs}s 再启动回来 ===")
    time.sleep(args.down_secs)

    docker_compose(["start", target], args.compose_dir)

    watch_until(admin, args.topic, args.partition,
                lambda m: leader_id in m["isr"],
                label=f"Broker {leader_id} 重新加入 ISR",
                timeout=120)

    print(f"\n=== STEP 3: 触发 Preferred Leader Election,把 Leader 还回去 ===")
    # 简单做法:通过 kafka-leader-election.sh in container 触发
    docker_compose(
        ["exec", "-T", "kafka2",
         "/opt/kafka/bin/kafka-leader-election.sh",
         "--bootstrap-server", "kafka1:19092",
         "--election-type", "preferred",
         "--topic", args.topic, "--partition", str(args.partition)],
        args.compose_dir,
    )

    watch_until(admin, args.topic, args.partition,
                lambda m: m["leader"] == leader_id,
                label=f"Leader 切回 {leader_id}",
                timeout=60)

    print(f"\n=== DONE ===")
    print(f"[summary] 全程 produce 成功 {bg.ok} 条,失败 {bg.fail} 条")
    print("(idempotent + acks=all + min.isr=2 + 3 副本,理论上不会丢;"
          "失败的那部分是 Leader 切换瞬间的瞬时报错,会被 Producer 重试覆盖)")
    bg.stop_flag.set()
    bg.join(timeout=15)


if __name__ == "__main__":
    main()
python
#!/usr/bin/env python3
"""
min_isr_demo.py — 演示 min.insync.replicas 起作用的全过程

剧本:
    1) 创建 Topic:3 副本,min.insync.replicas=2;
    2) Producer 配置 acks=all;
    3) 正常情况下 ISR=3,写入成功;
    4) 用 docker compose stop 停掉两个 Broker,让 ISR 缩到 1;
    5) 此时 Producer 写入会拿到 NotEnoughReplicas / NotEnoughReplicasAfterAppend 错误;
    6) 把 Broker 起回来,ISR 恢复到 ≥ 2,写入恢复成功;
    7) 输出全过程的 ack / 错误统计。

用法:
    python min_isr_demo.py
    python min_isr_demo.py --bootstrap localhost:9092 --topic learn.09.minisr
"""

from __future__ import annotations

import argparse
import subprocess
import time
from datetime import datetime

from confluent_kafka import KafkaException, Producer
from confluent_kafka.admin import AdminClient, NewTopic, ConfigResource, ResourceType


def parse_args():
    p = argparse.ArgumentParser()
    p.add_argument("--bootstrap", default="localhost:9092,localhost:9094,localhost:9096")
    p.add_argument("--topic", default="learn.09.minisr")
    p.add_argument("--compose-dir", default="/data/workspace/learnNote/kafka")
    return p.parse_args()


def ensure_topic(admin: AdminClient, topic: str):
    md = admin.list_topics(timeout=10)
    if topic in md.topics and md.topics[topic].error is None:
        print(f"[init] topic {topic} 已存在")
        return
    nt = NewTopic(topic, num_partitions=1, replication_factor=3,
                  config={"min.insync.replicas": "2"})
    fs = admin.create_topics([nt])
    for t, f in fs.items():
        f.result()
        print(f"[init] 创建 topic {t} (1 part, RF=3, min.isr=2)")


def show_isr(admin, topic):
    md = admin.list_topics(topic=topic, timeout=10)
    p = md.topics[topic].partitions[0]
    print(f"[meta] {topic}-P0  Leader={p.leader}  ISR={list(p.isrs)}  "
          f"Replicas={list(p.replicas)}")
    return p


def docker_compose(cmd, cwd):
    full = ["docker", "compose"] + cmd
    print(f"[shell] cd {cwd} && {' '.join(full)}")
    return subprocess.run(full, cwd=cwd, check=False,
                          stdout=subprocess.PIPE, stderr=subprocess.STDOUT)


class Counter:
    def __init__(self):
        self.ok = 0
        self.fail = 0
        self.errors = {}

    def on_delivery(self, err, msg):
        if err is None:
            self.ok += 1
        else:
            self.fail += 1
            key = err.name() if hasattr(err, "name") else str(err)
            self.errors[key] = self.errors.get(key, 0) + 1

    def summary(self, tag):
        print(f"[stat:{tag}] 成功={self.ok}  失败={self.fail}  错误明细={self.errors}")


def burst(producer: Producer, topic, n, c: Counter):
    for i in range(n):
        try:
            producer.produce(topic, value=f"{datetime.now().isoformat()}-{i}".encode(),
                             on_delivery=c.on_delivery)
            producer.poll(0)
        except BufferError:
            producer.poll(0.5)
        except KafkaException as e:
            c.fail += 1
            c.errors[str(e)] = c.errors.get(str(e), 0) + 1
    producer.flush(15)


def main():
    args = parse_args()
    admin = AdminClient({"bootstrap.servers": args.bootstrap})

    ensure_topic(admin, args.topic)
    p = show_isr(admin, args.topic)

    producer = Producer({
        "bootstrap.servers": args.bootstrap,
        "acks": "all",
        "enable.idempotence": True,
        "delivery.timeout.ms": 15000,
        "retries": 3,                       # 故意调小,观察失败更明显
        "request.timeout.ms": 5000,
    })

    print("\n=== 阶段 1:稳态写入(ISR=3,min.isr=2 满足) ===")
    c1 = Counter()
    burst(producer, args.topic, 50, c1)
    c1.summary("steady")
    show_isr(admin, args.topic)

    # 停 2 个 follower,让 ISR 缩到 1(只剩 Leader)
    leader_id = p.leader
    others = [b for b in [1, 2, 3] if b != leader_id]
    container_map = {1: "kafka1", 2: "kafka2", 3: "kafka3"}
    targets = [container_map[i] for i in others]

    print(f"\n=== 阶段 2:停掉两个 Follower 容器 {targets},让 ISR 缩到 1 ===")
    for c in targets:
        docker_compose(["stop", c], args.compose_dir)

    # 等 Controller 把它们踢出 ISR(默认 30s)
    print("[wait] 等待 35s 让 ISR 收缩到 1…")
    for i in range(35):
        time.sleep(1)
        if i % 5 == 0:
            show_isr(admin, args.topic)

    print("\n=== 阶段 3:再写入,应当大量收到 NotEnoughReplicas 错误 ===")
    c2 = Counter()
    burst(producer, args.topic, 50, c2)
    c2.summary("after-shrink")
    show_isr(admin, args.topic)

    print(f"\n=== 阶段 4:把 Broker 起回来,等 ISR 扩张回去 ===")
    for c in targets:
        docker_compose(["start", c], args.compose_dir)
    print("[wait] 等待 60s 让 ISR 重新扩张…")
    for i in range(60):
        time.sleep(1)
        if i % 10 == 0:
            show_isr(admin, args.topic)

    print("\n=== 阶段 5:再写入,应当全部成功 ===")
    c3 = Counter()
    burst(producer, args.topic, 50, c3)
    c3.summary("recovered")
    show_isr(admin, args.topic)

    print("\n=== 总结 ===")
    print(f"  正常阶段 OK={c1.ok} FAIL={c1.fail}")
    print(f"  ISR<min  OK={c2.ok} FAIL={c2.fail}    ← 这里应当大量 NotEnoughReplicas")
    print(f"  恢复阶段 OK={c3.ok} FAIL={c3.fail}")
    print("\n💡 结论:min.insync.replicas 是「acks=all 真正能保住数据」的兜底,"
          "ISR 不够时宁可拒绝写入,也不让数据有丢的风险。")


if __name__ == "__main__":
    main()
python
#!/usr/bin/env python3
"""
observe_isr.py — 实时打印每个分区的 Leader / ISR / Replicas

用法:
    pip install confluent-kafka
    python observe_isr.py --bootstrap localhost:9092 --topic learn.09.orders
    python observe_isr.py --bootstrap localhost:9092 --topic learn.09.orders --interval 2

观察重点:
    1) 在另一个终端 `docker compose stop kafka1`,看本脚本输出里 Leader/ISR 的变化;
    2) 起回来后看 ISR 重新扩张;
    3) 配合 09_replication.md 中「故障场景剧本」对照阅读。

依赖:
    confluent-kafka >= 2.0(AdminClient.describe_topics 在 2.2 之后行为更直观;
    若只有旧版本,可改用 list_topics(...) 的 cluster_metadata,本脚本两种 API 都兼容)。
"""

from __future__ import annotations

import argparse
import sys
import time
from datetime import datetime

from confluent_kafka.admin import AdminClient


COLOR_RED = "\033[31m"
COLOR_YEL = "\033[33m"
COLOR_GRN = "\033[32m"
COLOR_DIM = "\033[2m"
COLOR_RST = "\033[0m"


def parse_args():
    p = argparse.ArgumentParser(description="实时观察 Kafka Topic 的 ISR / Leader 变化")
    p.add_argument("--bootstrap", default="localhost:9092",
                   help="bootstrap.servers")
    p.add_argument("--topic", required=True, help="要观察的 Topic 名")
    p.add_argument("--interval", type=float, default=1.0,
                   help="刷新间隔秒数(默认 1.0)")
    p.add_argument("--once", action="store_true",
                   help="只查一次后退出(用于脚本编排)")
    return p.parse_args()


def fetch_topic_meta(admin: AdminClient, topic: str):
    """优先使用 list_topics()(AdminClient 与 Producer/Consumer 都通用)。"""
    md = admin.list_topics(topic=topic, timeout=10)
    if topic not in md.topics:
        raise RuntimeError(f"Topic {topic!r} 不存在")
    t = md.topics[topic]
    if t.error is not None:
        raise RuntimeError(f"获取 Topic {topic} 元数据失败:{t.error}")
    parts = []
    for pid in sorted(t.partitions.keys()):
        p = t.partitions[pid]
        parts.append({
            "id": pid,
            "leader": p.leader,
            "replicas": list(p.replicas),
            "isr": list(p.isrs),
        })
    cluster = {
        "broker_ids": sorted(md.brokers.keys()),
        "controller_id": md.controller_id,
        "cluster_id": md.cluster_id,
    }
    return cluster, parts


def fmt_replicas(replicas, isr, leader):
    """把副本列表用颜色标注:Leader=黄、ISR 内=绿、OSR=红。"""
    chips = []
    for r in replicas:
        if r == leader:
            chips.append(f"{COLOR_YEL}{r}*{COLOR_RST}")
        elif r in isr:
            chips.append(f"{COLOR_GRN}{r}{COLOR_RST}")
        else:
            chips.append(f"{COLOR_RED}{r}{COLOR_RST}")
    return ",".join(chips)


def render(cluster, parts, topic):
    sys.stdout.write("\033[2J\033[H")  # 清屏
    print(f"=== Kafka ISR Watcher === {datetime.now().strftime('%H:%M:%S')}")
    print(f"cluster_id     : {cluster['cluster_id']}")
    print(f"controller     : {cluster['controller_id']}")
    print(f"alive brokers  : {cluster['broker_ids']}")
    print(f"topic          : {topic}")
    print()
    print(f"{'Part':<6}{'Leader':<8}{'Replicas (黄=Leader, 绿=ISR, 红=OSR)':<55}{'ISR 数'}")
    print("-" * 85)
    for p in parts:
        leader = p["leader"]
        leader_str = f"{COLOR_YEL}{leader}{COLOR_RST}" if leader >= 0 else f"{COLOR_RED}NONE{COLOR_RST}"
        rep_str = fmt_replicas(p["replicas"], p["isr"], leader)
        n_isr = len(p["isr"])
        n_rep = len(p["replicas"])
        warn = ""
        if n_isr < n_rep:
            warn = f"  {COLOR_YEL}⚠ under-replicated{COLOR_RST}"
        if leader < 0:
            warn = f"  {COLOR_RED}⚠ NO LEADER{COLOR_RST}"
        print(f"{p['id']:<6}{leader_str:<16}{rep_str:<70}{n_isr}/{n_rep}{warn}")
    print()
    print(f"{COLOR_DIM}按 Ctrl+C 退出{COLOR_RST}")


def main():
    args = parse_args()
    admin = AdminClient({"bootstrap.servers": args.bootstrap})

    try:
        while True:
            try:
                cluster, parts = fetch_topic_meta(admin, args.topic)
                render(cluster, parts, args.topic)
            except Exception as e:
                sys.stdout.write("\033[2J\033[H")
                print(f"{COLOR_RED}[ERROR] {e}{COLOR_RST}")
                print(f"{COLOR_DIM}下一次刷新:{args.interval}s 后重试…{COLOR_RST}")
            if args.once:
                break
            time.sleep(args.interval)
    except KeyboardInterrupt:
        print("\nbye.")


if __name__ == "__main__":
    main()

kill_leader_demo.py ↗ · min_isr_demo.py ↗ · observe_isr.py ↗