主题
第 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_metadataTopic 里也有「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 关键结论
- HW 一定 ≤ LEO(不可能比已经写入的还高)。
- HW 推进至少需要一轮 Follower 的 Fetch(Leader 必须看到 Follower 的最新 fetch.offset)。
- HW 由 Leader 单点计算,再通过 FetchResponse 同步给 Follower。Follower 拿到最新 HW 之后才能更新自己本地 HW。
- 消费者的可见性取决于 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 回复 | 网络丢包就丢 |
1 | Leader 写入本地就回 ack | Leader 写入但 Follower 还没拉,Leader 死了就丢 |
all(或 -1) | ISR 中所有副本都写入后 才回 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=3,min.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=1 或 min.insync.replicas=1,仍然可能丢。三件套必须同时满足:副本数 ≥ 3、min.insync.replicas ≥ 2、acks=all。
7. 故障场景剧本
下面用三个剧本演练真实集群的故障行为。每个剧本前先约定环境:
集群:3 Broker(B1, B2, B3)
Topic:learn.09.orders,3 分区,每分区 3 副本,min.insync.replicas=2
Producer:acks=all、enable.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 涨幅只有 30T2: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 当新 LeaderT3:Controller 把新元数据广播给所有 Broker。
新 ISR = [B2, B3]
新 Leader = B2,Epoch +1T4:客户端缓存的元数据失效。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 脑裂」,因为:
- Controller 是单点(KRaft 下是 Active Controller,ZK 下争
/controller节点),所有 Leader/ISR 变更必须经过它。 - Broker 必须与 Controller 通信才能确认自己是 Leader;如果 Broker 与 Controller 失联(fenced),它会主动 step down,不再接受写入。
- 即使发生「老 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.size | message.max.bytes | replica.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.ordersTopic: 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.factor | 1(建 Topic 时指定) | 3 | 副本数 |
min.insync.replicas | 1 | 2(与 RF=3 配套) | acks=all 写入门槛 |
unclean.leader.election.enable | false | false(金融场景必须) | 是否允许从 OSR 选 Leader |
replica.lag.time.max.ms | 30000 | 视集群网络/磁盘抖动调 | 多久不追上就踢出 ISR |
num.replica.fetchers | 1 | 4~8(多 Broker 时) | Follower 拉取并发线程数 |
replica.fetch.max.bytes | 1048576 | 与 message.max.bytes 联动 | 每分区每次拉取上限 |
auto.leader.rebalance.enable | true | true | 自动把 Leader 还给 Preferred |
leader.imbalance.check.interval.seconds | 300 | 300 | 自动均衡的检查频率 |
leader.imbalance.per.broker.percentage | 10 | 10 | 触发均衡的不均衡百分比 |
Producer 端:
| 配置项 | 推荐 | 说明 |
|---|---|---|
acks | all | 配合 min.isr 才有意义 |
enable.idempotence | true | 重试不会产生重复 |
max.in.flight.requests.per.connection | ≤ 5(idempotence 限制) | 顺序保障 |
retries | Integer.MAX_VALUE | 配合 delivery.timeout.ms |
delivery.timeout.ms | 120000 | 总超时(包含所有重试) |
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 句话送给你:
- 不要把
acks=1当成「至少一次」——Leader 还没把数据同步到 Follower 就死掉,消息就没了。 - 不要在
min.insync.replicas=1的情况下声称「3 副本不丢消息」——一个副本的同步等于没同步。 - 不要轻易开
unclean.leader.election.enable=true——开了相当于「丢已 ack 消息换可用性」。 - 不要忽视 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 推进流程:
- Producer 写入 Leader,Leader 的 LEO 推进;
- Follower 发起 Fetch 请求(请求里带上自己当前的 fetch.offset),Leader 把数据返给 Follower,Follower 写入本地、LEO 推进;
- Follower 下一次 Fetch 时带上新的 fetch.offset,Leader 据此更新「Follower 的 LEO」记录;
- Leader 重新计算 HW = min(所有 ISR 副本的 LEO),HW 推进;
- 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。 - 常见原因:
- 网络抖动 / 跨机房延迟过大——fetch 请求来回慢;
- 磁盘 IO 打满——Follower 写本地 log 慢;
- GC 长时间停顿——fetcher 线程被卡住;
num.replica.fetchers太小——一个线程跨多个 Broker 串行 fetch,被慢的 Source 拖累;- 大消息超过
replica.fetch.max.bytes——Follower 永远拉不动; - Leader 被打满——Producer 写入压力过大,Leader 自顾不暇,fetch 请求处理慢。
- 排查工具:JMX
UnderReplicatedPartitions、server.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 中所有」),理由:
- 吞吐优先:Kafka 设计目标是高吞吐日志,要求「全员同步」可以让 Leader 快速推进 HW,配合批量 + 压缩压榨吞吐;
- 故障识别精细:用
replica.lag.time.max.ms时间窗自动识别掉队,掉队的 Follower 不进 ISR,就不会拖累 commit; - 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 + 序列号去重,不会重复; - 否则可能出现「同一条消息被重试多次写入」的重复。
- Producer 收到错误(如
加分项:能讲出「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 ↗