前置回顾
- 第 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 反其道而行:所有消息一进来就落磁盘,吞吐反而甩开内存型中间件.原因在于磁盘的两种用法差距悬殊:
|
消息队列的负载特征天然适合顺序写:写是追加,读大多是顺序消费.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 下),目录里是若干组段文件:
|
为什么要分段
单文件无限追加有三个问题:
| 问题 | 单文件后果 | 分段后 |
|---|---|---|
| 清理困难 | 过期消息混在新消息中间,删不掉 | 整段删除,一次删一个文件 |
| 定位困难 | 在几 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.直接进去看:
|
输出类似:
|
只有一组段文件,因为数据量远没到 1GB,还在第一个段里.如果用过幂等 Producer,还会有 .snapshot 文件(生产者状态).
用 kafka-dump-log.sh 看段内容:
|
|
baseOffset 是逻辑序号,position 是消息在 .log 文件里的物理字节位置.索引做的事就是在这两者之间搭桥.
再看索引文件:
|
|
索引条目远少于消息条数(每隔若干 KB 才记一条),这就是稀疏索引的设计.
稀疏索引:从 Offset 到物理位置
Consumer 发来 FetchRequest(offset=5),Broker 怎么知道这条消息在文件哪个字节?
查找流程
|
索引文件是 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 后:
|
三个关键点:
- 写入只到 page cache 就返回.Kafka 默认不主动 fsync(
log.flush.interval.messages/log.flush.interval.ms默认不限制),刷盘时机交给操作系统.内存写入的延迟是微秒级. - 顺序追加没有寻道.写磁头几乎不用移动,HDD 也能跑满带宽.
- 单写者模型.每个 Partition 同一时刻只有 leader 一个写入者,没有并发写竞争,不需要加锁和随机定位.
Page Cache = 前台草稿本:记录员把流水先记在桌上的草稿本(内存),动作飞快.
打烊后由操作系统按水位统一誊抄到正式账本(磁盘).
代价:店面失火(断电)会丢草稿本上未誊抄的部分,所以可靠性靠多副本兜底,不靠 fsync.
那"写 page cache 就返回"是不是不可靠?单机看是的,但 Kafka 的可靠性模型建立在跨机器副本上(第 06 篇的 acks=all 链条):Leader 内存里的数据同时抄给了 Follower,一台机器断电,副本顶上.对每条消息都 fsync 会摧毁吞吐,只有少数强一致性场景才开 flush.messages=1.
读取:零拷贝
Consumer 拉取时,Broker 用上一节的索引定位后,大头是传输:一次 fetch 动辄几 MB.传统网络服务的读取路径:
|
Kafka Broker 只是"转发",不解包消息内容(压缩 batch 原样转发),数据根本没必要进用户态:
|
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:按时间/大小整段删(默认)
|
- 判断依据是段的最后一条消息时间戳(超过
log.retention.hours默认 7 天)或总大小超过log.retention.bytes - 删除粒度是整段,一个 unlink 系统调用,成本极低
- 正在写的 active segment 永远不会被删
- 后台线程每
log.retention.check.interval.ms(默认 5 分钟)检查一次
compact:按 Key 保留最新值
cleanup.policy=compact 时,Broker 定期重写旧段,每个 Key 只留最后一条消息:
|
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 重写旧段保留最新值
动手练习
- 观察段滚动:用
kafka-configs.sh --alter --entity-type topics --entity-name order-events --add-config segment.bytes=1048576把段上限改成 1MB,持续生产几分钟,观察目录里出现多个段文件. - 对照索引与消息:分别 dump
.log和.index,数一数消息条数和索引条数,验证稀疏比例(约每 4KB 一条). - 体验时间索引:按第 07 篇的步骤用
--to-datetime把 Group offset 重置到 5 分钟前,体会 timeindex 的按时间定位. - 观察 retention:建一个
retention.ms=60000的 Topic,生产一批消息后停止,等 1 分钟,看段文件被整段删除. - 压测对比:用
kafka-producer-perf-test.sh压 10 万条,对比acks=1和acks=all的吞吐差,理解"等副本"的代价. - 估算文件句柄:
ls /proc/<broker_pid>/fd | wc -l看 Broker 打开的段文件数量,估算集群规模对应的 ulimit 需求.
下一篇进入副本机制与 ISR:leader/follower 怎么同步,什么时候消息才算"已提交",ISR 为什么缩容.