13851 字
69 分钟
Kafka 存储引擎与高性能原理:一条 Record 如何落盘并被读取

Kafka 存储引擎与高性能原理:一条 Record 如何落盘并被读取#

从 RecordBatch、LogSegment、稀疏索引到 Page Cache,沿完整数据路径解释 Kafka 为什么快

阅读目标#

第一篇建立了 Kafka 的核心世界观:Topic 由多个 Partition 组成,每个 Partition 是一条有序、可持久化、可按 Offset 读取的 Log。

但这仍然只是逻辑模型。

第二篇要把一个 Partition 放大到磁盘和内核层,回答两个真正决定 Kafka 工程能力的问题:

  1. 一批 Record 如何从 Broker 网络请求进入 Partition Log,并最终形成磁盘文件?
  2. Consumer 请求 Offset=36891 时,Broker 如何找到对应 RecordBatch,并高效返回数据?

完成本篇后,你应该能够:

  • 区分 RecordRecordBatchLogSegmentPartition Log 四个层次;
  • 解释 .log.index.timeindex.txnindex 的职责;
  • 手工推导一次 Offset → Segment → Relative Offset → Physical Position → RecordBatch
  • 画出 Broker 侧 Produce 和 Fetch 的核心执行链;
  • 区分 Kafka Offset 与文件字节 Position;
  • 区分“写入 Page Cache”和“同步到物理介质”;
  • 准确解释 Page Cache、批处理、压缩和高效文件传输在 Kafka 数据路径中的作用;
  • 不再用“因为顺序写和零拷贝”这一句空泛结论回答 Kafka 为什么快。

本文以 Apache Kafka 4.3 文档和当前 4.x 源码结构为主线。Kafka 4.3 的消息格式仍以 RecordBatch 为磁盘和协议数据单元;生产级源码阅读时,应固定到实际使用的 4.3.x Tag,避免把不断变化的 trunk 行号当作稳定接口。110


0. Offset=36891 到底存在哪里#

假设 Consumer 向 Broker 发起以下 Fetch:

Topic = order-events
Partition = 3
Offset = 36891
MaxBytes = 1 MiB

从第一篇的逻辑模型看,它是在表达:

order-events-3 这条 Partition Log 的 Offset 36891 开始读取一批数据。

但磁盘只认识目录、文件和字节位置,并不认识业务上的 Offset。

假设 Broker 的数据目录中存在:

order-events-3/
├── 00000000000000000000.log
├── 00000000000000000000.index
├── 00000000000000000000.timeindex
├── 00000000000000320000.log
├── 00000000000000320000.index
├── 00000000000000320000.timeindex
├── 00000000000000640000.log
├── 00000000000000640000.index
└── 00000000000000640000.timeindex

Broker 至少需要连续回答以下问题:

  1. order-events-3 对应哪个本地 Partition?
  2. Offset 36891 位于哪个 Segment?
  3. 目标 Offset 在 Segment 内的相对位置是多少?
  4. .index 是否精确保存 Offset 36891?
  5. 如果索引只记录了附近位置,还需要扫描多远?
  6. 找到 RecordBatch 后,数据需要复制进 JVM Heap 吗?
  7. 如果数据刚刚被写入,它是否仍在操作系统 Page Cache 中?
  8. 如果数据已经很旧,读取路径又会发生什么变化?

Offset 36891 的磁盘定位悬念

这一串问题揭示了 Kafka 存储引擎的核心矛盾:

Consumer 使用逻辑 Offset 读取
磁盘只能按文件与字节 Position 读取
Kafka 必须在两种坐标系之间建立高效映射

接下来,我们从第一篇中的 Partition 开始,把它逐层展开。


1. 从 Partition 到 Kafka 存储引擎#

1.1 逻辑 Log 与物理文件不是同一个概念#

第一篇中的 Partition 是一条逻辑上连续增长的 Log:

Partition 3
Offset
0 → 1 → 2 → 3 → 4 → 5 → 6 → ...

这种表示非常适合解释顺序、Replay 和 Consumer Position,但它隐藏了真实物理结构。

在磁盘上,一个 Partition 通常对应一个独立目录,目录中由多个 Segment 文件集合共同组成这条 Log:

Partition Directory
├── Segment 0
├── Segment 1
├── Segment 2
└── Active Segment

因此需要建立第一组映射:

逻辑概念物理实现
Topic一组 Partition
Partition独立目录和一条逻辑 Log
Log按 Base Offset 排序的多个 Segment
Segment数据文件与辅助索引文件集合
RecordRecordBatch 内部的一条业务记录
OffsetPartition 内的逻辑位置
Physical Position.log 文件中的字节位置

这里最容易混淆的是:

Kafka Offset ≠ 文件字节 Position

Offset=36891 并不意味着文件的第 36891 个字节。

因为每条 Record 的 Key、Value、Header 长度都可能不同;RecordBatch 还包含批头、压缩数据、CRC 和变长编码。逻辑 Offset 与物理字节位置之间不存在简单的一一等长关系。

1.2 为什么不能只使用一个无限增长的文件#

假设一个 Partition 永远只写:

partition.log

随着业务运行,这个文件可能达到几百 GB 甚至数 TB。此时系统需要处理:

  • 删除七天前的数据;
  • 重建损坏索引;
  • 崩溃后检查尾部不完整写入;
  • 对历史 Key 做 Log Compaction;
  • 把冷数据迁移到远程存储;
  • 迁移或重新分配 Partition;
  • 对文件做关闭、回收和生命周期管理。

如果只有一个大文件,删除旧数据就意味着修改或重写文件前半部分;索引损坏时可能扫描整条超大 Log;Compaction 也缺少可控的处理边界。

Kafka 因而把 Partition 拆成多个 Segment:

Partition 3
Segment 0 [0 … 31999]
Segment 1 [32000 … 63999]
Segment 2 [64000 … 95999]
Active Segment [96000 … ]

Segment 为以下工程操作提供了独立边界:

Roll
Retention Delete
Compaction
Recovery
Transfer
Tiering

Segment 不是为了改变 Partition 的逻辑连续性,而是为了让持续增长的 Log 可以被分块管理。

1.3 Active Segment 与 Closed Segment#

一个 Partition 在任意时刻通常只有一个当前写入 Segment:

[Closed Segment]
[Closed Segment]
[Closed Segment]
[Active Segment] ← 新 RecordBatch 追加到这里

Active Segment#

当前接受追加写入的 Segment。Kafka 不需要在所有历史文件中寻找写入位置,只需持续向当前文件尾部追加。

Closed Segment#

已经停止追加的历史 Segment。它的数据范围基本稳定,适合参与:

  • Retention 删除;
  • Log Compaction;
  • 远程分层迁移;
  • 文件校验和恢复;
  • 迁移与复制。

Kafka 的 Tiered Storage 也以完成写入的 Segment 为重要迁移单位:常见实时 Tail Read 由本地层和 OS Page Cache 提供,较旧历史数据才更多用于回填、恢复或远程读取。7

1.4 Partition 是“有序 Segment 集合”#

所以,一个更接近实现的 Partition 模型是:

Partition
Ordered Segment Collection
Closed Segments + One Active Segment
Each Segment = Data File + Index Files

第一篇说 Partition 是一条 Log;第二篇需要补充:

Partition 在逻辑上连续,在物理上由多个按照 Base Offset 排序的 Segment 拼接而成。


2. RecordBatch:Kafka 真正处理的数据单元#

2.1 错误认知:Kafka 一次只处理一条消息#

初学者很容易把 Kafka 想象成:

Record 1 → 一次网络请求 → 一次磁盘写入
Record 2 → 一次网络请求 → 一次磁盘写入
Record 3 → 一次网络请求 → 一次磁盘写入

如果真的如此,每条小消息都可能重复承担:

  • 请求头;
  • 协议解析;
  • 系统调用;
  • CRC 处理;
  • 压缩上下文;
  • 网络包与磁盘写入调度。

Kafka 的当前消息格式以 RecordBatch 为核心:消息始终以 Batch 形式写入;即便 Batch 只有一条 Record,它仍是一个 RecordBatch。1

Record 1
Record 2
Record 3
Record 4
RecordBatch
├── 共享批头
├── 共享 CRC 边界
├── 共享压缩边界
├── 连续网络传输
└── 连续磁盘追加

2.2 四层数据结构必须分清#

Record
RecordBatch
LogSegment
Partition Log

Record#

一条业务记录,包含经过序列化的:

Key
Value
Timestamp Delta
Offset Delta
Headers

RecordBatch#

一组被统一编码、校验、压缩、传输和存储的 Record。

LogSegment#

包含多个连续 RecordBatch 的数据文件与索引文件集合。

Partition Log#

由多个按 Base Offset 排序的 Segment 构成的完整逻辑日志。

如果把四者混成“消息文件”,就无法准确理解压缩、索引、Offset 分配和 Fetch。

2.3 RecordBatch Header#

Kafka 4.3 文档给出的当前 RecordBatch 磁盘格式包含以下核心字段:1

RecordBatch
├── baseOffset: int64
├── batchLength: int32
├── partitionLeaderEpoch: int32
├── magic: int8
├── crc: uint32
├── attributes: int16
│ ├── compressionType
│ ├── timestampType
│ ├── isTransactional
│ └── isControlBatch
├── lastOffsetDelta: int32
├── baseTimestamp: int64
├── maxTimestamp: int64
├── producerId: int64
├── producerEpoch: int16
├── baseSequence: int32
├── recordsCount: int32
└── records[]

RecordBatch 数据结构

本篇重点关注:

  • baseOffset:Batch 第一条 Record 的绝对 Offset;
  • lastOffsetDelta:Batch 最后一条 Record 相对 Base Offset 的增量;
  • baseTimestamp:Batch 时间戳基准;
  • attributes:压缩类型等 Batch 属性;
  • crc:Batch 数据完整性校验;
  • recordsCount:Batch 中 Record 数量。

而以下字段只先认识,不展开机制:

producerId
producerEpoch
baseSequence
isTransactional
isControlBatch

它们与幂等 Producer、事务和 Control Batch 密切相关,将在第五篇完整展开。

2.4 Base Offset 与 Offset Delta#

假设:

baseOffset = 1000

Batch 内部四条 Record 可以保存:

Record A: offsetDelta = 0 → Offset 1000
Record B: offsetDelta = 1 → Offset 1001
Record C: offsetDelta = 2 → Offset 1002
Record D: offsetDelta = 3 → Offset 1003

最终 Offset 的计算是:

absoluteOffset = baseOffset + offsetDelta

这样不需要每条 Record 都重复保存完整 64 位绝对 Offset。

时间戳也使用类似思路:

baseTimestamp = T₀
Record A: timestampDelta = 0
Record B: timestampDelta = 12 ms
Record C: timestampDelta = 24 ms

这体现了一种典型的批数据编码方式:

Batch 共享元数据
+
Record 保存局部 Delta
=
减少重复字段和编码体积

2.5 为什么压缩边界是 Batch#

如果逐条压缩:

compress(record1)
compress(record2)
compress(record3)

压缩算法很难充分利用不同 Record 之间的重复模式。

而业务事件通常具有高度相似的结构:

{
"orderId": "...",
"userId": "...",
"eventType": "OrderCreated",
"timestamp": "..."
}

字段名、Schema 和许多枚举值会反复出现。将多条 Record 放进同一个 Batch 后再压缩:

compress(record1 + record2 + record3 + ...)

可以获得更好的压缩比,同时减少 Broker 磁盘占用与网络字节数。

Kafka 设计文档强调:Record 可以成批分组、压缩后发送,Broker 校验 Batch 后仍可以把压缩后的 Batch 写入 Log,并继续以压缩形式传给 Consumer。2

但这里存在明确权衡:

Batch 更充分
→ 吞吐与压缩率通常更好
→ 可能等待更多数据
→ 单条消息端到端延迟可能增加

本篇只解释存储与协议为什么以 Batch 为边界;batch.sizelinger.msRecordAccumulatorProducerBatch 如何真正形成 Batch,将在第三篇沿 Producer 执行链展开。

2.6 CRC 校验边界#

RecordBatch Header 中包含 CRC。它主要用于发现:

  • 网络传输损坏;
  • 磁盘数据损坏;
  • 不完整或非法 Batch;
  • 结构字段与数据不一致。

Kafka 的 CORRUPT_MESSAGE 错误也明确包含 CRC 失败等消息损坏场景。13

这里必须区分:

CRC 完整性校验
副本一致性保证

CRC 回答的是“这批字节是否被损坏”;副本、ISR、Leader Election 和 ACK 回答的是“在机器故障后哪些数据仍然算已提交”。后者属于第五篇。

2.7 RecordBatch 为什么是 Kafka 性能的核心#

RecordBatch 同时贯穿:

Producer 聚合
ProduceRequest
Broker 校验
Segment Append
磁盘保存
FetchResponse
Consumer 解压和迭代

因此 Kafka 的批处理不是 Producer 端的局部优化,而是一条跨客户端、协议、Broker、磁盘和 Consumer 的端到端数据边界。


3. LogSegment 与磁盘文件组织#

3.1 一个 Segment 不是一个文件#

一个 Partition 目录可能是:

order-events-3/
├── 00000000000000000000.log
├── 00000000000000000000.index
├── 00000000000000000000.timeindex
├── 00000000000000000000.txnindex
├── 00000000000000320000.log
├── 00000000000000320000.index
├── 00000000000000320000.timeindex
├── 00000000000000320000.txnindex
├── 00000000000000640000.log
├── 00000000000000640000.index
├── 00000000000000640000.timeindex
└── 00000000000000640000.txnindex

拥有相同文件名前缀的一组文件构成一个 Segment:

00000000000000320000.log
00000000000000320000.index
00000000000000320000.timeindex
00000000000000320000.txnindex

它们共享:

Segment Base Offset = 32000

Partition、Segment 与文件集合

从对象模型看,可以把 LogSegment 简化理解为:

final class LogSegment {
FileRecords log;
OffsetIndex offsetIndex;
TimeIndex timeIndex;
TransactionIndex transactionIndex;
long baseOffset;
}

这不是 Kafka 原始源码,而是为了保留关键组成关系的概念伪代码。

当前 Kafka 源码中的 LogSegmentOffsetIndexTimeIndexTransactionIndex 正是这一层的重要阅读入口。10

3.2 文件名为什么是 Base Offset#

文件名:

00000000000000320000.log

不是时间戳,也不是随机 UUID,而是 Segment Base Offset 的补零表示:

Base Offset = 32000

假设 Segment 集合为:

Segment A: baseOffset = 0
Segment B: baseOffset = 32000
Segment C: baseOffset = 64000

目标 Offset=36891 时,可以先选择:

最大的 baseOffset ≤ 36891

即 Segment B。

Base Offset 同时承担:

  • Segment 排序键;
  • 文件名;
  • 绝对 Offset 与相对 Offset 的换算基准;
  • 定位某个 Offset 所属 Segment 的边界。

3.3 .log 文件保存什么#

.log 不是文本日志,也不是一行一条 JSON。

它保存二进制编码的连续 RecordBatch:

.log
┌──────────────────┐
│ RecordBatch 32000│
├──────────────────┤
│ RecordBatch 32071│
├──────────────────┤
│ RecordBatch 32146│
├──────────────────┤
│ RecordBatch 32219│
└──────────────────┘

两个相邻 Batch 的 Base Offset 不一定只相差 1,因为一个 Batch 内可能包含数十、数百甚至更多 Record。

.log 是 Segment 的主要事实数据;其他索引结构都是为了加速对它的读取。

3.4 .index:Offset 到 Physical Position#

Offset Index 的核心映射是:

Relative Offset → Physical Position

假设当前 Segment Base Offset=32000:

Relative OffsetAbsolute Offset.log Physical Position
0320000
128321284096
275322758320
4183241812480

绝对 Offset 与相对 Offset 的关系:

relativeOffset = absoluteOffset - segmentBaseOffset
absoluteOffset = segmentBaseOffset + relativeOffset

例如:

Target Offset = 36891
Base Offset = 32000
Relative = 4891

Offset Index 不保存业务 Value,也不一定为每条 Record 保存一项。它只提供一个接近目标的物理文件起点。

3.5 为什么使用 Relative Offset#

Partition Offset 使用长整型,是全 Partition 范围内持续增长的逻辑位置。

但一个 Segment 内只需要表示:

目标 Offset 离 Base Offset 有多远

使用 Relative Offset 可以采用更紧凑的索引条目,减少索引文件体积和内存映射压力。

同时,这也说明 Segment 不能无限增长:当相对 Offset 的可表示范围、索引容量或其他边界接近限制时,需要 Roll 出新的 Segment。

3.6 .timeindex:Timestamp 到 Relative Offset#

Time Index 的核心映射是:

Timestamp → Relative Offset

而不是:

Timestamp → `.log` Physical Position

因此按时间查找通常需要两次映射:

Timestamp
TimeIndex
Approximate Offset
OffsetIndex
Physical Position
.log Scan

Kafka Topic 配置文档还说明,时间索引只有在时间戳大于上一个已索引时间戳时才插入对应项。3

3.7 .txnindex:事务读取的辅助结构#

Transaction Index 与被中止事务范围有关,用于事务隔离读取时识别不应向 read_committed Consumer 暴露的事务数据。

本篇只建立它的文件职责:

.txnindex = 事务相关的派生辅助索引

以下内容留到第五篇:

Transaction Coordinator
Commit / Abort Marker
Last Stable Offset
Aborted Transaction
read_committed

3.8 Primary Data 与 Derived Index#

需要建立一个非常重要的存储系统认知:

.log
=
Primary Data / Source of Truth
.index / .timeindex
=
Derived Acceleration Structures

索引的职责是加速查找;如果索引损坏,系统可以通过扫描 .log 中的 RecordBatch 重建索引。

因此,Kafka 没有把所有信息都不可恢复地压进索引,而是把事实数据和查询加速结构分开。


4. OffsetIndex 与 TimeIndex:Kafka 如何定位数据#

这一章完整回答开篇问题:

Consumer 请求 Offset=36891 时,Broker 如何找到对应 RecordBatch?

4.1 第一阶段:定位 Segment#

假设 Partition 的 Segment Base Offset 为:

0
32000
64000
96000

目标:

Offset = 36891

Kafka 需要找到:

最大的 Base Offset ≤ 36891

即:

floorEntry(36891) = 32000

因此选中:

00000000000000320000.log

这一步把搜索范围从整个 Partition 缩小到了一个 Segment。

4.2 第二阶段:计算 Relative Offset#

Target Offset = 36891
Segment Base = 32000
Relative Offset
= 36891 - 32000
= 4891

OffsetIndex 在 Segment 内保存的是 Relative Offset,因此需要先完成这一坐标转换。

4.3 第三阶段:查询稀疏 OffsetIndex#

假设 .index 中存在:

Relative OffsetPhysical Position
00
102438720
204877410
4096153800
5120192410

目标 Relative Offset=4891。

不能选择 5120,因为它可能已经越过目标位置;应寻找:

最大的 indexedOffset ≤ 4891

因此得到:

Relative Offset = 4096
Physical Position = 153800

4.4 第四阶段:从附近 Position 顺序扫描#

索引并没有精确保存 Offset=4891 的字节位置。

Kafka 只知道:

.log 的 Position 153800 开始扫描,目标 RecordBatch 应该位于其后不远处。

于是顺序解析 RecordBatch:

Batch Base Offset = 36096
Batch Base Offset = 36320
Batch Base Offset = 36672
Batch Base Offset = 36864

如果最后一个 Batch 的范围覆盖 36891,就完成定位。

Offset 四阶段定位

完整过程:

Offset 36891
Segment Floor Lookup
Segment Base Offset 32000
Relative Offset 4891
Sparse OffsetIndex
Nearest Position 153800
Sequential RecordBatch Scan
Target RecordBatch

4.5 为什么不是稠密索引#

如果每条 Record 都维护:

Offset 0 → Position 0
Offset 1 → Position 125
Offset 2 → Position 284
...

查找会更精确,但代价也非常明确:

  • 索引条目数量接近 Record 数量;
  • 索引文件显著膨胀;
  • 写入时需要更频繁更新索引;
  • Page Cache 和地址空间消耗增加;
  • 恢复与重建成本上升。

Kafka 选择:

小型稀疏索引
+
邻近定位
+
短距离顺序扫描

本质上是在权衡:

Index Size ↔ Scan Distance

由于 .log 本身按照 Offset 顺序保存 RecordBatch,在找到邻近物理位置后做有限顺序扫描非常自然。

因此,稀疏索引不是“查询不够快的低配方案”,而是充分利用有序 Log 特征后的工程选择。

4.6 index.interval.bytes 不是“每隔多少条消息”#

Kafka Topic 配置中的 index.interval.bytes 控制大约累计写入多少消息字节后添加一条 Offset Index,并在满足时间条件时添加 Time Index。Kafka 4.3 默认大约每 4096 字节创建索引项。3

正确理解:

每累计约 N 字节数据
创建一个稀疏索引点

错误理解:

每 N 条消息创建一个索引点

因为不同 Record 和 Batch 大小并不固定。

索引更密:

查询更接近目标
扫描更少
索引文件更大
写索引更频繁

索引更稀:

索引更小
查询后需要多扫描一点数据

通常没有充分基准测试和明确瓶颈时,不应随意修改这一参数。

4.7 TimeIndex 为什么还要经过 OffsetIndex#

假设业务希望:

查找时间戳大于等于 10:30:00 的第一条 Record。

读取链是:

Target Timestamp
TimeIndex Binary Search
Approximate Relative Offset
Absolute Offset
OffsetIndex Lookup
Physical Position
.log Sequential Scan
Target Record

TimeIndex 二级定位

TimeIndex 解决的是:

时间坐标 → Offset 坐标

OffsetIndex 解决的是:

Offset 坐标 → 文件字节坐标

两个索引不是重复设计,而是连接三种不同坐标系:

Timestamp
Offset
Physical Position

4.8 mmap 在索引层的意义#

Kafka 的索引文件适合:

  • 条目固定、体积相对可控;
  • 按位置随机访问;
  • 频繁二分查找;
  • 不希望完整复制进 JVM Heap。

索引实现会使用内存映射文件相关机制,让操作系统按页加载需要访问的区域。

需要避免两个极端误解:

错误一:

用了 mmap,索引就一定永久全部驻留物理内存。

实际页面仍由 OS 根据内存压力调度。

错误二:

mmap 等于 Kafka 自己在 JVM Heap 中缓存全部索引。

内存映射文件并不等价于普通 Java Heap 对象数组。

当前源码阅读可从 OffsetIndexTimeIndex 及其共同抽象进入,重点观察查找、扩展、截断、重映射与关闭逻辑。10

4.9 本章完整算例#

最终把一次定位压缩成可脱稿复述的算例:

Target Offset = 36891
Segment Base Offsets:
0, 32000, 64000
Step 1:
选择 floor segment = 32000
Step 2:
relativeOffset = 36891 - 32000 = 4891
Step 3:
OffsetIndex 中查找 <= 4891 的最大条目
得到 4096 → position 153800
Step 4:
从 .log position 153800 顺序解析 RecordBatch
直到找到覆盖 Offset 36891 的 Batch

如果只能说“Kafka 用二分查找找消息”,说明理解还不够精确。

Kafka 实际上组合了:

Segment Floor Lookup
+
Sparse Index Lookup
+
Sequential Batch Scan

5. Broker 写入链:RecordBatch 如何进入磁盘#

本章从一个清晰边界开始:

Producer Client 已经形成 ProduceRequest

Producer 如何序列化、选择 Partition、在 RecordAccumulator 中组 Batch,以及 Sender 如何发送请求,全部属于第三篇。

本篇只追踪 ProduceRequest 进入 Broker 之后的数据路径。

5.1 Broker 写入主链#

为了建立源码地图,可以先记住一条经过压缩的主链:

SocketServer / Processor
RequestChannel
KafkaRequestHandler
KafkaApis.handleProduceRequest
ReplicaManager.appendRecords
Partition Leader
UnifiedLog.appendAsLeader
Active LogSegment.append
FileRecords.append

当前 Kafka Broker 的请求处理、Replica 管理和存储层仍可以沿 SocketServerKafkaRequestHandlerKafkaApisReplicaManagerPartitionUnifiedLogLogSegment 建立阅读路线。911

Broker Produce 写入执行链

注意:类名和核心职责相对稳定,但方法签名、异步封装、错误处理和模块位置会随版本演进。阅读实际生产版本时应固定 Tag。

5.2 为什么网络线程不直接执行完整存储写入#

Broker 接受大量客户端连接。如果网络线程在读完请求后直接完成全部校验、Leader 检查、Log Append 和等待逻辑:

某个请求处理变慢
Network Thread 长时间被占用
同一 Processor 上其他连接无法及时读写
网络层吞吐和延迟抖动

因此 Broker 将职责拆分:

Network Plane#

负责:

接受连接
读取 Socket 数据
协议帧处理
Request 入队
Response 写回

Request Handling Plane#

负责:

API 分发
授权与参数校验
调用 ReplicaManager
执行 Produce / Fetch 等业务请求逻辑

Storage Plane#

负责:

Partition Leader 检查
Log Append / Read
Segment / Index 管理
副本相关状态

RequestChannel 在网络接入和请求处理之间形成显式边界。

这类设计的价值不是“多线程越多越快”,而是把不同阻塞特征、资源模型和职责的工作隔离开。

5.3 Produce 请求进入 Broker 后校验什么#

真实源码中校验和错误分支很多,本篇只保留与存储主线直接相关的类别:

Topic / Partition 是否存在
当前 Broker 是否承载对应 Replica
当前 Replica 是否为可写 Leader
RecordBatch 格式是否合法
CRC 与长度是否一致
Batch / Record 大小是否超限
时间戳规则是否满足
幂等与事务元数据是否合法

这里不能把所有校验都归为“消息格式校验”。

例如:

  • NOT_LEADER_OR_FOLLOWER 属于 Partition 角色和 Metadata 问题;
  • CORRUPT_MESSAGE 可能来自 CRC、尺寸或 Batch 结构;
  • MESSAGE_TOO_LARGE 属于消息大小边界;
  • 幂等序列检查属于 Producer 状态验证。

本篇只关注校验之后,Batch 如何进入 Leader Log。

5.4 最终 Offset 为什么由 Leader Log 决定#

Producer 可以在客户端构造 Batch,但 Partition 的最终 Offset 必须符合 Leader Log 当前顺序。

假设 Leader 当前:

Log End Offset = 36891

新 Batch 有 100 条 Record,则可抽象为:

firstOffset = 36891
lastOffset = 36990
newLogEndOffset = 36991

简化伪代码:

AppendResult appendAsLeader(Records records) {
long firstOffset = currentLogEndOffset();
assignOffsets(records, firstOffset);
activeSegment().append(records);
updateLogEndOffset(records);
return result;
}

这段代码不是 Kafka 原始实现,只表达三个关键顺序:

  1. 根据 Leader Log 当前末尾确定新 Batch 的位置;
  2. 对 Batch 分配或校正最终 Offset;
  3. 追加后推进 Log End Offset。

为什么不能由每个 Producer 自己决定最终 Offset?

因为多个 Producer 会并发写同一 Partition。Offset 表达的是 Leader Log 中的全局追加顺序,只有掌握该 Log 写入权的一方才能确定最终顺序。

5.5 为什么只向 Active Segment 追加#

Partition 中的写入目标不是随机 Segment:

Closed Segment 0
Closed Segment 1
Closed Segment 2
Active Segment 3 ← append here

这使主要写入路径集中在:

当前文件尾部

避免:

  • 修改历史文件中间位置;
  • 在大量 Segment 间随机选择写入点;
  • 维护复杂的可变页更新;
  • 对历史记录执行原地覆盖。

Kafka 所谓“顺序写”首先来自这个数据模型:

Append-only Partition Log
+
One Active Segment
=
Tail Append Path

顺序写不是后来附加的磁盘优化,而是 Log 抽象的自然结果。

5.6 什么时候需要 Roll#

在追加前后,Broker 需要判断当前 Active Segment 是否应当关闭并创建新 Segment。

Roll 条件可归纳为:

Segment Size
Segment Age
Index Capacity
Relative Offset Range
特定管理条件

若需要 Roll:

Old Active Segment
↓ close
Closed Segment
Current Log End Offset
↓ becomes new base offset
New Active Segment

Roll 的完整生命周期放在第八章。

5.7 .log 与索引如何共同推进#

简化写入过程:

解析与校验 RecordBatch
确定 Batch Offset 范围
必要时 Roll Segment
追加 Batch 到 .log 文件尾部
达到 index.interval.bytes 时追加 OffsetIndex
满足时间索引条件时追加 TimeIndex
更新 Log End Offset 与状态

关键点:

不是每写一条 Record 就一定新增一条 OffsetIndex 或 TimeIndex。

.log 保存完整 Batch;索引按间隔采样,因此形成稀疏结构。

5.8 写入 Page Cache 不等于物理落盘#

应用程序调用文件写入后,数据通常先进入操作系统 Page Cache:

Kafka Broker
↓ write
OS Page Cache / Dirty Pages
↓ flush / writeback
Storage Device

需要区分三个时刻:

1. Java 写入调用返回
2. 数据进入内核 Page Cache
3. 脏页真正写入物理介质

它们不一定同时发生。

所以,不能看到 Broker Append 成功就简单断言:

每条消息已经立即 fsync 到磁盘硬件

但也不能反过来断言 Kafka 因而不可靠。Kafka 的可靠性模型组合了:

Leader Append
Replica Fetch
ISR
acks
min.insync.replicas
故障选主

这些属于第五篇。

本篇只需要记住:

Kafka 的高吞吐写路径通常依赖 OS 缓存和异步刷盘,而不是每条 Record 都执行一次同步物理写入。

5.9 为什么 Kafka 不在 JVM Heap 里缓存全部消息#

假设 Broker 自己维护一个巨大 Java 消息缓存:

Disk File
↓ copy
OS Page Cache
↓ copy
JVM Heap Message Cache
↓ copy
Socket Output

会产生:

  • Page Cache 与 JVM Cache 双份占用;
  • 大量长生命周期大对象;
  • GC 扫描、晋升和停顿压力;
  • 缓存淘汰逻辑重复;
  • Broker 重启后 JVM Cache 全部重建;
  • 数据仍然要经过内核网络栈。

操作系统已经非常擅长:

  • 缓存最近读写文件页;
  • 在内存压力下回收页面;
  • 预读连续文件;
  • 在文件和网络之间组织高效传输。

Kafka 设计文档明确强调利用文件系统 Page Cache,而不是维护等量的应用层内存缓存。2

这并不表示 Kafka 不使用 Heap。Broker 仍需要大量对象管理请求、Metadata、索引状态、网络和控制逻辑。真正的区别是:

Kafka 不试图把完整消息数据集再次复制成一套巨大 JVM Heap Cache。

5.10 Broker 写入链伪代码#

把真实代码压缩为学习伪代码:

ProduceResponse handleProduceRequest(ProduceRequest request) {
Map<TopicPartition, Records> batches = request.partitionRecords();
for (var entry : batches.entrySet()) {
Partition partition = replicaManager.requirePartition(entry.getKey());
partition.appendRecordsToLeader(entry.getValue());
}
return buildResponse();
}

继续向下:

AppendResult appendRecordsToLeader(Records records) {
ensureCurrentReplicaIsLeader();
return unifiedLog.appendAsLeader(records);
}

存储主路径:

AppendResult appendAsLeader(Records records) {
validate(records);
assignOffsets(records, logEndOffset());
if (shouldRoll(records)) {
roll();
}
activeSegment().append(records);
updateLogState(records);
return appendResult();
}

阅读源码时,应该不断问:

输入是什么?
输出是什么?
这个对象持有什么状态?
线程边界在哪里?
异常由谁转换为协议错误?
什么时候只完成本地 Append?
什么时候等待副本 ACK?

最后一个问题涉及延迟请求和可靠性,将在第五篇再展开。


6. Broker 读取链:从 Fetch Offset 到网络响应#

本章只分析 Broker 侧 Fetch。

Consumer Client 的 poll()、Fetch Session、Group Coordinator、Rebalance、Position 与 Committed Offset 关系将在第四篇展开。

6.1 Consumer 请求的不是“一条消息”#

Consumer 可能发送:

TopicPartition = order-events-3
Fetch Offset = 36891
Max Bytes = 1 MiB
Isolation = read_uncommitted / read_committed

这里的 Offset 表达:

从这个位置开始读取。

Broker 通常返回满足协议和字节限制的一批连续 RecordBatch:

36891
36892
36893
...

不是像数据库主键查询一样,只返回 Offset=36891 对应的一条 Record。

Kafka Consumer 配置文档明确说明 Consumer 以 Batch 形式 Fetch Record;即使第一个非空 Partition 的首个 Batch 超过配置限制,Broker 也可能返回它,以保证 Consumer 能继续前进。5

6.2 Broker Fetch 主链#

可以先建立以下主路径:

FetchRequest
KafkaApis
ReplicaManager.fetchMessages
Partition.readRecords
UnifiedLog.read
Locate Segment
OffsetIndex Lookup
LogSegment.read
FileRecords
FetchResponse

Broker Fetch 读取链

这一链路和写入链形成对称:

写入读取
ProduceRequestFetchRequest
appendRecordsfetchMessages / readRecords
UnifiedLog.appendUnifiedLog.read
Active Segment AppendSegment Selection
FileRecords.appendFileRecords read / transfer

6.3 读取前需要确定可见边界#

并不是 Log 中所有已存在字节都一定能对当前 Consumer 可见。

Broker 还需要根据请求类型和隔离级别确定读取上界,例如后续文章会区分:

Log End Offset
High Watermark
Last Stable Offset

本篇只保留抽象:

requestedOffset
validate range
determine max readable offset
read within visible boundary

如果当前请求 Offset 不在 Broker 保留的范围内,可能产生 OFFSET_OUT_OF_RANGE13

不要在本篇提前展开 HW、LSO 与事务可见性的精确计算,它们属于第四、五篇。

6.4 Segment 与 Index 再次参与#

UnifiedLog.read 收到目标 Offset 后,复用第四章的定位过程:

Target Offset
Segment Floor Lookup
Relative Offset
OffsetIndex Lookup
Physical Position
LogSegment / FileRecords Read

读取返回的通常是从目标位置开始、受字节上限和可见边界约束的一段连续 Records。

因此,索引层并不负责“反序列化业务对象”。它只负责找到二进制 RecordBatch 在文件中的范围。

6.5 热数据 Tail Read#

Kafka 最典型的实时流场景是:

Producer 正在写 Partition 尾部
Consumer 也正在追赶 Partition 尾部

写入刚刚进入 Page Cache 的数据页,可能立即被 Consumer Fetch 复用:

Producer Append
OS Page Cache
Consumer Tail Fetch

这时不需要每次重新从物理介质读取相同数据。

Kafka 4.3 Tiered Storage 文档明确描述:流式消费大多属于 Tail Read,利用 OS Page Cache 提供数据;较旧数据通常在回填或故障恢复时才更频繁地从磁盘读取。7

6.6 冷数据历史读取#

如果 Consumer:

  • 长时间离线后重新追赶;
  • 执行大规模历史 Replay;
  • 进行数据回填;
  • 读取已经被 Page Cache 淘汰的 Segment;

则可能发生:

Fetch
Page Cache Miss
Storage Read / Page Fault
Load Pages into Page Cache
Return Data

读取延迟会更明显地受:

  • 磁盘介质;
  • I/O 队列;
  • Segment 位置;
  • 其他工作负载;
  • 远程 Tier;
  • 网络带宽;
  • 文件系统状态;

影响。

所以:

Kafka 快不意味着读取任意历史 Offset 都等价于纯内存读取。

必须区分:

Hot Tail Read
Cold Historical Read
Remote Tier Read

6.7 传统用户态数据复制路径#

一个传统的读取再发送模型可能是:

Storage Device
Kernel Page Cache
↓ copy
Application User Buffer / JVM Buffer
↓ copy
Kernel Socket Buffer
NIC

其中应用程序把文件数据读入自己的用户态 Buffer,然后再写入 Socket。

代价包括:

  • 数据进入用户态再返回内核态;
  • 大型 Buffer 管理;
  • JVM Heap 或 Direct Buffer 压力;
  • 更多 CPU Copy;
  • 更多上下文切换和系统调用开销。

6.8 高效文件传输路径#

更高效的目标是:

File / Page Cache
Kernel Network Path
NIC

应用程序主要描述:

哪个文件
从什么 Position
传输多少字节

而尽量避免把完整 Batch 再复制到大型 JVM Heap Buffer 中转。

Kafka 设计文档将 Page Cache 与 Zero-Copy 优化联系起来:数据可以在 Page Cache 中复用,并在消费时减少应用层复制,使消费吞吐更接近网络上限。2

传统复制与高效文件传输路径

6.9 “Zero Copy”不是“物理世界零复制”#

这是 Kafka 面试和技术文章中最容易被过度简化的概念。

更准确的表述是:

Zero Copy 优化的核心是减少不必要的用户态数据复制和 JVM 中转。

实际硬件与内核路径仍可能发生:

  • Storage Device 与内存之间的 DMA;
  • 内核页面与 NIC 队列之间的数据描述或搬运;
  • TLS 加密 Buffer 处理;
  • 网络设备内部复制;
  • 特定操作系统和文件系统的实现差异。

所以不要写:

Kafka 数据从磁盘直接飞到网卡,全程一个字节都不复制。

更好的回答是:

Kafka 借助操作系统的文件传输能力,
尽量绕过“文件 → JVM 用户态大 Buffer → Socket”的重复中转。

6.10 哪些情况会破坏理想路径#

Kafka 的高效读取路径有明确边界:

TLS / SSL#

数据需要加密,可能无法完全沿用理想的纯文件到网络路径。

消息格式转换#

旧格式兼容或需要转换时,Broker 可能必须解析、重建 Batch。

Broker 重新压缩#

某些转换场景会增加 CPU 和 Buffer 工作。

Remote Tier#

数据需要先从远程对象存储或其他远程层取回,再进入本地返回路径。

Page Cache Miss#

历史数据必须实际从磁盘加载。

Fetch 太小#

大量小 Fetch 会削弱批量传输和协议摊销效果。

随机历史读取#

访问模式偏离连续 Tail Read 时,预读和缓存局部性下降。

专家级结论不应该是:

Kafka 永远零拷贝,所以一定快。

而应该是:

Kafka 的消息格式和存储布局为高效文件传输创造了条件,但实际路径仍受协议、加密、数据冷热、格式兼容和存储层影响。

6.11 Broker 读取链伪代码#

Records read(TopicPartition tp, long offset, int maxBytes) {
Partition partition = replicaManager.requirePartition(tp);
return partition.readRecords(offset, maxBytes);
}

存储层:

Records read(long offset, int maxBytes) {
long maxReadableOffset = determineVisibleBoundary();
ensureOffsetInRange(offset, maxReadableOffset);
LogSegment segment = segments.floorSegment(offset);
OffsetPosition start = segment.offsetIndex().lookup(offset);
return segment.read(
offset,
start.position(),
maxBytes,
maxReadableOffset
);
}

仍然是概念伪代码,目的在于保留核心决策顺序。

阅读源码时应追问:

Segment 在哪里被选择?
读取上界在哪里确定?
返回的是 Records 还是业务对象?
什么时候只返回部分 Batch?
什么时候触发延迟 Fetch?
数据如何进入 Network Send?

7. Kafka 为什么快:不是一个神奇优化点#

前面已经分别研究了数据格式、磁盘组织、索引、Broker 写入和读取路径。现在才能回答:

Kafka 为什么能获得高吞吐?

7.1 Log 模型把主要写入变成 Tail Append#

关系型数据库的某些更新需要:

查找数据页
修改页中记录
维护多个索引
处理随机位置更新

Kafka 的 Partition Log 主要做:

确定当前 Active Segment
向文件尾部追加 RecordBatch

所以,顺序访问首先来自数据模型,而不是一句独立调优技巧。

Append-only Log
+
Active Segment
=
集中尾部写入

7.2 Batch 贯穿整个 Data Path#

Business Records
RecordBatch
ProduceRequest
Broker Append
Segment File
FetchResponse
Consumer

批处理可以摊薄:

协议头
请求次数
系统调用
校验边界
压缩上下文
磁盘追加操作
网络发送调度

Kafka 的性能优势不是“Broker 收到单条消息后再偷偷拼 Batch”,而是客户端、协议和存储格式从设计上共享批处理边界。1

7.3 压缩后的 Batch 尽量保持压缩#

理想路径:

Producer Compression
Network
Broker Validation
Compressed Log Storage
Network
Consumer Decompression

避免每一层都:

解压 → 重建 → 再压缩

这同时减少:

  • 网络流量;
  • 磁盘空间;
  • Broker CPU 重复压缩;
  • 用户态数据重建。

但要保留边界:Broker 仍需要进行必要校验,特殊兼容和转换场景可能改变路径。

7.4 Page Cache 让写后即读非常高效#

Producer Append
Page Cache
Consumer Tail Read

实时流场景天然具有较好的时间局部性:刚刚被写入的数据很快被下游读取。

Kafka 不需要额外复制一份完整消息缓存,就能让 OS 缓存热数据。

7.5 稀疏索引避免“索引比数据还重”#

Kafka 不为每条 Record 维护复杂索引树,而是:

Segment Base Offset
+
Sparse OffsetIndex
+
Short Sequential Scan

这使索引文件相对紧凑,查询仍然能迅速接近目标位置。

7.6 减少用户态数据复制#

读取路径尽量复用:

File / Page Cache → Kernel Network Path

而不是:

File → JVM 大 Buffer → Socket

这降低 CPU 和内存中转成本。

7.7 Partition 提供水平并行#

单 Partition 保持局部顺序:

P0: ordered append and fetch

多个 Partition 可以分布到多个 Broker:

P0 → Broker 1
P1 → Broker 2
P2 → Broker 3
P3 → Broker 4

形成:

  • 并行写入;
  • 并行读取;
  • 存储容量扩展;
  • 网络吞吐扩展。

Partition 数量和分布不是越多越好,容量设计与治理放在第七篇。

7.8 一个概念性能模型#

可以用以下关系帮助记忆:

Kafka Throughput
Partition Parallelism
× Batch Efficiency
× Sequential Access Efficiency
× Cache Hit Efficiency
× Network Transfer Efficiency

它不是可以直接代入生产容量的数学公式,只表达:

Kafka 吞吐取决于多层设计共同工作,而不是单点优化。

7.9 正确回答“Kafka 为什么快”#

不合格回答:

因为 Kafka 顺序写、Page Cache、零拷贝。

更完整的回答:

Kafka 先用 Append-only Partition Log 把主写路径收束为 Active Segment 尾部追加,再用 RecordBatch 让压缩、网络协议、CRC、磁盘写入和 Fetch 共享批处理边界;Segment 与稀疏索引在可管理文件大小和查询效率之间平衡;消息数据主要复用 OS Page Cache,并在读取时尽量减少 JVM 用户态中转;最后通过多个 Partition 将存储和吞吐水平分布到多个 Broker。

可以进一步收束为:

Log Data Model
+
Segmented Storage
+
Batch-oriented Format
+
Sparse Index
+
OS Page Cache
+
Efficient File Transfer
+
Partition Parallelism

8. Segment Roll、Retention 与 Compaction#

8.1 为什么需要 Roll#

Active Segment 不能无限增长。

达到一定条件后:

Active Segment
↓ Roll
Closed Segment
Current Log End Offset
New Segment Base Offset
New Active Segment

LogSegment 生命周期

8.2 Roll 条件#

Roll 判断可从以下维度理解:

Segment Size#

segment.bytes

当 Segment 文件达到配置大小附近时,创建新 Segment。

Segment Age#

segment.ms
segment.jitter.ms

即使流量不大,也需要在时间到达后 Roll,避免一个低流量 Partition 长期不产生可独立清理的 Closed Segment。

Index Capacity#

索引文件达到容量或不能继续安全追加时,需要 Roll。

Relative Offset Range#

Segment 内使用相对 Offset 表示索引,不能让单 Segment 跨度无限增长。

管理与恢复条件#

实际实现还会考虑恢复、快照、文件状态和其他边界。

8.3 Roll 后发生什么#

概念链:

停止向旧 Segment 追加
旧 Segment 成为 Closed
以当前 Log End Offset 创建新文件组
新 Segment 成为 Active
后续 RecordBatch 继续追加

例如:

Old active baseOffset = 32000
Current LEO = 64000
New files:
00000000000000640000.log
00000000000000640000.index
00000000000000640000.timeindex
00000000000000640000.txnindex

8.4 Retention Delete 以 Segment 为主要边界#

Kafka Log 保留不是:

Consumer ACK 一条,就立刻删除一条。

而是:

Closed Segments
达到 retention 时间或大小条件
标记并删除整个 Segment 文件集合

cleanup.policy=delete 是默认清理策略,可按照保留时间或总大小清理旧 Segment。3

这种设计再次体现 Queue 与 Log 的差异:

消费进度
数据保留生命周期

Consumer 已经读过的数据,仍可在 Retention 窗口内供其他 Consumer 或 Replay 使用。

8.5 为什么 Active Segment 通常不能直接按 Retention 删除#

Active Segment 仍在持续追加,其文件和索引边界尚未稳定。

如果直接删除其中“过期的前半部分”,又会回到修改大文件内部的问题。

因此 Kafka 通过 Roll 产生 Closed Segment,再以 Segment 为单位进行清理。

这也说明:

segment.ms 不仅影响文件大小,还会间接影响低流量 Topic 旧数据形成可清理 Segment 的及时性。

8.6 Log Compaction 不是按时间删除#

假设:

Offset 0: Key=A Value=1
Offset 1: Key=B Value=3
Offset 2: Key=A Value=2
Offset 3: Key=A Value=5

Compaction 目标是保留每个 Key 的最新值:

Key=B Value=3
Key=A Value=5

cleanup.policy=compact 用于保留每个 Key 的最后已知值,也可以和 delete 组合。3

必须避免以下错误:

Compaction = 删除超过 N 天的数据

错误。那是 Retention Delete 的主要语义。

Compaction 后 Offset 会重新连续编号

错误。Offset 作为 Log 中的位置标识不会因为中间 Record 被压缩清理而重新编号。

可能出现:

Offset 100 存在
Offset 101 已被压缩移除
Offset 102 存在

Consumer 从 101 开始 Fetch 时,可以继续前进到后续有效 Batch,而不是要求所有 Offset 对应一条仍存在的业务 Record。

8.7 Tiered Storage#

Tiered Storage 将数据分成:

Local Tier
+
Remote Tier

当前尾部和热 Segment 保留在 Broker 本地;完成的冷 Segment 可以迁移到远程存储。Kafka 4.3 官方文档将本地层描述为现有 Broker Local Disk,将远程层描述为 HDFS、S3 等外部存储实现所承载的 Completed Log Segments。7

本篇只需要建立:

Segment 也是远程迁移的重要边界

RemoteLogManager、远程 Metadata Topic、Follower 从远程层恢复和生产参数治理将在第七篇生产架构中再讨论。


9. Kafka 存储源码阅读地图#

源码阅读不是从某个 3000 行类第一行硬读到最后一行。

正确方式是先建立数据结构,再沿一条请求链向上下游扩展。

9.1 第一层:消息格式#

优先理解:

Records
RecordBatch
DefaultRecordBatch
Record
DefaultRecord
MemoryRecords
FileRecords

重点问题:

  1. Batch Header 如何解析?
  2. baseOffsetoffsetDelta 在哪里组合?
  3. Compression Type 存在哪里?
  4. CRC 覆盖哪些区域?
  5. FileRecords 如何表达文件中的 Records 范围?
  6. Batch 如何写入或传输到 NIO Channel?

Kafka 4.3 官方 Message Format 文档应当和源码一起阅读,而不是只看类名猜格式。1

9.2 第二层:索引#

阅读:

OffsetIndex
TimeIndex
TransactionIndex
AbstractIndex

重点问题:

  1. 为什么 OffsetIndex 保存 Relative Offset?
  2. lookup() 返回的是最终 Record 还是邻近 Position?
  3. 索引如何做二分查找?
  4. Index File 如何扩展、截断和重新映射?
  5. TimeIndex 为什么保存 Timestamp 到 Relative Offset?
  6. Index 损坏后如何依赖 .log 重建?

9.3 第三层:Segment#

阅读:

LogSegment

重点问题:

  1. 一个 Segment 持有哪些 File 与 Index?
  2. append 如何同时推进数据和索引?
  3. read 如何使用开始 Position?
  4. recover 如何扫描 RecordBatch?
  5. shouldRoll 的条件来自哪里?
  6. Segment 关闭、截断和删除如何协调文件?

9.4 第四层:逻辑 Log#

阅读:

UnifiedLog
LocalLog
LogSegments
LogManager

重点问题:

  1. Partition 的 Segment 集合如何管理?
  2. floorSegment(offset) 如何实现?
  3. Active Segment 如何取得?
  4. appendAsLeader 的 Offset 分配和验证顺序是什么?
  5. read 如何确定可见上界?
  6. Roll 如何创建新 Segment?
  7. Log End Offset、恢复点和快照如何维护?

当前 4.x 的 UnifiedLogLocalLogLogSegment 已位于存储模块,阅读旧教程时要警惕类所在模块和实现语言已经发生变化。10

9.5 第五层:Partition 与 ReplicaManager#

阅读:

Partition
ReplicaManager

重点问题:

  1. TopicPartition 如何映射到本地 Partition?
  2. 如何判断当前 Replica 是 Leader 还是 Follower?
  3. Produce 如何进入 appendRecordsToLeader
  4. Fetch 如何进入本地 Log?
  5. 哪些请求要等待副本条件满足?
  6. Delayed Produce / Delayed Fetch 如何与主链配合?

第五和第六个问题牵涉可靠性与 Purgatory,在第五篇展开。

9.6 第六层:Broker 请求入口#

阅读:

SocketServer
KafkaRequestHandler
KafkaApis
RequestChannel

重点问题:

  1. Socket 数据如何变成 Request?
  2. Processor 与 Handler Thread 如何分工?
  3. Produce、Fetch 在哪个 API Handler 分发?
  4. 错误如何转成协议 Response?
  5. Records 如何进入 Network Send?

9.7 推荐顺序#

Message Format 文档
DefaultRecordBatch / FileRecords
OffsetIndex / TimeIndex
LogSegment
UnifiedLog / LocalLog
Partition
ReplicaManager
KafkaApis
SocketServer

先从底层数据结构向调用者反推,比从网络入口一路陷入所有分支更容易形成稳定模型。

9.8 源码阅读规范#

建议每条链都使用固定模板记录:

入口方法
输入对象
关键状态
线程边界
核心数据结构变化
输出对象
错误边界

例如 OffsetIndex.lookup(targetOffset)

输入:Absolute / Relative Offset?
状态:Base Offset、Mapped Index File
动作:Binary Search
输出:OffsetPosition
边界:结果是否精确等于目标?
下游:LogSegment 从该 Position 顺序扫描

这种阅读方式比复制源码注释更接近工程专家需要的能力。


10. 可运行实验:亲手观察 Kafka 磁盘结构#

以下实验用于把抽象结构落到真实文件。

环境定位:单节点、单副本、本地学习。该配置不代表生产集群建议。

10.1 实验目录#

kafka-storage-lab/
├── docker-compose.yml
└── data/

10.2 Docker Compose#

services:
kafka:
image: apache/kafka:4.3.0
container_name: kafka-ch02
hostname: kafka
ports:
- "9092:9092"
environment:
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:9093
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
KAFKA_LOG_DIRS: /var/lib/kafka/data
volumes:
- ./data:/var/lib/kafka/data

启动:

Terminal window
docker compose up -d
docker logs -f kafka-ch02

看到 Broker 启动完成后,再执行后续命令。

不同 Docker Desktop、文件权限和镜像小版本可能需要调整挂载目录权限。若本地挂载失败,可以先改用 Docker Named Volume,确认 Kafka 正常运行后再观察容器内部文件。

10.3 创建容易发生 Roll 的 Topic#

Terminal window
docker exec kafka-ch02 \
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 \
--create \
--topic order-events \
--partitions 1 \
--replication-factor 1 \
--config segment.bytes=1048576 \
--config index.interval.bytes=4096 \
--config retention.ms=86400000

这里故意将:

segment.bytes = 1 MiB

设得很小,以便少量测试数据就生成多个 Segment。

验证:

Terminal window
docker exec kafka-ch02 \
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 \
--describe \
--topic order-events

查看 Topic 配置:

Terminal window
docker exec kafka-ch02 \
/opt/kafka/bin/kafka-configs.sh \
--bootstrap-server localhost:9092 \
--entity-type topics \
--entity-name order-events \
--describe

10.4 批量发送测试数据#

Linux / macOS / WSL:

Terminal window
python3 - <<'PY' | docker exec -i kafka-ch02 \
/opt/kafka/bin/kafka-console-producer.sh \
--bootstrap-server localhost:9092 \
--topic order-events
import json
for i in range(30000):
print(json.dumps({
"eventId": i,
"orderId": f"ORD-{i:08d}",
"eventType": "OrderCreated",
"payload": "x" * 180
}))
PY

如果 30,000 条仍不足以产生多个 Segment,可继续增加消息数量或 Payload 大小。

10.5 观察 Partition 目录#

Terminal window
docker exec kafka-ch02 sh -lc \
'ls -lh /var/lib/kafka/data/order-events-0'

预期看到多组:

*.log
*.index
*.timeindex

并观察:

  1. 文件名前缀是否递增;
  2. .log 是否接近配置的 1 MiB 后生成新文件;
  3. 最后一个 Active Segment 是否可能明显小于前面 Segment;
  4. .index 是否远小于 .log
  5. 多个文件是否共享相同 Base Offset 前缀。

如果目录名称或数据路径与示例不同,先执行:

Terminal window
docker exec kafka-ch02 \
/opt/kafka/bin/kafka-log-dirs.sh \
--bootstrap-server localhost:9092 \
--describe \
--topic-list order-events

再根据实际路径检查。

10.6 使用 kafka-dump-log.sh 查看 RecordBatch#

先找到某个 .log 文件:

Terminal window
docker exec kafka-ch02 sh -lc \
'find /var/lib/kafka/data/order-events-0 -name "*.log" | sort | head -1'

假设输出:

/var/lib/kafka/data/order-events-0/00000000000000000000.log

执行:

Terminal window
docker exec kafka-ch02 \
/opt/kafka/bin/kafka-dump-log.sh \
--files /var/lib/kafka/data/order-events-0/00000000000000000000.log \
--print-data-log

观察:

baseOffset
lastOffset
position
CreateTime
isvalid
size
producerId
compression codec

不同 Kafka 小版本的输出字段和工具参数可能略有差异,以容器内:

Terminal window
/opt/kafka/bin/kafka-dump-log.sh --help

为准。

实验目标不是记住工具输出,而是验证:

磁盘里保存的是 RecordBatch
每个 Batch 有 Base Offset 和文件 Position
Batch 可以包含多条 Record

10.7 查看索引#

对 Offset Index 做 sanity check:

Terminal window
docker exec kafka-ch02 \
/opt/kafka/bin/kafka-dump-log.sh \
--files /var/lib/kafka/data/order-events-0/00000000000000000000.index \
--index-sanity-check

也可以根据当前工具帮助信息,查看 Offset / Position 条目。

重点验证:

OffsetIndex 文件明显小于 .log
索引条目不是每条 Record 一项
Position 指向 .log 中的字节位置

10.8 手工完成一次 Offset 定位#

选择某个目标 Offset,例如:

targetOffset = 12000

执行步骤:

  1. 列出所有 .log 文件前缀;
  2. 选择最大且不超过 12000 的 Base Offset;
  3. 计算 Relative Offset;
  4. 使用 Dump Log 工具查看该 Segment 中 Batch 的 Base Offset 与 Position;
  5. 找到不大于目标 Offset 的邻近 Batch;
  6. 验证目标是否落在该 Batch Offset 范围内。

把结果记录成:

Target Offset:
Selected Segment Base Offset:
Relative Offset:
Nearest Indexed Offset:
Physical Position:
Containing RecordBatch:

这是本篇最重要的实践。

10.9 观察 Retention 与 Roll#

可临时创建更短保留时间的测试 Topic:

Terminal window
docker exec kafka-ch02 \
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 \
--create \
--topic retention-lab \
--partitions 1 \
--replication-factor 1 \
--config segment.bytes=262144 \
--config retention.ms=120000 \
--config file.delete.delay.ms=10000

发送足够数据形成多个 Closed Segment,再定时观察目录变化。

需要理解:

retention.ms 到期
对应 Record 在该毫秒立即消失

清理还受:

  • Segment 是否 Closed;
  • Retention 检查周期;
  • 文件删除延迟;
  • Broker 调度;

影响。

10.10 清理实验#

Terminal window
docker compose down

同时删除本地测试数据:

Terminal window
rm -rf data

请确认当前目录正确,避免误删其他文件。


11. 常见错误认知与纠偏#

11.1 “Kafka 一个 Partition 就是一个 .log 文件”#

错误。

正确模型:

Partition
=
多个按 Base Offset 排序的 Segment

每个 Segment 又是一组 .log 与索引文件。

11.2 “Offset 就是文件字节偏移量”#

错误。

Offset
=
Partition 内逻辑位置
Physical Position
=
Segment `.log` 中字节位置

两者通过 OffsetIndex 建立近似映射。

11.3 “.index 保存每一条消息的位置”#

错误。

OffsetIndex 是稀疏索引,按字节间隔采样;查到邻近 Position 后仍需顺序扫描 Batch。

11.4 “TimeIndex 可以直接定位消息文件字节”#

错误。

Timestamp → Offset → Position

TimeIndex 通常先得到近似 Offset,再经过 OffsetIndex。

11.5 “Kafka 顺序写,所以完全没有随机 I/O”#

过度简化。

主要消息写入路径是 Active Segment Tail Append,但系统仍会发生:

  • 索引查找;
  • 历史 Fetch;
  • 恢复扫描;
  • Compaction;
  • Segment 删除;
  • Replica Fetch;
  • Metadata 与内部 Topic I/O;
  • 文件系统和存储设备内部行为。

更准确的说法是:

Kafka 将高频业务消息主写路径设计为顺序追加,并使常见 Tail Read 具有良好局部性。

11.6 “写入成功意味着已经同步到物理磁盘”#

错误。

应用写入、进入 Page Cache、脏页 Flush、物理介质持久化是不同阶段。

消息可靠性必须结合 ACK、Replica 和 ISR 分析。

11.7 “Kafka 不用内存”#

错误。

Kafka 大量使用内存处理请求、Metadata、网络、索引映射和状态;只是消息数据主要利用 OS Page Cache,而不是全部复制为 JVM Heap Cache。

11.8 “Zero Copy 表示绝对没有数据复制”#

错误。

它主要表示减少不必要的用户态复制和 JVM Buffer 中转。

11.9 “Batch 越大越好”#

错误。

更大的 Batch 可能提高吞吐和压缩率,但也可能:

  • 增加聚合等待;
  • 增大单次请求;
  • 增加内存占用;
  • 增加单 Batch 失败和重试成本;
  • 触发消息大小限制;
  • 增加尾延迟。

第三篇会系统分析 Batch、Latency 和 Memory 的权衡。

11.10 “Segment 越小越好,删除更灵活”#

错误。

过小 Segment 会带来:

  • 文件数量暴增;
  • 更频繁 Roll;
  • 更多文件句柄与索引;
  • Retention 和恢复管理开销;
  • Page Cache 局部性和磁盘元数据压力。

过大 Segment 又会让:

  • Retention 粒度变粗;
  • 单文件恢复和迁移成本增大;
  • 低流量 Topic 旧数据延迟形成 Closed Segment。

需要结合吞吐、保留、恢复和运维成本权衡。


12. 故障与性能排查思路#

12.1 Produce 延迟升高#

排查树:

Produce Latency ↑
├── 网络排队 / Request Queue
├── Broker CPU 或 GC
├── 磁盘写入 / Page Cache 压力
├── Segment Roll 频繁
├── Batch 校验或格式转换
├── Replica / ACK 等待
├── Hot Partition
└── 消息过大或压缩 CPU 高

本篇重点检查:

  • 是否频繁生成 Segment;
  • segment.bytes 是否被异常设小;
  • 磁盘吞吐和延迟;
  • Page Cache 与系统内存压力;
  • Request Handler 是否饱和;
  • 单 Partition 流量是否倾斜。

ACK 与 Replica 等待第五篇展开。

12.2 Fetch 延迟升高#

Fetch Latency ↑
├── Page Cache Miss
├── 历史回填 / Random Read
├── Remote Tier Read
├── 磁盘队列拥塞
├── Fetch 太小、请求过多
├── TLS / CPU
├── 消息格式转换
├── Broker 网络瓶颈
└── Hot Partition

首先区分:

实时 Tail Consumer 变慢
还是
历史 Replay Consumer 变慢

两者的根因完全可能不同。

12.3 磁盘空间快速增长#

Disk Usage ↑
├── Produce Rate 上升
├── Retention 时间过长
├── retention.bytes 未限制
├── Segment 尚未满足可删除条件
├── Compaction 跟不上
├── Key 分布导致压缩收益低
├── Replica 数量增加
├── Remote Tier 未及时迁移
└── 删除延迟或 Broker 清理异常

不要看到目录大就立即手工删除 .log 文件。

手工删除会破坏 Kafka 对 Segment、索引、Checkpoint 和 Replica 状态的一致管理。

12.4 Segment 数量异常多#

可能原因:

  • segment.bytes 过小;
  • segment.ms 过短;
  • Topic / Partition 数量过多;
  • 大量低流量 Partition 定时 Roll;
  • 测试配置误带入生产。

影响:

  • 文件句柄;
  • 文件系统元数据;
  • Broker 启动恢复时间;
  • Retention 扫描;
  • Partition 迁移和备份;
  • Page Cache 管理。

12.5 Index 损坏或不一致#

基本原则:

.log 是事实数据
.index / .timeindex 是派生结构

Kafka 启动恢复流程可以扫描 Log 重建索引,但操作前应:

  • 确认 Broker 版本;
  • 确认 Replica 与 Leader 状态;
  • 备份现场;
  • 参考官方工具和恢复流程;
  • 避免在生产环境直接试验手工修改文件。

13. 第二篇自测与验收#

完成本篇后,应当能够脱稿回答:

  1. Partition 为什么不能只使用一个无限增长文件?
  2. Partition、LogSegment、RecordBatch 和 Record 是什么关系?
  3. Active Segment 与 Closed Segment 分别承担什么职责?
  4. Segment 文件名为什么使用 Base Offset?
  5. .log.index.timeindex.txnindex 分别保存什么?
  6. Kafka Offset 与 .log Physical Position 有什么区别?
  7. 为什么 OffsetIndex 保存 Relative Offset?
  8. Kafka 为什么选择稀疏索引?
  9. index.interval.bytes 为什么不是消息条数?
  10. Offset=36891 如何经过四个阶段定位到 RecordBatch?
  11. TimeIndex 为什么还需要 OffsetIndex?
  12. 为什么索引可以从 .log 重建?
  13. RecordBatch Header 中哪些字段与本篇最相关?
  14. 为什么 Compression 以 Batch 为边界?
  15. Broker 收到 ProduceRequest 后的主写入链是什么?
  16. 为什么网络线程不直接完成全部存储工作?
  17. 最终 Partition Offset 为什么由 Leader Log 决定?
  18. 为什么只写 Active Segment?
  19. 写入 Page Cache 与物理落盘有什么区别?
  20. Kafka 为什么不把全部消息维护成 JVM Heap Cache?
  21. Broker Fetch 的主要读取链是什么?
  22. Consumer 请求一个 Offset 时为什么会返回连续 Batch?
  23. Hot Tail Read 与 Cold Historical Read 有什么区别?
  24. Zero Copy 更准确的含义是什么?
  25. TLS、格式转换和 Remote Tier 为什么可能改变理想读取路径?
  26. Segment Roll、Retention Delete 与 Log Compaction 有什么区别?
  27. 为什么 Compaction 不会重新编号 Offset?
  28. 为什么“Kafka 快是因为顺序写”是不完整答案?

最终口述验收#

你应该能够不看文章,用 20~30 分钟完整讲清:

RecordBatch
Active LogSegment Append
.log / Sparse Index
Page Cache
Fetch Offset Lookup
File Transfer

同时能够画出:

Offset
→ Segment
→ Relative Offset
→ OffsetIndex
→ Physical Position
→ RecordBatch

如果只能背出“顺序写、Page Cache、零拷贝”,本篇仍未真正掌握。


14. 全文收束:一条 RecordBatch 的完整旅程#

14.1 写入侧#

ProduceRequest
SocketServer / RequestChannel
KafkaRequestHandler
KafkaApis
ReplicaManager
Partition Leader
UnifiedLog
必要时 Segment Roll
Active LogSegment
RecordBatch Append to .log
Sparse OffsetIndex / TimeIndex Update
OS Page Cache
OS Writeback
Storage Device

14.2 读取侧#

Consumer Fetch(offset=36891)
KafkaApis / ReplicaManager
Partition / UnifiedLog
确定可见读取边界
Segment Floor Lookup
Base Offset=32000
Relative Offset=4891
OffsetIndex Nearest Entry
Physical Position=153800
Sequential RecordBatch Scan
Page Cache Hit or Storage Read
FetchResponse
Efficient Network Path
Consumer

Kafka 并不是简单地在一个普通文件上执行“顺序写”。

它将 Partition 拆成可独立管理的 Segment,把 Record 组织成贯穿协议和磁盘的 RecordBatch,使用稀疏索引在逻辑 Offset 与物理 Position 之间搭桥,再利用操作系统 Page Cache 和高效文件传输路径减少数据在磁盘、内核与 JVM 之间不必要的移动。

所以 Kafka 的高性能不是一个神奇技巧,而是一套一致的数据路径设计:

Log 决定写入形态,Batch 决定传输边界,Segment 决定管理边界,Index 决定定位方式,Page Cache 与内核路径决定数据如何移动。


15. 下一篇预告:Kafka Producer 完整执行链#

到目前为止,我们都从 ProduceRequest 已经到达 Broker 之后开始观察。

下一篇会把视角移动到 Java Client:

  • KafkaProducer.send() 为什么通常不会立即发送网络请求?
  • Serializer、Partitioner 和 Metadata 分别在什么位置工作?
  • RecordAccumulator 为什么按 TopicPartition 组织 Batch?
  • BufferPool 如何限制 Producer 内存?
  • Sender Thread 如何判断哪些 Broker Ready?
  • batch.sizelinger.ms 到底影响哪个执行阶段?
  • buffer.memory 耗尽后业务线程会发生什么?
  • max.in.flight.requests.per.connection 如何影响吞吐和顺序?
  • Retry、ACK 和 Idempotence 如何接入发送状态机?

下一篇:

《Kafka Producer:一条 Record 从 send() 到 Partition Leader 的完整旅程》


参考资料#

  1. Apache Kafka 4.3 Message Format
  2. Apache Kafka 4.3 Design
  3. Apache Kafka 4.3 Topic Configs
  4. Apache Kafka 4.3 Log Implementation
  5. Apache Kafka 4.3 Consumer Configs
  6. Apache Kafka 4.3 Broker Configs
  7. Apache Kafka 4.3 Tiered Storage
  8. Apache Kafka 4.3 Messages Implementation
  9. Apache Kafka KafkaApis Source
  10. Apache Kafka Storage Log Source Package
  11. Apache Kafka ReplicaManager Source
  12. Apache Kafka Developer Guide
  13. Apache Kafka Protocol Errors Source
Kafka 存储引擎与高性能原理:一条 Record 如何落盘并被读取
https://jupiter-ws.cn/posts/backend/kafka/02_kafka_storage_engine_high_performance/
作者
Jupiter
发布于
2026-07-17
许可协议
CC BY-NC-SA 4.0