Kafka 存储引擎与高性能原理:一条 Record 如何落盘并被读取
从 RecordBatch、LogSegment、稀疏索引到 Page Cache,沿完整数据路径解释 Kafka 为什么快
阅读目标
第一篇建立了 Kafka 的核心世界观:Topic 由多个 Partition 组成,每个 Partition 是一条有序、可持久化、可按 Offset 读取的 Log。
但这仍然只是逻辑模型。
第二篇要把一个 Partition 放大到磁盘和内核层,回答两个真正决定 Kafka 工程能力的问题:
- 一批 Record 如何从 Broker 网络请求进入 Partition Log,并最终形成磁盘文件?
- Consumer 请求 Offset=36891 时,Broker 如何找到对应 RecordBatch,并高效返回数据?
完成本篇后,你应该能够:
- 区分
Record、RecordBatch、LogSegment、Partition 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-eventsPartition = 3Offset = 36891MaxBytes = 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.timeindexBroker 至少需要连续回答以下问题:
order-events-3对应哪个本地 Partition?- Offset 36891 位于哪个 Segment?
- 目标 Offset 在 Segment 内的相对位置是多少?
.index是否精确保存 Offset 36891?- 如果索引只记录了附近位置,还需要扫描多远?
- 找到 RecordBatch 后,数据需要复制进 JVM Heap 吗?
- 如果数据刚刚被写入,它是否仍在操作系统 Page Cache 中?
- 如果数据已经很旧,读取路径又会发生什么变化?

这一串问题揭示了 Kafka 存储引擎的核心矛盾:
Consumer 使用逻辑 Offset 读取 ↓磁盘只能按文件与字节 Position 读取 ↓Kafka 必须在两种坐标系之间建立高效映射接下来,我们从第一篇中的 Partition 开始,把它逐层展开。
1. 从 Partition 到 Kafka 存储引擎
1.1 逻辑 Log 与物理文件不是同一个概念
第一篇中的 Partition 是一条逻辑上连续增长的 Log:
Partition 3
Offset0 → 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 | 数据文件与辅助索引文件集合 |
| Record | RecordBatch 内部的一条业务记录 |
| Offset | Partition 内的逻辑位置 |
| Physical Position | .log 文件中的字节位置 |
这里最容易混淆的是:
Kafka Offset ≠ 文件字节 PositionOffset=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 为以下工程操作提供了独立边界:
RollRetention DeleteCompactionRecoveryTransferTieringSegment 不是为了改变 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 1Record 2Record 3Record 4 │ ▼RecordBatch ├── 共享批头 ├── 共享 CRC 边界 ├── 共享压缩边界 ├── 连续网络传输 └── 连续磁盘追加2.2 四层数据结构必须分清
Record ↓RecordBatch ↓LogSegment ↓Partition LogRecord
一条业务记录,包含经过序列化的:
KeyValueTimestamp DeltaOffset DeltaHeadersRecordBatch
一组被统一编码、校验、压缩、传输和存储的 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[]
本篇重点关注:
baseOffset:Batch 第一条 Record 的绝对 Offset;lastOffsetDelta:Batch 最后一条 Record 相对 Base Offset 的增量;baseTimestamp:Batch 时间戳基准;attributes:压缩类型等 Batch 属性;crc:Batch 数据完整性校验;recordsCount:Batch 中 Record 数量。
而以下字段只先认识,不展开机制:
producerIdproducerEpochbaseSequenceisTransactionalisControlBatch它们与幂等 Producer、事务和 Control Batch 密切相关,将在第五篇完整展开。
2.4 Base Offset 与 Offset Delta
假设:
baseOffset = 1000Batch 内部四条 Record 可以保存:
Record A: offsetDelta = 0 → Offset 1000Record B: offsetDelta = 1 → Offset 1001Record C: offsetDelta = 2 → Offset 1002Record D: offsetDelta = 3 → Offset 1003最终 Offset 的计算是:
absoluteOffset = baseOffset + offsetDelta这样不需要每条 Record 都重复保存完整 64 位绝对 Offset。
时间戳也使用类似思路:
baseTimestamp = T₀
Record A: timestampDelta = 0Record B: timestampDelta = 12 msRecord 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.size、linger.ms、RecordAccumulator 和 ProducerBatch 如何真正形成 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.log00000000000000320000.index00000000000000320000.timeindex00000000000000320000.txnindex它们共享:
Segment Base Offset = 32000
从对象模型看,可以把 LogSegment 简化理解为:
final class LogSegment { FileRecords log; OffsetIndex offsetIndex; TimeIndex timeIndex; TransactionIndex transactionIndex; long baseOffset;}这不是 Kafka 原始源码,而是为了保留关键组成关系的概念伪代码。
当前 Kafka 源码中的 LogSegment、OffsetIndex、TimeIndex 和 TransactionIndex 正是这一层的重要阅读入口。10
3.2 文件名为什么是 Base Offset
文件名:
00000000000000320000.log不是时间戳,也不是随机 UUID,而是 Segment Base Offset 的补零表示:
Base Offset = 32000假设 Segment 集合为:
Segment A: baseOffset = 0Segment B: baseOffset = 32000Segment 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 Offset | Absolute Offset | .log Physical Position |
|---|---|---|
| 0 | 32000 | 0 |
| 128 | 32128 | 4096 |
| 275 | 32275 | 8320 |
| 418 | 32418 | 12480 |
绝对 Offset 与相对 Offset 的关系:
relativeOffset = absoluteOffset - segmentBaseOffsetabsoluteOffset = segmentBaseOffset + relativeOffset例如:
Target Offset = 36891Base Offset = 32000Relative = 4891Offset 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 ScanKafka Topic 配置文档还说明,时间索引只有在时间戳大于上一个已索引时间戳时才插入对应项。3
3.7 .txnindex:事务读取的辅助结构
Transaction Index 与被中止事务范围有关,用于事务隔离读取时识别不应向 read_committed Consumer 暴露的事务数据。
本篇只建立它的文件职责:
.txnindex = 事务相关的派生辅助索引以下内容留到第五篇:
Transaction CoordinatorCommit / Abort MarkerLast Stable OffsetAborted Transactionread_committed3.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 为:
0320006400096000目标:
Offset = 36891Kafka 需要找到:
最大的 Base Offset ≤ 36891即:
floorEntry(36891) = 32000因此选中:
00000000000000320000.log这一步把搜索范围从整个 Partition 缩小到了一个 Segment。
4.2 第二阶段:计算 Relative Offset
Target Offset = 36891Segment Base = 32000
Relative Offset= 36891 - 32000= 4891OffsetIndex 在 Segment 内保存的是 Relative Offset,因此需要先完成这一坐标转换。
4.3 第三阶段:查询稀疏 OffsetIndex
假设 .index 中存在:
| Relative Offset | Physical Position |
|---|---|
| 0 | 0 |
| 1024 | 38720 |
| 2048 | 77410 |
| 4096 | 153800 |
| 5120 | 192410 |
目标 Relative Offset=4891。
不能选择 5120,因为它可能已经越过目标位置;应寻找:
最大的 indexedOffset ≤ 4891因此得到:
Relative Offset = 4096Physical Position = 1538004.4 第四阶段:从附近 Position 顺序扫描
索引并没有精确保存 Offset=4891 的字节位置。
Kafka 只知道:
从
.log的 Position 153800 开始扫描,目标 RecordBatch 应该位于其后不远处。
于是顺序解析 RecordBatch:
Batch Base Offset = 36096Batch Base Offset = 36320Batch Base Offset = 36672Batch Base Offset = 36864如果最后一个 Batch 的范围覆盖 36891,就完成定位。

完整过程:
Offset 36891 ↓Segment Floor Lookup ↓Segment Base Offset 32000 ↓Relative Offset 4891 ↓Sparse OffsetIndex ↓Nearest Position 153800 ↓Sequential RecordBatch Scan ↓Target RecordBatch4.5 为什么不是稠密索引
如果每条 Record 都维护:
Offset 0 → Position 0Offset 1 → Position 125Offset 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 解决的是:
时间坐标 → Offset 坐标OffsetIndex 解决的是:
Offset 坐标 → 文件字节坐标两个索引不是重复设计,而是连接三种不同坐标系:
TimestampOffsetPhysical Position4.8 mmap 在索引层的意义
Kafka 的索引文件适合:
- 条目固定、体积相对可控;
- 按位置随机访问;
- 频繁二分查找;
- 不希望完整复制进 JVM Heap。
索引实现会使用内存映射文件相关机制,让操作系统按页加载需要访问的区域。
需要避免两个极端误解:
错误一:
用了 mmap,索引就一定永久全部驻留物理内存。实际页面仍由 OS 根据内存压力调度。
错误二:
mmap 等于 Kafka 自己在 JVM Heap 中缓存全部索引。内存映射文件并不等价于普通 Java Heap 对象数组。
当前源码阅读可从 OffsetIndex、TimeIndex 及其共同抽象进入,重点观察查找、扩展、截断、重映射与关闭逻辑。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 Scan5. Broker 写入链:RecordBatch 如何进入磁盘
本章从一个清晰边界开始:
Producer Client 已经形成 ProduceRequestProducer 如何序列化、选择 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 管理和存储层仍可以沿 SocketServer、KafkaRequestHandler、KafkaApis、ReplicaManager、Partition、UnifiedLog 与 LogSegment 建立阅读路线。911

注意:类名和核心职责相对稳定,但方法签名、异步封装、错误处理和模块位置会随版本演进。阅读实际生产版本时应固定 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 / ReadSegment / Index 管理副本相关状态RequestChannel 在网络接入和请求处理之间形成显式边界。
这类设计的价值不是“多线程越多越快”,而是把不同阻塞特征、资源模型和职责的工作隔离开。
5.3 Produce 请求进入 Broker 后校验什么
真实源码中校验和错误分支很多,本篇只保留与存储主线直接相关的类别:
Topic / Partition 是否存在当前 Broker 是否承载对应 Replica当前 Replica 是否为可写 LeaderRecordBatch 格式是否合法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 = 36891lastOffset = 36990newLogEndOffset = 36991简化伪代码:
AppendResult appendAsLeader(Records records) { long firstOffset = currentLogEndOffset(); assignOffsets(records, firstOffset); activeSegment().append(records); updateLogEndOffset(records); return result;}这段代码不是 Kafka 原始实现,只表达三个关键顺序:
- 根据 Leader Log 当前末尾确定新 Batch 的位置;
- 对 Batch 分配或校正最终 Offset;
- 追加后推进 Log End Offset。
为什么不能由每个 Producer 自己决定最终 Offset?
因为多个 Producer 会并发写同一 Partition。Offset 表达的是 Leader Log 中的全局追加顺序,只有掌握该 Log 写入权的一方才能确定最终顺序。
5.5 为什么只向 Active Segment 追加
Partition 中的写入目标不是随机 Segment:
Closed Segment 0Closed Segment 1Closed Segment 2Active 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 SizeSegment AgeIndex CapacityRelative Offset Range特定管理条件若需要 Roll:
Old Active Segment ↓ closeClosed Segment
Current Log End Offset ↓ becomes new base offsetNew Active SegmentRoll 的完整生命周期放在第八章。
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 ↓ writeOS Page Cache / Dirty Pages ↓ flush / writebackStorage Device需要区分三个时刻:
1. Java 写入调用返回2. 数据进入内核 Page Cache3. 脏页真正写入物理介质它们不一定同时发生。
所以,不能看到 Broker Append 成功就简单断言:
每条消息已经立即 fsync 到磁盘硬件但也不能反过来断言 Kafka 因而不可靠。Kafka 的可靠性模型组合了:
Leader AppendReplica FetchISRacksmin.insync.replicas故障选主这些属于第五篇。
本篇只需要记住:
Kafka 的高吞吐写路径通常依赖 OS 缓存和异步刷盘,而不是每条 Record 都执行一次同步物理写入。
5.9 为什么 Kafka 不在 JVM Heap 里缓存全部消息
假设 Broker 自己维护一个巨大 Java 消息缓存:
Disk File ↓ copyOS Page Cache ↓ copyJVM Heap Message Cache ↓ copySocket 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-3Fetch Offset = 36891Max Bytes = 1 MiBIsolation = read_uncommitted / read_committed这里的 Offset 表达:
从这个位置开始读取。
Broker 通常返回满足协议和字节限制的一批连续 RecordBatch:
368913689236893...不是像数据库主键查询一样,只返回 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
这一链路和写入链形成对称:
| 写入 | 读取 |
|---|---|
| ProduceRequest | FetchRequest |
| appendRecords | fetchMessages / readRecords |
| UnifiedLog.append | UnifiedLog.read |
| Active Segment Append | Segment Selection |
| FileRecords.append | FileRecords read / transfer |
6.3 读取前需要确定可见边界
并不是 Log 中所有已存在字节都一定能对当前 Consumer 可见。
Broker 还需要根据请求类型和隔离级别确定读取上界,例如后续文章会区分:
Log End OffsetHigh WatermarkLast Stable Offset本篇只保留抽象:
requestedOffset ↓validate range ↓determine max readable offset ↓read within visible boundary如果当前请求 Offset 不在 Broker 保留的范围内,可能产生 OFFSET_OUT_OF_RANGE。13
不要在本篇提前展开 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 ReadCold Historical ReadRemote Tier Read6.7 传统用户态数据复制路径
一个传统的读取再发送模型可能是:
Storage Device ↓Kernel Page Cache ↓ copyApplication User Buffer / JVM Buffer ↓ copyKernel 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 1P1 → Broker 2P2 → Broker 3P3 → 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 Parallelism8. Segment Roll、Retention 与 Compaction
8.1 为什么需要 Roll
Active Segment 不能无限增长。
达到一定条件后:
Active Segment ↓ RollClosed Segment
Current Log End Offset ↓New Segment Base Offset ↓New Active Segment
8.2 Roll 条件
Roll 判断可从以下维度理解:
Segment Size
segment.bytes当 Segment 文件达到配置大小附近时,创建新 Segment。
Segment Age
segment.mssegment.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 = 32000Current LEO = 64000
New files:00000000000000640000.log00000000000000640000.index00000000000000640000.timeindex00000000000000640000.txnindex8.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=1Offset 1: Key=B Value=3Offset 2: Key=A Value=2Offset 3: Key=A Value=5Compaction 目标是保留每个 Key 的最新值:
Key=B Value=3Key=A Value=5cleanup.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 第一层:消息格式
优先理解:
RecordsRecordBatchDefaultRecordBatchRecordDefaultRecordMemoryRecordsFileRecords重点问题:
- Batch Header 如何解析?
baseOffset和offsetDelta在哪里组合?- Compression Type 存在哪里?
- CRC 覆盖哪些区域?
FileRecords如何表达文件中的 Records 范围?- Batch 如何写入或传输到 NIO Channel?
Kafka 4.3 官方 Message Format 文档应当和源码一起阅读,而不是只看类名猜格式。1
9.2 第二层:索引
阅读:
OffsetIndexTimeIndexTransactionIndexAbstractIndex重点问题:
- 为什么 OffsetIndex 保存 Relative Offset?
lookup()返回的是最终 Record 还是邻近 Position?- 索引如何做二分查找?
- Index File 如何扩展、截断和重新映射?
- TimeIndex 为什么保存 Timestamp 到 Relative Offset?
- Index 损坏后如何依赖
.log重建?
9.3 第三层:Segment
阅读:
LogSegment重点问题:
- 一个 Segment 持有哪些 File 与 Index?
append如何同时推进数据和索引?read如何使用开始 Position?recover如何扫描 RecordBatch?shouldRoll的条件来自哪里?- Segment 关闭、截断和删除如何协调文件?
9.4 第四层:逻辑 Log
阅读:
UnifiedLogLocalLogLogSegmentsLogManager重点问题:
- Partition 的 Segment 集合如何管理?
floorSegment(offset)如何实现?- Active Segment 如何取得?
appendAsLeader的 Offset 分配和验证顺序是什么?read如何确定可见上界?- Roll 如何创建新 Segment?
- Log End Offset、恢复点和快照如何维护?
当前 4.x 的 UnifiedLog、LocalLog 和 LogSegment 已位于存储模块,阅读旧教程时要警惕类所在模块和实现语言已经发生变化。10
9.5 第五层:Partition 与 ReplicaManager
阅读:
PartitionReplicaManager重点问题:
- TopicPartition 如何映射到本地 Partition?
- 如何判断当前 Replica 是 Leader 还是 Follower?
- Produce 如何进入
appendRecordsToLeader? - Fetch 如何进入本地 Log?
- 哪些请求要等待副本条件满足?
- Delayed Produce / Delayed Fetch 如何与主链配合?
第五和第六个问题牵涉可靠性与 Purgatory,在第五篇展开。
9.6 第六层:Broker 请求入口
阅读:
SocketServerKafkaRequestHandlerKafkaApisRequestChannel重点问题:
- Socket 数据如何变成 Request?
- Processor 与 Handler Thread 如何分工?
- Produce、Fetch 在哪个 API Handler 分发?
- 错误如何转成协议 Response?
- 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启动:
docker compose up -d
docker logs -f kafka-ch02看到 Broker 启动完成后,再执行后续命令。
不同 Docker Desktop、文件权限和镜像小版本可能需要调整挂载目录权限。若本地挂载失败,可以先改用 Docker Named Volume,确认 Kafka 正常运行后再观察容器内部文件。
10.3 创建容易发生 Roll 的 Topic
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。
验证:
docker exec kafka-ch02 \ /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server localhost:9092 \ --describe \ --topic order-events查看 Topic 配置:
docker exec kafka-ch02 \ /opt/kafka/bin/kafka-configs.sh \ --bootstrap-server localhost:9092 \ --entity-type topics \ --entity-name order-events \ --describe10.4 批量发送测试数据
Linux / macOS / WSL:
python3 - <<'PY' | docker exec -i kafka-ch02 \ /opt/kafka/bin/kafka-console-producer.sh \ --bootstrap-server localhost:9092 \ --topic order-eventsimport jsonfor 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 目录
docker exec kafka-ch02 sh -lc \ 'ls -lh /var/lib/kafka/data/order-events-0'预期看到多组:
*.log*.index*.timeindex并观察:
- 文件名前缀是否递增;
.log是否接近配置的 1 MiB 后生成新文件;- 最后一个 Active Segment 是否可能明显小于前面 Segment;
.index是否远小于.log;- 多个文件是否共享相同 Base Offset 前缀。
如果目录名称或数据路径与示例不同,先执行:
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 文件:
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执行:
docker exec kafka-ch02 \ /opt/kafka/bin/kafka-dump-log.sh \ --files /var/lib/kafka/data/order-events-0/00000000000000000000.log \ --print-data-log观察:
baseOffsetlastOffsetpositionCreateTimeisvalidsizeproducerIdcompression codec不同 Kafka 小版本的输出字段和工具参数可能略有差异,以容器内:
/opt/kafka/bin/kafka-dump-log.sh --help为准。
实验目标不是记住工具输出,而是验证:
磁盘里保存的是 RecordBatch每个 Batch 有 Base Offset 和文件 PositionBatch 可以包含多条 Record10.7 查看索引
对 Offset Index 做 sanity check:
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执行步骤:
- 列出所有
.log文件前缀; - 选择最大且不超过 12000 的 Base Offset;
- 计算 Relative Offset;
- 使用 Dump Log 工具查看该 Segment 中 Batch 的 Base Offset 与 Position;
- 找到不大于目标 Offset 的邻近 Batch;
- 验证目标是否落在该 Batch Offset 范围内。
把结果记录成:
Target Offset:Selected Segment Base Offset:Relative Offset:Nearest Indexed Offset:Physical Position:Containing RecordBatch:这是本篇最重要的实践。
10.9 观察 Retention 与 Roll
可临时创建更短保留时间的测试 Topic:
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 清理实验
docker compose down同时删除本地测试数据:
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 → PositionTimeIndex 通常先得到近似 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. 第二篇自测与验收
完成本篇后,应当能够脱稿回答:
- Partition 为什么不能只使用一个无限增长文件?
- Partition、LogSegment、RecordBatch 和 Record 是什么关系?
- Active Segment 与 Closed Segment 分别承担什么职责?
- Segment 文件名为什么使用 Base Offset?
.log、.index、.timeindex和.txnindex分别保存什么?- Kafka Offset 与
.logPhysical Position 有什么区别? - 为什么 OffsetIndex 保存 Relative Offset?
- Kafka 为什么选择稀疏索引?
index.interval.bytes为什么不是消息条数?- Offset=36891 如何经过四个阶段定位到 RecordBatch?
- TimeIndex 为什么还需要 OffsetIndex?
- 为什么索引可以从
.log重建? - RecordBatch Header 中哪些字段与本篇最相关?
- 为什么 Compression 以 Batch 为边界?
- Broker 收到 ProduceRequest 后的主写入链是什么?
- 为什么网络线程不直接完成全部存储工作?
- 最终 Partition Offset 为什么由 Leader Log 决定?
- 为什么只写 Active Segment?
- 写入 Page Cache 与物理落盘有什么区别?
- Kafka 为什么不把全部消息维护成 JVM Heap Cache?
- Broker Fetch 的主要读取链是什么?
- Consumer 请求一个 Offset 时为什么会返回连续 Batch?
- Hot Tail Read 与 Cold Historical Read 有什么区别?
- Zero Copy 更准确的含义是什么?
- TLS、格式转换和 Remote Tier 为什么可能改变理想读取路径?
- Segment Roll、Retention Delete 与 Log Compaction 有什么区别?
- 为什么 Compaction 不会重新编号 Offset?
- 为什么“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 Device14.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 ↓ConsumerKafka 并不是简单地在一个普通文件上执行“顺序写”。
它将 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.size与linger.ms到底影响哪个执行阶段?buffer.memory耗尽后业务线程会发生什么?max.in.flight.requests.per.connection如何影响吞吐和顺序?- Retry、ACK 和 Idempotence 如何接入发送状态机?
下一篇:
《Kafka Producer:一条 Record 从
send()到 Partition Leader 的完整旅程》
参考资料
- Apache Kafka 4.3 Message Format
- Apache Kafka 4.3 Design
- Apache Kafka 4.3 Topic Configs
- Apache Kafka 4.3 Log Implementation
- Apache Kafka 4.3 Consumer Configs
- Apache Kafka 4.3 Broker Configs
- Apache Kafka 4.3 Tiered Storage
- Apache Kafka 4.3 Messages Implementation
- Apache Kafka KafkaApis Source
- Apache Kafka Storage Log Source Package
- Apache Kafka ReplicaManager Source
- Apache Kafka Developer Guide
- Apache Kafka Protocol Errors Source