高水位与 Leader Epoch
定义与作用
Kafka 使用高水位(High Watermark,HW) 和 Leader Epoch 两个机制来保证 Leader 切换时的数据一致性。HW 定义了"哪些消息对所有 ISR 可见(已确认)",Leader Epoch 用于防止旧 Leader 恢复后覆盖新 Leader 的数据。
这两个机制是 Kafka 在 CAP 理论中实现 CP(Consistency + Partition Tolerance)的关键。
核心原理
LEO 与 HW 的关系
| 术语 | 定义 | 谁维护 |
|---|---|---|
| LEO (Log End Offset) | 下一条待写入的 Offset | 每个副本各自维护 |
| HW (High Watermark) | Consumer 可读到的最大 Offset | Leader 取所有 ISR 中最小 LEO 作为 HW |
| Committed | 已提交消息 = Offset < HW 的消息 | Broker 不负责恢复已 Committed 消息 |
HW 的计算
Leader HW = min(Leader LEO, min(ISR 中各副本的 LEO))
Consumer 可见范围 = [0, HW)
Leader Epoch:防止数据回滚
Leader Epoch 防止的问题:旧 Leader 恢复后可能重新接收写入,覆盖新 Leader 已确认的数据。
Leader Epoch 机制流程
Leader Epoch Sequence File (存储在 Leader 端):
[LeaderEpoch=0, StartOffset=0]
[LeaderEpoch=5, StartOffset=100]
[LeaderEpoch=6, StartOffset=106]
→ Follower 恢复连接时:
发送 LeaderEpochRequest(当前 Epoch=5)
→ Leader 返回 Epoch=6 的 StartOffset=106
→ Follower 截断 >=106 的本地数据
→ 从 106 开始重新同步
完整示例
示例一:验证 HW 与 Consumer 可见性
场景:3 副本,ISR=[1,2,3],观察 Consumer 何时能读到消息。
# 终端 1:Consumer(正常消费)
bin/kafka-console-consumer.sh --topic demo-hw --group g1 \
--bootstrap-server localhost:9092
# 终端 2:Producer(acks=all,确保同步到 ISR)
bin/kafka-console-producer.sh --topic demo-hw \
--bootstrap-server localhost:9092 --producer-property acks=all
> msg1 # Consumer 收到(ISR 已同步 → HW 推进)
> msg2 # Consumer 收到
# 终端 3:观察 ISR 和 Leader
bin/kafka-topics.sh --describe --topic demo-hw --bootstrap-server localhost:9092
# Isr: 1,2,3 → HW 推进顺利
发送 3 条消息后:
| Offset | Status | Consumer 可见? |
|---|---|---|
| 0 | Committed (HW≥1) | ✅ 是 |
| 1 | Committed (HW≥2) | ✅ 是 |
| 2 | Committed (HW≥3) | ✅ 是 |
示例二:Leader Epoch 防止 Split-Brain
场景:模拟网络分区导致旧 Leader 存活但失去联系。
# 初始状态
bin/kafka-topics.sh --describe --topic demo-epoch --bootstrap-server localhost:9092
# Partition: 0 Leader: 1 Isr: 1,2,3
# 1. 网络隔离 Broker 1(模拟分区)
# iptables -A INPUT -p tcp --dport 9092 -j DROP (仅示例)
# 2. Controller 检测到 Broker 1 不可达 → 新 Leader = Broker 2
bin/kafka-topics.sh --describe --topic demo-epoch --bootstrap-server localhost:9097
# Partition: 0 Leader: 2 Isr: 2,3
# 3. 继续写入(Producer 自动重连到 Broker 2)
# Offset 100-120 写入到 Broker 2
# 4. 恢复 Broker 1 网络
# Broker 1 尝试重新连接
# Leader Epoch 机制:Broker 1 发现自己的 Epoch 已过期
# → 截断本地可能不一致的数据
# → 从 Broker 2 同步 Offset 100+
易错场景
易错 1:混淆 LEO 和 Consumer Offset
场景:用 Consumer Offset 判断数据是否"安全"。
正确理解:
- LEO = 下一条即将写入的位置(内部概念,Consumer 不关心)
- HW = 所有 ISR 已确认的位置(Consumer 只能读到 HW 之前)
- Consumer Offset = 该 Consumer 已消费到的位置(消费端概念)
三者完全独立。
易错 2:acks=1 时 HW 不推进的场景
场景:Producer 使用 acks=1(只等 Leader 确认),但 Follower 全部故障。
后果:Leader 的 LEO 在增长(数据持续写入),但 HW 停滞不前进(没有 Follower 同步)。Consumer 看不到新消息。
面试高频考点
Q:Kafka 如何保证"已提交的消息在 Leader 切换后不丢失"?
A:
- HW 的定义:只有所有 ISR 副本都同步的消息才算 Committed(即 < HW 的消息)
- Leader 只在 ISR 中选:新 Leader 一定在 ISR 中,意味着它一定有所有 Committed 消息
- 新 Leader 以自身 HW 为起点:Follower 成为 Leader 后,消费者只能读到其 HW 之前的数据
但存在一个边界情况(Kafka 0.11 之前的 HW 截断问题):如果 Follower 的 HW 比 Leader 低,且 Leader 故障,Follower 的 HW 虽然是 ISR 中最小的,但可能截断了 Leader 的一些日志。Kafka 0.11+ 通过 Leader Epoch 解决了这个问题。