这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 多副本是 Kafka 不丢消息的根基。本章讲清楚副本如何同步、什么算“同步”、Leader 挂了如何选举,以及经典的一致性问题如何被 leader epoch 解决。
11.1 副本的角色
每个分区有一个 Leader 和若干 Follower:
- Leader:处理该分区所有生产与消费请求;
- Follower:唯一工作是向 Leader 发送 Fetch 请求,拉取数据写进本地日志。
注意这个反直觉的设计: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 条)
- LEO(Log End Offset):日志末端位移,下一条待写入消息的 offset;
- HW(High Watermark):所有 ISR 副本都已同步到的位移。消费者只能读到 HW 之前的消息。
HW 的推进过程:
- Follower 发送 Fetch,请求
fetchOffset = 自身 LEO; - Leader 处理 Fetch 时更新“该 Follower 已追上到哪”的视图;
- 若 ISR 全部追上某位移,Leader 推进 HW;
- Follower 从 Fetch 响应中获知新 HW,更新本地 HW。
11.3 ISR:什么算“同步”
ISR = 与 Leader 保持同步的副本集合(含 Leader)。 判定标准由 replica.lag.time.max.ms(默认 30 秒)决定:
- Follower 持续向 Leader 发 Fetch,且每次都能在超时时间内追上 Leader 的 LEO,就留在 ISR;
- 落后超过
replica.lag.time.max.ms,被 Leader 移出 ISR(收缩); - 重新追上后,再加入 ISR(扩张)。
相关命令观察:
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 把选择权交给你:
- 要一致性:RF=3 + min.insync.replicas=2 + acks=all;
- 要可用性:min.insync.replicas=1,宕机更多时仍可写,但有丢数据风险;
- 绝不推荐:RF=1 的“生产”主题。
11.5 Leader 选举
Leader 所在 Broker 宕机后,Controller 负责为受影响分区选新 Leader:
- Controller 通过元数据日志感知 Broker 下线;
- 对该 Broker 上的所有分区,从 ISR 中按顺序选择第一个存活副本作为新 Leader;
- 写入元数据并广播给所有 Broker;
- 客户端元数据刷新后,把请求发往新 Leader。
因为新 Leader 来自 ISR,它拥有 HW 之前的全部数据,已提交消息不会丢失。
Unclean Leader Election
如果 ISR 全部不可用,但某个非同步 Follower 还活着怎么办?
unclean.leader.election.enable=false(默认):拒绝选举,分区不可用,保数据;true:让落后副本上位,恢复可用性,但丢掉未同步的已提交消息。
这是经典的 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 副本均衡与机架感知
- preferred replica:AR(Assigned Replicas)列表中的第一个副本。Kafka 倾向让分区的 preferred replica 成为 Leader,从而使 Leader 均匀分布;
auto.leader.rebalance.enable=true(默认)时,Broker 定期检查并触发平衡; 也可手动执行:
bin/kafka-leader-election.sh --bootstrap-server localhost:9092 \
--election-type preferred --all-topic-partitions
- 配置
broker.rack后,Kafka 会尽量把副本分散到不同机架,防止机架交换机故障导致分区全灭。
11.8 动手实验:亲历一次选举
- 伪集群创建 topic:
test-ha,3 分区 RF=3; describe记录每个分区的 Leader/ISR;kill -9某 Leader 所在 Broker;- 立刻
describe:Leader 变化、ISR 收缩; - 期间持续发消息(注意
min.insync.replicas),观察写入是否正常; - 重启被 kill 的 Broker,观察它作为 Follower 追赶、ISR 扩张;
- 打开该分区的
leader-epoch-checkpoint,看 epoch 历史。
本章小结
- Follower 通过 Fetch 拉取同步;LEO 是日志末端,HW 是消费者可见边界;
- ISR 由
replica.lag.time.max.ms动态判定,acks=all 只等当前 ISR; RF=3 + min.insync.replicas=2是可靠性的标准配置;- Controller 从 ISR 选新 Leader;unclean election 用可用性换丢失风险;
- Leader epoch 解决了基于 HW 截断的分叉问题,是副本一致性的关键机制。
思考题
- 为什么消费者只能读到 HW 之前的消息?如果允许读到 HW 之后会怎样?
replica.lag.time.max.ms调大或调小,分别对可用性与数据安全有什么影响?- 什么情况下
acks=all仍然可能丢消息?给出至少两种场景与对应的防范配置。