乐途乐途
主页
  • 计算机基础

    • TCP/IP
    • Linux
    • HTTP
  • 数据库

    • SQL
    • MySQL 5.7
  • 编程语言

    • C
    • C++
    • Java SE
    • Python2
    • Python3
  • 数据格式

    • JSON
    • XML
  • 认证与安全

    • JWT
  • 工具

    • Markdown
  • Git

    • GitFlow
  • Quartz

    • Quartz
  • Java

    • Maven 入门
    • Maven 进阶
    • MyBatis
    • Spring
    • Spring MVC
  • Java

    • Spring Boot
    • Spring Cloud
    • Spring Cloud Alibaba
    • Spring Security
    • Spring AI
    • Spring Batch
    • Kafka
    • Java 设计模式
  • 缓存

    • Redis
  • 搜索引擎

    • Elasticsearch
  • 分布式协调

    • ZooKeeper
联系
阿里云
主页
  • 计算机基础

    • TCP/IP
    • Linux
    • HTTP
  • 数据库

    • SQL
    • MySQL 5.7
  • 编程语言

    • C
    • C++
    • Java SE
    • Python2
    • Python3
  • 数据格式

    • JSON
    • XML
  • 认证与安全

    • JWT
  • 工具

    • Markdown
  • Git

    • GitFlow
  • Quartz

    • Quartz
  • Java

    • Maven 入门
    • Maven 进阶
    • MyBatis
    • Spring
    • Spring MVC
  • Java

    • Spring Boot
    • Spring Cloud
    • Spring Cloud Alibaba
    • Spring Security
    • Spring AI
    • Spring Batch
    • Kafka
    • Java 设计模式
  • 缓存

    • Redis
  • 搜索引擎

    • Elasticsearch
  • 分布式协调

    • ZooKeeper
联系
阿里云
  • 学习路径
  • 第1章 消息队列与 Kafka 概述

    • 章节导读
    • 消息队列基础
    • Kafka 概述
    • Kafka 为什么快
  • 第2章 快速上手:单机环境搭建

    • 章节导读
    • 环境准备与安装
    • 快速启动
    • Topic 管理
  • 第3章 核心概念:主题、分区与日志

    • 章节导读
    • Topic
    • Partition
    • Offset
    • Segment 与存储结构
  • 第4章 生产者详解

    • 章节导读
    • Producer 概述
    • Producer 发送机制
    • Producer 分区策略
    • Producer 幂等与事务
    • Producer 配置
  • 第5章 消费者与消费组

    • 章节导读
    • Consumer 概述
    • Consumer Group 消费组
    • Consumer 分区分配策略
    • Consumer Offset 提交
    • Consumer 多线程
    • Consumer 配置
  • 第6章 Broker 与控制器

    • 章节导读
    • Broker 概述
    • Broker 配置
    • Controller
    • KRaft 共识协议
  • 第7章 副本与数据可靠性

    • 章节导读
    • 副本机制
    • ISR 与副本同步
    • Leader 选举
    • ACK 与一致性保证
    • 高水位与 Leader Epoch
  • 第8章 存储与性能优化

    • 章节导读
    • 存储架构
    • 日志清理与压缩
    • Page Cache 与零拷贝
    • Producer 性能优化
    • Consumer 性能优化
    • Broker 性能优化
  • 第9章 生产环境运维与监控

    • 章节导读
    • Topic 管理
    • Kafka 运维工具
    • 监控
    • 常见故障排查
  • 第10章 Kafka生态与面试考点

    • 章节导读
    • Kafka 生态全景
    • 面试高频 30 题

高水位与 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 可读到的最大 OffsetLeader 取所有 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 条消息后:

OffsetStatusConsumer 可见?
0Committed (HW≥1)✅ 是
1Committed (HW≥2)✅ 是
2Committed (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:

  1. HW 的定义:只有所有 ISR 副本都同步的消息才算 Committed(即 < HW 的消息)
  2. Leader 只在 ISR 中选:新 Leader 一定在 ISR 中,意味着它一定有所有 Committed 消息
  3. 新 Leader 以自身 HW 为起点:Follower 成为 Leader 后,消费者只能读到其 HW 之前的数据

但存在一个边界情况(Kafka 0.11 之前的 HW 截断问题):如果 Follower 的 HW 比 Leader 低,且 Leader 故障,Follower 的 HW 虽然是 ISR 中最小的,但可能截断了 Leader 的一些日志。Kafka 0.11+ 通过 Leader Epoch 解决了这个问题。

上一页
ACK 与一致性保证