KafkaNotes

第 11 章:副本机制:ISR、LEO 与高可用

zjc 于 2026-01-11 发布

这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 多副本是 Kafka 不丢消息的根基。本章讲清楚副本如何同步、什么算“同步”、Leader 挂了如何选举,以及经典的一致性问题如何被 leader epoch 解决。

11.1 副本的角色

每个分区有一个 Leader 和若干 Follower:

注意这个反直觉的设计:Follower 拉,而不是 Leader 推。好处是每个 Follower 按自己的节奏同步,慢的 Follower 不会拖累 Leader 的写入性能。

11.2 LEO 与 HW

两个核心指针:

Leader 日志:  [0][1][2][3][4][5]   LEO=6
                              ^
                              HW=5  (已被 ISR 全部同步的位置)

Follower A 日志: [0][1][2][3][4][5]  LEO=6  (已请求到 6)
Follower B 日志: [0][1][2][3][4]     LEO=5  (落后 1 条)

HW 的推进过程:

  1. Follower 发送 Fetch,请求 fetchOffset = 自身 LEO
  2. Leader 处理 Fetch 时更新“该 Follower 已追上到哪”的视图;
  3. 若 ISR 全部追上某位移,Leader 推进 HW;
  4. Follower 从 Fetch 响应中获知新 HW,更新本地 HW。

11.3 ISR:什么算“同步”

ISR = 与 Leader 保持同步的副本集合(含 Leader)。 判定标准由 replica.lag.time.max.ms(默认 30 秒)决定:

相关命令观察:

kafka-topics.sh --describe --topic order-events
# Isr 字段少了某个副本 = 收缩

kafka-run-class.sh kafka.tools.IsrChangeDetector   # 旧版工具
# 新版直接看 JMX 指标 UnderReplicatedPartitions

ISR 是动态的,这带来一个重要性质:acks=all 等待的是“当前 ISR”的确认,不是所有副本。如果 ISR 收缩到只剩 Leader,acks=all 实际退化为 acks=1。这就是必须配置 min.insync.replicas 的原因。

11.4 min.insync.replicas 的保护作用

假设 RF=3,min.insync.replicas=2

场景 ISR 数量 acks=all 写入
三台健康 3 成功,等全部 ISR
一台宕机 2 成功,至少两副本持久化
两台宕机 1 失败,抛 NotEnoughReplicasException

最后一种行为是故意的:宁可暂时不可写,也不允许只有一副本就确认“安全”。可用性与一致性的权衡,Kafka 把选择权交给你:

11.5 Leader 选举

Leader 所在 Broker 宕机后,Controller 负责为受影响分区选新 Leader:

  1. Controller 通过元数据日志感知 Broker 下线;
  2. 对该 Broker 上的所有分区,从 ISR 中按顺序选择第一个存活副本作为新 Leader;
  3. 写入元数据并广播给所有 Broker;
  4. 客户端元数据刷新后,把请求发往新 Leader。

因为新 Leader 来自 ISR,它拥有 HW 之前的全部数据,已提交消息不会丢失。

Unclean Leader Election

如果 ISR 全部不可用,但某个非同步 Follower 还活着怎么办?

这是经典的 CAP 取舍。支付、订单类系统应保持 false;某些允许少量丢失的指标日志场景可以考虑 true,但要有清晰的业务确认。

11.6 Leader Epoch:解决截断不一致

老问题

HW 更新有延迟。Follower 重启后需要截断日志到 HW 以保证与 Leader 一致,但在极端时序下(Follower 的 HW 尚未从 Leader 获取、而它曾短暂被当成 Leader),可能出现副本间数据分叉:

场景:Follower A 先被选为 Leader,写入 2 条后宕机未同步;
     B 上位成为新 Leader,A 恢复后按旧 HW 截断,可能多删或少删。

Leader Epoch 机制

每次 Leader 变更,epoch(纪元)+1。Follower 重启或追赶时,先向 Leader 发 OffsetsForLeaderEpoch 请求:

Follower: "我这里有 epoch=5, endOffset=87,你们呢?"
Leader:   "epoch=6 从 offset=87 开始,6 的 endOffset=90"
-> Follower 精确截断到 87,再从 87 开始同步

副本不再依赖可能滞后的 HW 做截断,而是以 Leader 的权威 epoch 记录为准,消除了数据分叉。leader-epoch-checkpoint 文件保存的正是这组 (epoch, startOffset) 历史。

11.7 副本均衡与机架感知

bin/kafka-leader-election.sh --bootstrap-server localhost:9092 \
  --election-type preferred --all-topic-partitions

11.8 动手实验:亲历一次选举

  1. 伪集群创建 topic:test-ha,3 分区 RF=3;
  2. describe 记录每个分区的 Leader/ISR;
  3. kill -9 某 Leader 所在 Broker;
  4. 立刻 describe:Leader 变化、ISR 收缩;
  5. 期间持续发消息(注意 min.insync.replicas),观察写入是否正常;
  6. 重启被 kill 的 Broker,观察它作为 Follower 追赶、ISR 扩张;
  7. 打开该分区的 leader-epoch-checkpoint,看 epoch 历史。

本章小结

思考题

  1. 为什么消费者只能读到 HW 之前的消息?如果允许读到 HW 之后会怎样?
  2. replica.lag.time.max.ms 调大或调小,分别对可用性与数据安全有什么影响?
  3. 什么情况下 acks=all 仍然可能丢消息?给出至少两种场景与对应的防范配置。