阶段二 · 架构原理与调优

存储引擎:日志分段,索引,零拷贝

前置回顾

  • 第 02 篇定义了数据模型:Partition 是 append-only log,Offset 是消息地址
  • 第 06 篇讲过消息落到 Broker 磁盘才算"不丢",第 07 篇搭好了 3-Broker 本地集群
  • 本篇拆开磁盘看里面存了什么,回答阶段二的核心问题:Kafka 的吞吐为什么能到百万级

一句话总结

Kafka 的存储引擎 = 日志分段 + 稀疏索引 + 零拷贝.
消息顺序追加到 segment 文件,page cache 吸收读写,sendfile 把数据从磁盘直达网络.理解"顺序写 + 内存缓存 + 绕过用户态",就理解了 Kafka 快在哪里.

为什么 Kafka 把消息放磁盘

直觉上磁盘慢.一块机械盘随机写 4KB 只有 100 次/s 左右,算下来约 0.4MB/s,所以大多数消息中间件把队列放内存,磁盘只做持久化兜底.

Kafka 反其道而行:所有消息一进来就落磁盘,吞吐反而甩开内存型中间件.原因在于磁盘的两种用法差距悬殊:

4KB 随机写:100 IOPS × 4KB ≈ 0.4 MB/s
顺序写:一次寻道后持续写,HDD 也能到 100+ MB/s

差距:两个数量级

消息队列的负载特征天然适合顺序写:写是追加,读大多是顺序消费.Kafka 把这一特征用到了极致.

顺序写 = 流水线装货:传送带不停,货箱一个接一个放上去.
随机写 = 找停车位:每次都先绕场一圈找到空位才能停,时间全花在"找"上.

传统数据库为什么不行?MySQL 用 B+ 树组织数据,写一条消息 = 定位叶子页(随机 IO)+ 更新索引(更多随机 IO).在"追加日志"这个场景里,B+ 树的结构维护全是浪费.

Kafka 官方早期的基准测试:3 台廉价服务器,200 万次写入/s,300 万次读取/s.支撑这个数字的是本篇要拆的三个机制:日志分段、稀疏索引、零拷贝.


Partition 目录:日志分段

日志分段(Log Segment):Partition 在磁盘上不是一个无限增长的文件,而是按固定大小切成的一段段文件,每段有一个 base offset.

第 02 篇说过 Partition 是 append-only log.落在磁盘上,一个 Partition 是一个目录(log.dirs 下),目录里是若干组段文件:

/bitnami/kafka/data/order-events-0/            ← 一个 Partition 一个目录
├── 00000000000000000000.log ← 段 1,base offset 0
├── 00000000000000000000.index ← 段 1 的偏移量索引
├── 00000000000000000000.timeindex ← 段 1 的时间索引
├── 00000000000000942650.log ← 段 2,base offset 942650
├── 00000000000000942650.index
├── 00000000000000942650.timeindex
├── partition.metadata ← 分区元数据(KRaft)
└── leader-epoch-checkpoint ← leader 任期检查点

为什么要分段

单文件无限追加有三个问题:

问题 单文件后果 分段后
清理困难 过期消息混在新消息中间,删不掉 整段删除,一次删一个文件
定位困难 在几 TB 文件里找 offset 只能线性扫 先定位段,再在段内找
管理成本 巨型文件的句柄、mmap、校验都重 段文件 1GB 左右,可控

分段后,正在写入的只有最后一个段(active segment),其余段全部只读.删除、归档、索引加载都只发生在旧段上,不碰写入路径.

段的命名与滚动

  • 文件名是段的 base offset(该段第一条消息的 offset),20 位补零
  • 写满 log.segment.bytes(默认 1GB)或到达 log.segment.ms(默认 7 天)就滚动,开新段
  • 段内 offset 用相对值编号:段的 base offset 记作 0

日志分段 = 账本分册:一本账本记满换下一本,每本从第 1 页重新编号.
查账先找第几册(段),再翻册内的页(offset).旧册整本归档或销毁,不碰正在写的那册.


动手验证:拆开一个 Partition 目录

第 07 篇的 docker-compose 集群数据在容器的 /bitnami/kafka/data.直接进去看:

docker compose exec kafka-0 bash

# 如果 order-events 还没建,按第 07 篇的步骤先创建
# 造点数据(交互式,回车即发一条,Ctrl+C 退出)
kafka-console-producer.sh --topic order-events \
--property "parse.key=true" --property "key.separator=:" \
--bootstrap-server localhost:19092
# 输入几十条: user-1:{"action":"purchase"}

# 退出 producer 后看分区目录
ls -lh /bitnami/kafka/data/order-events-0/

输出类似:

-rw-r--r-- 1 1001 root 7.6K ... 00000000000000000000.index
-rw-r--r-- 1 1001 root 12K ... 00000000000000000000.log
-rw-r--r-- 1 1001 root 12K ... 00000000000000000000.timeindex
-rw-r--r-- 1 1001 root 74B ... leader-epoch-checkpoint
-rw-r--r-- 1 1001 root 78B ... partition.metadata

只有一组段文件,因为数据量远没到 1GB,还在第一个段里.如果用过幂等 Producer,还会有 .snapshot 文件(生产者状态).

用 kafka-dump-log.sh 看段内容:

kafka-dump-log.sh --files /bitnami/kafka/data/order-events-0/00000000000000000000.log \
--print-data-log | head -5
baseOffset: 0 lastOffset: 0 ... position: 0   CreateTime: 1725468000000 ... isvalid: true
baseOffset: 1 lastOffset: 1 ... position: 356 CreateTime: 1725468012000 ... isvalid: true
baseOffset: 2 lastOffset: 2 ... position: 712 CreateTime: 1725468024000 ... isvalid: true

baseOffset 是逻辑序号,position 是消息在 .log 文件里的物理字节位置.索引做的事就是在这两者之间搭桥.

再看索引文件:

kafka-dump-log.sh --files /bitnami/kafka/data/order-events-0/00000000000000000000.index
offset: 0 position: 0
offset: 3 position: 1068
offset: 6 position: 2136

索引条目远少于消息条数(每隔若干 KB 才记一条),这就是稀疏索引的设计.


稀疏索引:从 Offset 到物理位置

Consumer 发来 FetchRequest(offset=5),Broker 怎么知道这条消息在文件哪个字节?

查找流程

FetchRequest offset=5
│
▼
1. 定位段:内存中的段列表按 base offset 有序,二分找到包含 offset=5 的段
│ (段 0 的区间是 [0, 942650) → 命中段 0)
▼
2. 查偏移索引:二分查找 ≤ 5 的最大条目 → offset=3, position=1068
│
▼
3. 顺序扫:从 position=1068 往后读,跳过 offset 3, 4,停在 5

索引文件是 mmap 进内存的,二分查找只读内存,不发起磁盘 IO.找到 position 后最多顺序扫 log.index.interval.bytes(默认 4KB)的数据量.

为什么索引是稀疏的

方案 索引条目 文件大小 查找成本
密集索引 每条消息一条 与消息量成正比,巨型 一步定位
稀疏索引(实际) 每 4KB 一条 段的 1/512 二分 + 扫 4KB

条目格式:relativeOffset(4 字节) + position(4 字节),8 字节一条.relativeOffset 是相对段 base offset 的差值,绝对值不重复占位.1GB 的段,索引文件才 2MB 左右,可以整个 mmap 到内存.

稀疏索引 = 教科书目录:目录只列章,不列每一行.
找"第 5 页的内容",先翻到目录指出的章节起点,再往后翻几行.省下给每行做目录的成本.

索引指向的是 batch 不是单条消息

第 03 篇讲过 Producer 按 batch 发送.落盘的最小单位也是 batch(RecordBatch),索引条目指向 batch 的 base offset.
查到 batch 后,在 batch 内部顺序扫到目标 offset,代价可以忽略.

时间索引:按时间找 Offset

.timeindex 存的是 timestamp(8 字节)+ relativeOffset(4 字节),同样稀疏.

它服务的是"按时间定位"的场景.第 07 篇用过的 --to-datetime 重置 offset,原理就是二分时间索引找到该时间点最近的 offset,再重置.


写入路径与读取路径

存储引擎的两条核心路径,决定了吞吐的上限.

写入:顺序追加 + Page Cache

Broker 收到 ProduceRequest 后:

Producer batch ──► Broker 网络线程
│
▼
active segment 末尾追加(顺序写)
│
▼
Linux Page Cache(内存) ← 到这里就返回 ack 给 Producer
│
▼
后台异步:OS 按水位刷盘;副本同步(下一篇讲)

三个关键点:

  1. 写入只到 page cache 就返回.Kafka 默认不主动 fsync(log.flush.interval.messages / log.flush.interval.ms 默认不限制),刷盘时机交给操作系统.内存写入的延迟是微秒级.
  2. 顺序追加没有寻道.写磁头几乎不用移动,HDD 也能跑满带宽.
  3. 单写者模型.每个 Partition 同一时刻只有 leader 一个写入者,没有并发写竞争,不需要加锁和随机定位.

Page Cache = 前台草稿本:记录员把流水先记在桌上的草稿本(内存),动作飞快.
打烊后由操作系统按水位统一誊抄到正式账本(磁盘).
代价:店面失火(断电)会丢草稿本上未誊抄的部分,所以可靠性靠多副本兜底,不靠 fsync.

那"写 page cache 就返回"是不是不可靠?单机看是的,但 Kafka 的可靠性模型建立在跨机器副本上(第 06 篇的 acks=all 链条):Leader 内存里的数据同时抄给了 Follower,一台机器断电,副本顶上.对每条消息都 fsync 会摧毁吞吐,只有少数强一致性场景才开 flush.messages=1.

读取:零拷贝

Consumer 拉取时,Broker 用上一节的索引定位后,大头是传输:一次 fetch 动辄几 MB.传统网络服务的读取路径:

传统路径:read() + write(),4 次拷贝,4 次上下文切换

磁盘 ──DMA──► 内核 page cache ──CPU──► 用户态 buffer ──CPU──► socket buffer ──DMA──► 网卡

数据必须进用户态,应用才能操作;再拷回内核,才能发出去

Kafka Broker 只是"转发",不解包消息内容(压缩 batch 原样转发),数据根本没必要进用户态:

sendfile 路径:2 次拷贝,2 次上下文切换

磁盘 ──DMA──► 内核 page cache ──DMA(scatter-gather)──► 网卡

数据全程在内核态流动,应用只调一次 sendfile(),不参与搬运

Java 里对应 FileChannel.transferTo,底层就是 Linux 的 sendfile.效果对比:

指标 传统路径 sendfile
CPU 拷贝次数 2 次 0 次(新内核 scatter-gather)
上下文切换 4 次 2 次
用户态内存占用 每次读取都要 buffer 0

消息不经过 JVM 堆,不产生 GC 压力,CPU 只负责调度和网络栈.这就是"3 台机器 300 万次读取/s"的来源.

传统路径 = 仓库 → 办公室 → 货车:快递先从仓库搬进办公室,再搬上车,经手两次.
sendfile = 仓库直通货车:传送带从仓库货架直接伸进货箱,货不进办公室.

零拷贝对压缩消息同样生效

Producer 压缩的 batch,Broker 不解压,只校验 CRC 后原样转发,所以 sendfile 照常工作.解压发生在 Consumer 端.
两个例外:1. 老客户端需要降级消息格式(downconversion),Broker 必须解压重编码,退回传统路径;2. 启用了 SSL 端点,加密必须经过用户态.


日志清理:Retention 与 Compaction

消息不会永久保留.清理有两种策略,由 Topic 的 cleanup.policy 决定.

delete:按时间/大小整段删(默认)

segment 0                  segment 1                  segment 2(active)
base=0,最后写入 8 天前 base=942650,3 天前 base=1885300,正在写
│ │ │
▼ ▼ ▼
整段删除,文件 unlink 保留 保留
  • 判断依据是段的最后一条消息时间戳(超过 log.retention.hours 默认 7 天)或总大小超过 log.retention.bytes
  • 删除粒度是整段,一个 unlink 系统调用,成本极低
  • 正在写的 active segment 永远不会被删
  • 后台线程每 log.retention.check.interval.ms(默认 5 分钟)检查一次

compact:按 Key 保留最新值

cleanup.policy=compact 时,Broker 定期重写旧段,每个 Key 只留最后一条消息:

重写前:offset 0  key=user-1  value=v1
offset 1 key=user-2 value=v1
offset 2 key=user-1 value=v2 ← 旧值
offset 3 key=user-1 value=v3 ← 最新,保留

重写后:offset 1 key=user-2 value=v1
offset 3 key=user-1 value=v3

compact 后消息仍然有序,但 offset 出现空洞(第 02 篇提过"offset 不连续也合法").__consumer_offsets 用的就是 compact 策略(第 04 篇).

选型:事件流、审计日志用 delete;状态表、offset 存储用 compact.


生产注意事项

文件句柄数

每个段最多 3 个文件(log/index/timeindex),大分区数集群的句柄数 = 分区数 × 段数 × 3 的量级.
默认 ulimit 通常不够,按 10 万起步调:ulimit -n 100000.segment.bytes 调小会让段更多、句柄更多,调之前先算账.

内存大部分要留给 Page Cache

Kafka 读写都走 page cache,命中率直接决定吞吐.
JVM 堆给 5-6GB 足够(消息不进堆),其余内存留给操作系统.堆配得过大反而挤压缓存.

磁盘:JBOD 而非 RAID,本地盘而非 NAS

副本机制已经提供冗余,RAID 的额外冗余是浪费,还拖慢重建.多块本地盘组 JBOD,分区目录会分散到各盘,并行吞吐更高.
网络存储(NAS/SAN)的延迟抖动对顺序写没有收益,官方明确不推荐.

lag 过大的冷读会把压力传导给磁盘

消息在 page cache 里是"热"的,读它不碰磁盘.
Consumer 落后太多时,要读的数据早已被挤出缓存,每次 fetch 都穿透到磁盘随机读,整个 Broker 的吞吐都会受影响.这也是第 07 篇强调 lag 告警的底层原因.

磁盘写满 = 分区下线

磁盘 100% 后,Broker 会拒绝该盘上分区的写入,甚至整机下线,比"慢"严重得多.
必须:监控磁盘使用率,设 log.retention.bytes 兜底,预留告警余量.


为什么不用内存存消息,像 Redis 那样?

内存成本是磁盘的几十倍,消息量级下不可持续.
而且 Kafka 靠 page cache 已经拿到了"内存读写"的大部分收益:写是写内存,读只要数据还在缓存里也是读内存.冷数据落到磁盘,总成本可控.

索引里存的是相对 offset,怎么换算绝对 offset?

绝对 offset = 段 base offset + relativeOffset.段文件名的 20 位数字就是 base offset,换算在加载索引时一次完成,查找时直接比较绝对 offset.

page cache 里的数据断电就丢,Kafka 凭什么说可靠?

单机断电确实会丢 page cache 中未落盘的部分.Kafka 的答案不是更勤快地 fsync,而是跨机器副本:acks=all + min.insync.replicas≥2(第 06 篇),消息在多个 Broker 的内存里同时存在,单机断电由副本兜底.
只有单副本集群才需要靠 fsync 保护数据.

为什么删除消息不能像数据库一样删单条?

append-only 文件中间删一条,后面的数据要整体前移,等于随机写,吞吐模型就崩了.整段删除只需要 unlink.
要"删旧值"的语义,用 compact 策略,它本质是后台重写旧段,同样不动 append 路径.

零拷贝这么好,什么场景会失效?

只要 Broker 需要"看懂"消息就得进用户态:老协议降级(downconversion)、SSL 加密.
现代客户端 + PLAINTEXT 端点下不会触发,保持客户端和 Broker 版本不落后太多即可.


快速回顾

  • 日志分段:Partition 目录里按 base offset 命名的段文件,只写 active segment,清理按段进行
  • 稀疏索引:每 4KB 一条 offset→position 条目,mmap 进内存,二分定位后小范围顺序扫
  • 写入路径:顺序追加 + page cache,默认不主动 fsync,可靠性交给跨机器副本
  • 读取路径:sendfile 零拷贝,数据不经过用户态,压缩消息原样转发
  • 日志清理:delete 按时间/大小整段删,compact 按 Key 重写旧段保留最新值

动手练习

  1. 观察段滚动:用 kafka-configs.sh --alter --entity-type topics --entity-name order-events --add-config segment.bytes=1048576 把段上限改成 1MB,持续生产几分钟,观察目录里出现多个段文件.
  2. 对照索引与消息:分别 dump .log 和 .index,数一数消息条数和索引条数,验证稀疏比例(约每 4KB 一条).
  3. 体验时间索引:按第 07 篇的步骤用 --to-datetime 把 Group offset 重置到 5 分钟前,体会 timeindex 的按时间定位.
  4. 观察 retention:建一个 retention.ms=60000 的 Topic,生产一批消息后停止,等 1 分钟,看段文件被整段删除.
  5. 压测对比:用 kafka-producer-perf-test.sh 压 10 万条,对比 acks=1 和 acks=all 的吞吐差,理解"等副本"的代价.
  6. 估算文件句柄:ls /proc/<broker_pid>/fd | wc -l 看 Broker 打开的段文件数量,估算集群规模对应的 ulimit 需求.

下一篇进入副本机制与 ISR:leader/follower 怎么同步,什么时候消息才算"已提交",ISR 为什么缩容.