Kafka 的可靠性建立在「多副本」之上,但副本不是简单复制——ISR、HW、Leader 选举、控制器这些机制决定了数据一致性边界。本文深入副本与控制器:它们怎么协作、故障时怎么决策、以及数据到底会不会丢。
1. 副本与分区的层级关系
1.1 Leader 与 Follower
一个分区的副本数等于 replication.factor。其中:
# 副本角色
# Leader 副本: 唯一接收读写,Producer 只写 Leader
# Follower 副本: 从 Leader 拉取数据(fetch),保持同步
# 同步集 ISR: 与 Leader 保持同步的副本集合
# 消费者: 可读 Leader 或配置的特定副本(follower 读取是可选特性)
副本模型的要点:只有 Leader 提供读写服务,Follower 只做数据备份。某个 Follower 挂了,Leader 继续服务;Leader 挂了,从 ISR 中选新 Leader。
1.2 副本同步的两条线
Follower 从 Leader 拉数据的节奏与消费相似:Follower 发送 FetchRequest,Leader 返回新增消息并记录该 Follower 的读取位置。同步的「健康」判断有两层:跟上 LEO 的滞后阈值(replica.lag.time.max.ms,默认 30s)与从未落后超过阈值。满足即留在 ISR。
2. ISR:在同步副本集合
2.1 ISR 的收缩与扩张
**ISR(In-Sync Replicas)**是「与 Leader 保持同步」的副本集合,由 Leader 动态维护:
# ISR 的规则
# 收缩: Follower 落后超过 replica.lag.time.max.ms(默认 30s)→ 被移出 ISR
# 扩张: 落后的 Follower 重新追上 LEO → 被加回 ISR
# 关键: 只有 ISR 内的副本才有资格被选为 Leader
ISR 收缩时,Leader 会向 ISR Shrink 事件、指标与告警。频繁收缩说明有副本长期落后(网络、磁盘慢、broker 繁忙)。
2.2 ISR 与数据可靠性
acks=all 的生产端等待 ISR 内所有副本确认。因此:
- ISR 越小,可用性风险越高:ISR 只剩 Leader 一个时,acks=all 退化等同 acks=1。
- ISR 与
min.insync.replicas:生产端acks=all时,若 ISR <min.insync.replicas,写入报NotEnoughReplicasException——这是「宁可拒写也不丢」的保护。
# 关键组合: acks=all + min.insync.replicas=2 + replication.factor=3
# Leader 挂、剩下 ISR≥2 时仍可写;ISR 收缩到 1 时拒写,保护数据
# 这个组合是「不丢数据」场景的推荐配置
3. HW 与 LEO:一致性的刻度尺
3.1 LEO 与 HW 的定义
- LEO(Log End Offset):副本日志的末尾偏移量,即已写入的下一条。
- HW(High Watermark):ISR 中所有副本 LEO 的最小值。HW 之前的数据对所有 ISR 副本可见。
# HW 的推进
# Leader 根据 ISR 各 Follower 的 fetch 位置,取最小值作为 HW
# 消费端只能读到 < HW 的消息(HW 是可读边界)
# 位移提交同样以 HW 为界,防止「读未同步数据」
HW 机制回答一个关键问题:什么数据才算「稳定」? 答案是「HW 之前」——它保证任何 ISR 副本都能读到这条数据。未过 HW 的数据在 Leader 挂掉时可能「读不到了」(读后即丢,但不会读到「假数据」)。
3.2 HW 更新的两个方向
HW 随「写入推进」与「follower 追上」而更新。生产上 HW 推进慢的常见原因:慢 Follower 拖住 HW(ISR 里有落后的副本,HW 停在它的位置)。此时要么等它追上,要么它被移出 ISR 后 HW 恢复推进。
4. Leader 选举与故障恢复
4.1 谁有资格当 Leader
Leader 挂掉后,控制器从 ISR 内选新 Leader(unclean.leader.election.enable=false 时)。选出的 Leader 保证「已有数据不丢」。
# 选举优先级(控制器视角)
# 1) ISR 内第一个存活的副本(优先选 ISR)
# 2) 配置 unclean.leader.election.enable=true 时,允许选 ISR 外的副本
# → 但会丢「ISR 未同步的数据」,可用性换一致性
unclean 选举是最后手段:开启后 Leader 挂且 ISR 全挂时,系统还能选个「数据落后」的副本当 Leader 维持服务,但丢数据。默认关闭,只有「服务连续性 > 数据完整性」的场景(如日志非关键)才开。
4.2 故障恢复的时间线
# Leader 故障恢复
# 1) 检测: 控制器通过 zookeeper/KRaft 元数据发现 Leader 下线
# 2) 选举: 从 ISR 选新 Leader,更新分区的 leader epoch
# 3) 接管: 新 Leader 开始提供读写,Follower 重新同步
# 4) 位移: HW 以新 Leader 的 LEO 与 ISR 重建
# 恢复速度: 秒级;ISR 健康时消费者感知不到中断
Leader Epoch 是防「僵尸 Leader」的机制:旧 Leader 恢复后带着过期数据回来,epoch 更高者胜出,防止旧 Leader 把已废弃的数据写回。
5. 控制器:集群的元数据大脑
5.1 控制器职责
控制器(Controller)是集群里「管元数据变更」的节点:
# 控制器的职责
# 1) 分区 Leader 选举
# 2) 分区迁移/副本调整的执行
# 3) 集群成员变更(broker 上下线)的处理
# 4) 元数据(主题/分区/副本)的广播与维护
KRaft 时代:控制器从 Zookeeper 迁移到内置的 KRaft 元数据日志(Kafka 3.x 全面引入,3.4 起生产可用)。KRaft 简化部署(去 ZK)、用日志共识替代 ZK 的读一致性模型,控制器成为**仲裁组(Controller Quorum)**内的节点。
5.2 控制器选举与故障转移
KRaft 下控制器由仲裁组通过 KRaft 协议(Raft 变体)选举:某节点获得多数票成为 Active Controller。控制器故障时,另一个仲裁节点接任,元数据日志保证状态一致。整个故障转移在秒级完成,期间分区不迁移、只暂停元数据变更。
# 观察控制器健康
# 指标: kafka.controller:ActiveControllerCount(应为 1)
# 日志: 控制器变更、ISR 变更、Leader 变更,都是排查线索
# 告警: ActiveControllerCount != 1 或频繁 LeaderElection 需要关注
6. 分区迁移与副本调整
6.1 为什么需要迁移
扩容缩容、broker 均衡、磁盘打满,都需要把分区在 broker 间搬移:
# 触发迁移的典型场景
# 1) 新增 broker 后把部分分区迁过去均衡负载
# 2) 某 broker 磁盘告警,迁走其 Leader 副本
# 3) 调整 replication.factor(加副本/减副本)
# 工具: kafka-reassign-partitions.sh 配合 JSON 计划
6.2 迁移的过程与代价
迁移本质是新副本从旧 Leader 全量同步 → 追平 → 入 ISR → 切换:
# 迁移过程
# 1) 目标 broker 新建副本,从 Leader 全量拉取(catch-up)
# 2) 追平后进入 ISR,新旧副本并存
# 3) 按计划把 Leader 迁移到目标副本,更新元数据
# 4) 旧副本按 retention/迁移策略清理
# 代价: 全量复制吃网络与磁盘 IO,迁移期监控副本滞后
工程要点:迁移要错峰执行(分批次、限速)、迁移中监控 ISR 收缩(网络扛不住会掉 ISR)、迁移完成验证 Leader 分布均衡。
7. 副本不同步的排查
# 排查副本长期不同步(非 Leader 副本 lag)
# 1) 看 broker 日志: "Replica fetcher" 相关错误(网络、连接重试)
# 2) 看磁盘: 落后副本所在 broker 磁盘 IO 是否饱和
# 3) 看网络: broker 间带宽是否被打满(抓包/监控)
# 4) 看 JVM: 副本拉取线程是否被 GC 阻塞
# 5) 看配置: replica.fetch.max.bytes 与 fetch 线程数是否够大
常见根因:单 broker 磁盘慢拖累全分区、broker 间网络拥塞、数据量突增超过 follower 追赶能力。临时手段是把落后副本移出 ISR(它不再拖 HW),根本手段是修根因或迁走分区。
8. 总结
Kafka 的副本与控制器机制回答了可靠性工程的核心问题:什么时候会丢数据、丢了丢多少、丢了怎么恢复。ISR 划定「谁同步、谁能当 Leader」,HW 划定「什么数据稳定」,控制器与 KRaft 管元数据与选举,迁移与排查管日常维护。落地建议:可靠性敏感场景坚持 acks=all + min.insync.replicas≥2,Leader 选举关闭 unclean,迁移错峰执行并盯 ISR,故障恢复时先看 Leader 变更与 ISR 收缩。理解副本机制,才能在「可用性 vs 一致性」之间做出有依据的选择。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。