阶段二 · 架构原理与调优

副本机制与 ISR

前置回顾

  • 第 03 篇讲 acks 时 ISR 首次出现:acks=all 等的是 ISR 全体确认,不是全体副本
  • 第 06 篇给过标准配置:replication.factor=3 + min.insync.replicas=2,号称"最多容忍 1 个副本故障"
  • 第 08 篇说消息写进 page cache 就返回,可靠性交给副本
  • 本篇把"副本"这条线拆到底:同步怎么发生,已提交怎么定义,ISR 为什么缩容

一句话总结

每个 Partition 有一个 Leader 承担全部读写,Follower 像 Consumer 一样向 Leader 拉取同步.
消息被 ISR 全体副本写入后才算已提交(HW 推进);ISR 按时间阈值动态伸缩,Leader 只从 ISR 里选.
理解 ISR 与 HW,就理解了 Broker 侧"不丢"的全部机制.

为什么需要副本

第 08 篇说过,消息一进来就落在 Broker 磁盘.单台 Broker 挂掉,它上面所有分区的数据就没了:

单 Broker 部署:
Broker A 承载 P0, P1, P2
磁盘损坏 → 三个分区的数据全部消失,不可恢复

所以 Kafka 给每个 Partition 配 N 个副本(replication.factor),分散到不同 Broker:

replication.factor = 3:
P0: Leader 在 Broker 0,副本在 Broker 1, Broker 2
P1: Leader 在 Broker 1,副本在 Broker 2, Broker 0
P2: Leader 在 Broker 2,副本在 Broker 0, Broker 1

任何一台 Broker 宕机,每个 Partition 都还有 2 个副本可用

每个 Partition 的副本分两种角色:

角色 职责
Leader 接收 Producer 写入、响应 Consumer 拉取,分区唯一的读写入口
Follower 向 Leader 拉取数据保持同步,随时准备接替

读写只走 Leader.这是 Kafka 一致性模型简单的根源:一个分区同一时刻只有一个写者(第 08 篇的单写者模型),不存在多副本写冲突;Consumer 读到的也永远是 Leader 的提交视角.

Leader = 会议记录员,所有发言写进唯一的主记录本.
Follower = 抄写员,各持一本副本,不停照着主记录本抄.
记录员缺席,抄得最全的人带着本子顶上,记录不中断.

对比 MySQL 主从,相似但有关键差异:

维度 MySQL 主从 Kafka 副本
写入口 主库 Leader
从节点分担读 常见(读写分离) 默认不,读也走 Leader
同步方式 推 binlog Follower 拉取
故障切换 人工或中间件 自动(Controller 负责,下一篇讲)

Follower 怎么同步:拉,不是推

Follower 的同步动作和普通 Consumer 几乎一样:向 Leader 发 FetchRequest,按 offset 拉取:

                    ┌─ Leader 日志 ────────────────┐
│ ... [msg-7] [msg-8] [msg-9] │
│ │
Follower ── FetchRequest(offset=7) ───────────────►│
Follower ◄── FetchResponse(msg-7, msg-8, msg-9) ───┘

两个细节:

  1. FetchRequest 携带 follower 自己的 LEO,Leader 据此记录"这个 follower 追到哪了",这是后面 HW 计算的数据来源
  2. 没数据时请求挂起等待(第 04 篇 Consumer 的 long poll 是同一机制),新消息一写入立即返回,同步延迟毫秒级

为什么是拉而不是 Leader 推:

方案 问题
Leader 推 Leader 要跟踪每个 follower 的速率,慢副本拖垮 Leader,抖动时要做背压控制
Follower 拉 复用 Consumer 的协议和实现,follower 按自身能力拉,天然背压

已提交的定义:LEO 与 HW

两个水位线

  • LEO(Log End Offset):副本"写到哪了",即下一条消息的 offset
  • HW(High Watermark):分区"提交到哪了",Consumer 能读到的最大 offset + 1

HW 由 ISR 全体副本的 LEO 决定:

Leader:     [0] [1] [2] [3] [4] [5] [6] [7] [8] [9]     LEO = 10
Follower-1: [0] [1] [2] [3] [4] [5] [6] [7] [8] LEO = 9
Follower-2: [0] [1] [2] [3] [4] [5] [6] LEO = 7
▲
HW = min(10, 9, 7) = 7

offset < 7:已提交,Consumer 最多读到 msg-6
offset ≥ 7:未提交,Consumer 不可见

已提交(committed)的精确含义

一条消息被 ISR 中所有副本写入(它的 offset < HW)后,才算已提交.
已提交的消息:Consumer 可见,Leader 挂掉后由任何 ISR 副本接任都不丢.

与 acks 参数对账

第 03 篇的 acks 表格,用 LEO/HW 重新解释:

acks Leader 的行为 消息状态
0 不等回复 可能没到 Leader
1 写进 Leader 日志就回复 已写但未提交,在 HW 之后
all 等 ISR 所有副本的 LEO 越过该消息 已提交(HW 已推进)

acks=all 等的"ISR 确认",是 follower 下一次 FetchRequest 带来的 LEO 推进,不是 follower 落盘(follower 同样写 page cache,第 08 篇).

HW 的第二个作用:切换时的截断点

Leader 挂掉后,新 Leader 从 ISR 选出,所有副本把日志截断到 HW,丢弃 HW 之后的部分:

旧 Leader 日志: [0..6] [7] [8] [9]       ← 7,8,9 未提交,ISR 没同步完
新 Leader 日志: [0..6] [7] ← 它是 ISR 成员,至少有到 HW 的数据

截断后,所有副本统一到 HW=7:
新 Leader 日志: [0..6]

为什么截断不心疼:7,8,9 从未提交,Producer 没收到 ack,会重试

这正是第 06 篇"不丢"的存储层保证:acks=all 返回过的消息一定已提交,一定在所有 ISR 副本上,一定不会被截断.

leader-epoch-checkpoint 是干什么的

第 08 篇拆分区目录时见过 leader-epoch-checkpoint:记录每个 Leader 任期内写入的起始 offset.
精确的截断对齐靠它完成(新 Leader 与各 follower 对账,避免截断过多或过少),"截断到 HW"是简化后的等价描述.


ISR:动态伸缩的同步集合

定义

ISR(In-Sync Replicas)是"与 Leader 保持同步"的副本集合,Leader 永远在 ISR 里.第 07 篇 describe 输出里的 Isr 列就是它.

进出的标准:时间,不是条数

follower 落后超过 replica.lag.time.max.ms(默认 30s)没追上,被 Leader 踢出 ISR:

t=0s:  follower LEO=100, leader LEO=100,同步中
t=30s: follower LEO=100, leader LEO=100500,follower 期间零拉取
→ 踢出 ISR(它还在继续追,追上后自动回来)

为什么用时间阈值而不是"落后 N 条":

标准 问题
落后 N 条 消息速率波动大:N 条在低流量时是几小时,高流量时是一瞬间,阈值失效
30s 未追上 与流量无关,稳定表达"这个副本还能不能跟上"

老版本确实用过 replica.lag.max.messages(按条数),已被时间阈值取代.

缩容与扩容

ISR=[0,1,2] → follower-2 卡住 30s → ISR=[0,1] → follower-2 追上 → ISR=[0,1,2]

常见触发缩容的场景:

场景 机制
Follower Broker 长时间 GC 停顿 拉取停摆,30s 后踢出
网络抖动/分区 拉取超时,踢出
云盘 IO 抖动 拉得动但追不上,踢出
新增副本/分区迁移 追赶历史数据期间,暂不入 ISR

ISR = 车队编队:头车带队,跟车保持在后视镜可见距离内才算编队.
掉队的车继续开,追近就归队;头车抛锚,只有编队内的车有资格接任头车.

缩容对写入的影响

acks=all 时,Leader 等 ISR 全体确认.ISR 缩容,等待集合跟着变小:

ISR 大小 ≥ min.insync.replicas:  正常写入(甚至少等一个副本,延迟更低)
ISR 大小 < min.insync.replicas: Leader 拒绝写入,报 NotEnoughReplicas

这是第 06 篇"宁可不可用,也不丢数据"的实现:缩到危险线以下,写入口直接关闭.


Leader 选举与 Unclean Election

正常选举:只从 ISR 里选

Leader 挂掉后,Controller(下一篇的主角)从 ISR 里选一个新 Leader.ISR 里任意副本都满足两个条件:

  1. 拥有 HW 之前的所有消息(否则早被踢出)
  2. 30s 内同步过,数据最新

ISR 全挂:不干净的选举

极端情况:3 副本,2 台 Broker 同时宕机,只剩一个恰好被踢出 ISR 的落后副本:

replication.factor=3, min.insync.replicas=2

Broker 0: P0 Leader ✗ 宕机
Broker 1: P0 follower(ISR) ✗ 宕机
Broker 2: P0 follower(非ISR,落后 1000 条) ✓ 活着

ISR 空了,要恢复可用只能选 ISR 外的副本 → unclean election

由 unclean.leader.election.enable 决定(默认 false):

配置 行为 取舍
false(默认) 拒绝从 ISR 外选举,分区不可用,等 ISR 成员回来 用可用性换不丢
true 选落后的副本当 Leader,分区恢复可用 用不丢换可用性

unclean election 丢的是哪部分数据

丢的不是"落后那 1000 条"——那些从未提交,Producer 没收到 ack.
丢的是已提交的消息:新 Leader 上任后所有副本截断到它的位置,如果它缺了被踢出 ISR 前空档期的已提交消息,那些消息永久消失,而 Producer 早已收到 ack.

生产建议:保持 false.宁可分区短时不可用,也不静默丢已提交数据.只有"丢数据无所谓,可用性高于一切"的场景(如日志采集)才考虑打开.


动手验证:Leader 切换与 ISR 缩容

环境用第 07 篇的 docker-compose(3 Broker,replication.factor=3,min.insync.replicas=2).

观察 Leader 分布与 ISR

docker compose exec kafka-0 bash

kafka-topics.sh --describe --topic order-events --bootstrap-server localhost:19092
Topic: order-events  PartitionCount: 6  ReplicationFactor: 3
Partition: 0 Leader: 0 Replicas: 0,1,2 Isr: 0,1,2
Partition: 1 Leader: 1 Replicas: 1,2,0 Isr: 1,2,0
Partition: 2 Leader: 2 Replicas: 2,0,1 Isr: 2,0,1
...

6 个分区的 Leader 均匀分布在 3 个 Broker 上(第 02 篇讲过的负载均衡).

杀掉 Leader,看自动切换

# Partition 0 的 Leader 在 kafka-0 上,停掉它
docker compose stop kafka-0

# 等 controller 检测到故障(约几秒到十几秒),进 kafka-1 观察
docker compose exec kafka-1 bash
kafka-topics.sh --describe --topic order-events --bootstrap-server localhost:19092
Partition: 0  Leader: 1  Replicas: 0,1,2  Isr: 1,2

Leader 自动换成 1,死掉的 0 从 ISR 移除.此刻生产消费不受影响:

# 在 kafka-1 里继续生产
kafka-console-producer.sh --topic order-events \
--property "parse.key=true" --property "key.separator=:" \
--bootstrap-server localhost:19092
# 正常写入:ISR 还剩 2 个,满足 min.insync.replicas=2
# 输入几条后 Ctrl+C,然后 docker compose start kafka-0

用 Go 读 ISR

// kadm 是 franz-go 的 admin 子包
import "github.com/twmb/franz-go/pkg/kadm"

adm := kadm.NewClient(client) // client 是已建好的 *kgo.Client

md, err := adm.Metadata(ctx, "order-events")
if err != nil {
log.Fatal(err)
}
for _, p := range md.Topics["order-events"].Partitions {
fmt.Printf("partition=%d leader=%d replicas=%v isr=%v\n",
p.Partition, p.Leader, p.Replicas, p.ISR)
}

输出与 CLI 的 describe 一致.停 Broker 后再跑这段,能看到 ISR 收缩.


生产注意事项

标准配置(第 06 篇的落地)

replication.factor=3 + min.insync.replicas=2 + acks=all + unclean.leader.election.enable=false:

  • 容忍 1 个副本故障不丢数据
  • 2 个副本故障时 ISR 只剩 1 个 < 2,拒绝写入
  • ISR 全挂时拒绝 unclean 选举,完整覆盖"可用性换不丢"

副本必须分散到不同故障域

3 个副本挤在一台物理机/一个机架,复制等于没复制,一次断电全灭.
配置 broker.rack 标记机架,Kafka 会把同一 Partition 的副本尽量分到不同机架;云上把可用区当机架用.

Follower 假死时,写入会停摆约 30s

follower 宕机但还没被踢出 ISR 的窗口期(最长 replica.lag.time.max.ms),acks=all 的 Produce 会等这个"已死"的副本确认,直到它被踢出.
表现为写入延迟尖刺(约 30s)而非失败.不丢数据,但会冲击依赖低延迟的链路;调小 replica.lag.time.max.ms 可缩短窗口,代价是 ISR 对抖动更敏感.

副本同步吃带宽

follower 拉取流量 ≈ 写入流量 × (replication.factor - 1).
跨机架/跨可用区的同步流量更贵,容量规划要按这个公式算网卡,不能只看写入流量.

盯住 ISR 相关指标

第 07 篇列过 UnderReplicatedPartitions(>0 持续 5min 告警)和 IsrShrinksPerSec(突增告警).
这两个指标直接反映本篇机制:缩容频繁说明有副本在掉队,离"拒绝写入"不远了.


Follower 能分担读吗?

默认不能,读全走 Leader.好处是 Consumer 读到的永远是提交状态,不存在"从库读到旧数据"的问题.
Kafka 2.3+ 支持 consumer rack 感知(KIP-392):Consumer 与 follower 同机架时可就近读,用于削减跨机房读流量,不改变"读 Leader"的默认行为.

Leader 与 Follower 的日志完全一样吗?

HW 之后不一样:follower 可能还没追到最新写入,任何副本 HW 之后都可能有未提交的"脏数据".
所以切换 Leader 时先截断到 HW 再开始新一轮写入,保证所有副本的已提交前缀一致.

replica.lag.time.max.ms 调小,故障切换会更快吗?

不会.踢出 ISR 的是落后 follower,不触发选举;选举延迟由 Controller 的故障检测决定(下一篇讲).
调小的实际效果:缩容更敏感,抖动时 ISR 频繁进出,收益只是"选举资格审查更严".

acks=1 返回成功的消息,Consumer 什么时候能读到?

HW 推进之后.acks=1 只保证 Leader 写入了,消息还在 HW 之后,Consumer 拉不到;等 follower 追上、HW 越过它,才可见.
若 Leader 在 follower 追上之前挂了,这条消息被截断,Producer 已收到 ack 但消息消失.这就是 acks=1 丢消息的存储层解释.

副本数从 3 降到 2,min.insync.replicas 要不要跟着降?

不要.2 副本 + min.insync.replicas=2 意味着任何一台宕机都拒绝写入,可用性比 3 副本 + min.insync=2 还差.
3 副本的冗余不是浪费:多出的一个副本让 min.insync=2 在单台故障时依然成立.


快速回顾

  • 副本模型:每个 Partition 一个 Leader 管读写,follower 拉取同步,分散在不同 Broker
  • 同步方式:follower 发 FetchRequest 并上报 LEO,复用拉模型,延迟毫秒级
  • 已提交 = offset < HW:HW 是 ISR 全体 LEO 的最小值,Consumer 只能读到 HW 之前
  • 切换截断:新 Leader 上任,所有副本截断到 HW,未提交消息靠 Producer 重试
  • ISR 按时间伸缩:replica.lag.time.max.ms 默认 30s,缩到 min.insync 以下拒绝写入
  • 选举与 unclean:只从 ISR 选 Leader;ISR 全挂时默认拒绝 unclean 选举,宁可不可用不丢已提交

动手练习

  1. 观察 ISR 动态:describe 记录初始 ISR,停一个 follower Broker 等 40s 再看,重启后观察回归.
  2. 体验 30s 写停摆:建一个 rf=3、min.insync.replicas=1 的 Topic,停掉两个 follower,立即生产,观察 Produce 卡住约 30s(等 ISR 踢出)后恢复.
  3. 验证 min.insync.replicas:用 kafka-configs.sh 把 order-events 的 min.insync.replicas 改成 3,停一个 Broker 后生产消息,观察 NotEnoughReplicas 错误,再改回 2.
  4. 体验 Leader 切换:停 Leader Broker,在另一个 Broker 上 describe,记录新 Leader 与 ISR 变化;重启旧 Broker,观察它回归 ISR 的过程.
  5. 写 Go 观察者:用 kadm 每 2s 拉一次 Metadata 打印 Leader/ISR,配合步骤 1 实时观察缩容与恢复.
  6. 实验 unclean election:停掉两个 Broker 只剩一个非 ISR 副本,分别用 unclean.leader.election.enable=false/true 观察:前者分区不可用,后者恢复但丢已提交消息.

下一篇进入 Controller 与元数据管理:ZooKeeper 模式与 KRaft 模式的差别,Controller 怎么选举,谁来决定 Leader 换人.