KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA
01 · 存储模型 — keel 龙骨
第 00 章说「消息不会被消费掉」,但它到底躺在哪、长什么样,决定了 Kafka 的性能边界:为什么顺序写能扛住百万级吞吐,为什么按 timestamp 查历史消息很贵,为什么一个分区就是一个 CPU 核。这一章把分区目录拆开看。
第 00 章说「消息不会被消费掉」,但它到底躺在哪、长什么样,决定了 Kafka 的性能边界:为什么顺序写能扛住百万级吞吐,为什么按 timestamp 查历史消息很贵,为什么一个分区就是一个 CPU 核。这一章把分区目录拆开看。
现场
本机 topic orders-01(3 分区 3 副本)的分区 0 上,我先用固定 key 打了 8000 条、每条约 223 字节的消息,然后直接进数据目录看文件:
total 22346
-rw-r--r-- 1 ... 10485760 00000000000000000000.index
-rw-r--r-- 1 ... 1783691 00000000000000000000.log
-rw-r--r-- 1 ... 10485756 00000000000000000000.timeindex
-rw-r--r-- 1 ... 11 leader-epoch-checkpoint
-rw-r--r-- 1 ... 43 partition.metadata
第一眼会愣住:.log 是 1.78 MB(数据本体),.index 和 .timeindex 却各是 10 MB——索引比数据大。这不可能,除非它们不是「实打实写满」的。
一条消息从 send 到被按 offset 读到
flowchart TD
A["① 生产者 send(key=bulk)"] --> B["② 分区器 hash(key) % 分区数<br/>选中 orders-01-0"]
B --> C["③ 追加到该分区活动段的 .log<br/>顺序写,只 append"]
C --> D{"④ 段大小是否超过<br/>segment.bytes?"}
D -- "否" --> C
D -- "是" --> E["⑤ 滚动:关闭当前段<br/>开一个新的 baseOffset 段"]
C -. "每累计写入超过 log.index.interval.bytes<br/>就往 .index 落一条 offset->position" .-> F[".index 稀疏索引"]
C -. "同时往 .timeindex 落 timestamp->offset" .-> G[".timeindex 稀疏索引"]
H["⑥ 消费者 fetch(partition=0, offset=2900)"] --> F
F --> I["⑦ 先查索引拿到最近的一条<br/>offset 291 position 48828"]
I --> J["⑧ 从该 position 起顺序扫描 .log<br/>找到 offset 2900 返回"]
style A fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style B fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style C fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
style D fill:#fff4e6,stroke:#e8590c,color:#1a1a2e
style E fill:#fff4e6,stroke:#e8590c,color:#1a1a2e
style F fill:#fff9db,stroke:#f08c00,color:#1a1a2e
style G fill:#fff9db,stroke:#f08c00,color:#1a1a2e
style H fill:#f3f0ff,stroke:#7048e8,color:#1a1a2e
style I fill:#fff9db,stroke:#f08c00,color:#1a1a2e
style J fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
这里的 ⑦ 是理解存储模型的关键:索引是稀疏的,不是每条消息一行。所以读一个 offset 是「先近似定位、再小范围顺序扫描」,两次 IO 而不是一次。
精确定义
分区(partition) 是 Kafka 的并行与顺序单元。一个分区就是一串追加记录,天然有序;不同分区之间没有顺序保证。分区数决定一个消费组的并行上限(第 04 章)。
段(segment) 是分区的物理分片。文件名就是这段的 base offset,比如 00000000000000002165.log 表示这段从 offset 2165 开始。段存在的作用是让「删除旧数据」变成一个文件级操作(删整个段),而不是在文件中间挖洞。
稀疏索引 只在「自上次写索引以来累计写入超过 log.index.interval.bytes」时落一条,默认 4096 字节。索引条目是 8 字节:4 字节相对 offset + 4 字节物理 position。
producerId / baseSequence 是记录批头里的字段,幂等生产者靠它们去重(第 03 章)。magic: 2 是记录批格式版本,crc 是校验。
证据:稀疏索引到底有多稀疏
直接把索引 dump 出来看(dumplog 一个 .index 文件):
offset: 145 position: 16276
offset: 218 position: 32552
offset: 291 position: 48828
offset: 364 position: 65104
...
相邻两条索引的 offset 差 73,position 差 16276。换算一下:日志总长 1783691 字节,索引条目 109 条,平均每条索引覆盖 8000/109 ≈ 73.4 条记录、1783691/109 ≈ 16364 字节。.log 里每条记录约 223 字节,73 × 223 ≈ 16279,对得上。
再往上追一层:生产者每批最多装 batch.size(默认 16384 字节)——dump 出第一批的批头是 baseOffset: 0 lastOffset: 72 count: 73 size: 16276,正好 73 条、16276 字节。索引是按批边界落的:写入量跨过 4096 字节的下限后,要等到这一批写完才落一条索引,所以实测间隔(约 16 KB)远大于配置的 4096。
这就是为什么 .index 显示 10 MB:log.index.size.max.bytes 默认 10485760,索引文件是 mmap 预分配到这个大小的,实际有效内容只有前 109×8 = 872 字节。滚动后关闭的段会看到真实大小——第 08 章滚动出的关闭段,索引文件只有 520 字节。
证据:记录长什么样
dumplog --print-data-log 把 .log 里的真实字节翻译成字段:
baseOffset: 0 lastOffset: 5 count: 6 baseSequence: 0 lastSequence: 5
producerId: 1000 producerEpoch: 0 partitionLeaderEpoch: 0
isTransactional: false isControl: false ... size: 181 magic: 2 compresscodec: none crc: 2985226122 isvalid: true
| offset: 0 CreateTime: ... keySize: 6 valueSize: 7 sequence: 0 headerKeys: [] key: user-1 payload: order-A
| offset: 1 ... key: user-2 payload: order-B
逐字段读:
baseOffset: 0 lastOffset: 5 count: 6——这一批含 6 条记录,覆盖 offset 0..5。批是 Kafka 的最小存取单位,一次 IO 读一整批。baseSequence: 0 lastSequence: 5——批内每条记录的序号 0..5,幂等去重用的就是它。producerId: 1000——这条记录来自哪个生产者实例。ConsoleProducer在 4.x 默认开幂等,所以每次运行会分配一个新 PID(这里 1000,后来换了个生产者又变成 1001)。isTransactional: false isControl: false——非事务、非控制批。第 05 章的事务实验里,这里会变成true。magic: 2——记录批格式 v2(0.11 之后),支持 headers、幂等、事务。crc: ... isvalid: true——整批的 CRC 校验,读的时候会验;校验不过会走日志恢复流程。- 每条记录的
key: user-1 payload: order-A——key决定分区(hash(key) % 分区数),headerKeys: []是消息头,本机这次没设。
顺序写为什么快、随机读为什么贵
Kafka 的写路径只有一条:追加到活动段末尾。磁盘顺序写的吞吐接近内存带宽(本机实测约 22 MB/s,受限于这台机器的磁盘与三副本同步,第 09 章有数字),因为它不需要 seek,所有写都落在文件末尾,page cache 友好、预读有效。
读路径则分两种:
- 消费(从某个 offset 顺序往后读)是顺序读,同样便宜——消费者位点只会前进,配合操作系统的预读,实际吞吐很高。
- 按时间或按 offset 随机跳着读(比如「查 3 小时前那条消息」)要先在
.timeindex/.index上做二分,再进.log顺序扫描一小段。单次还能接受,但如果你把它当成数据库的二级索引来用(大量随机点查),就会撞上稀疏索引的固有限制:它只能帮你把范围缩小到一个段内的小窗,剩下的还是要顺序扫。
这就是「Kafka 能扛顺序写但扛不住随机读」的准确含义:不是随机读慢得夸张,而是它没有任何为随机点查优化的结构——.index 是稀疏的、只有 8 字节一条,.log 里没有按 key 的索引。要按 key 查历史值,那是 compact 日志或者一张真数据库的活。
一次完整运行:段滚动
第 08 章的实验把 topic 的 segment.bytes 调到最小允许值 1 MB(1048576),生产约 6 MB 数据,得到:
00000000000000000000.log size=1047556
00000000000000002165.log size=1047556
00000000000000004330.log size=1047556
00000000000000006495.log size=1047556
00000000000000008660.log size=1047556
00000000000000010825.log size=568546
六个段,前五个都是 1047556 字节(略小于 1 MB,因为攒满一批才滚),最后一个是没写满的活动段。dumplog 最后一段得到 Log starting offset: 10825——文件名就是它的 base offset。
同时目录里多出了 .snapshot 文件(00000000000000002165.snapshot 等)。它是生产者状态快照:记录每个 producerId 在这个段结束时的最后序列号,broker 重启后靠它恢复幂等去重状态。因为本机用的 kafka-python 默认开了幂等,所以每个滚动点都有快照。
注意一个对照:滚动完的段,索引文件是 520 字节(30 条真实索引);只有活动段的索引才是 10 MB 的预分配。看到 10 MB 不用慌,那是 mmap 的地址空间,不是磁盘占用。
主动破坏:把段调小会怎样
segment.bytes 想设成 10240(10 KB)会被直接拒绝:
Invalid value 10240 for configuration segment.bytes: Value must be at least 1048576
下限是 1 MB。把段调小的代价是元数据变多:段越多,.index、.timeindex、.snapshot 和 leader-epoch-checkpoint 就越多,broker 打开的文件句柄也越多。本机实验里 6 MB 数据就产生了 6 个段、18 个索引/快照文件。生产上段大小通常留在几百 MB 到 1 GB,别为了「删得快」把段调得过小。
生产边界
- 教学替身 vs 真实依赖:本机的日志目录直接放在 Windows 的普通磁盘上,索引用 mmap。真实集群会把
log.dirs配到多块盘上,broker 会并行使用多个目录。另外本机是 Windows,日志目录和索引句柄有个坑——任何让 broker「重命名/删除含打开索引的目录」的操作(删 topic、缩副本 reassign)都会失败并把 log dir 标记 offline,进而关掉 broker,细节见第 06 章。 - 要盯的指标:每块盘的可用空间、
log.segment.bytes与段数量、LogFlush相关延迟、单分区的大小。单个分区没写满的段不会被 retention 删除,长事务会拖住一个分区末尾的段(deleteHorizonMs非空时该段不删)。 - 失败策略:磁盘写满时 broker 会把 log dir 标记 offline,落在这个目录上的分区整体不可用。上线前给
log.retention.bytes或磁盘告警留冗余,别等写满。
动手
- 给一个主题打不同的 key,用
dumplog印证「同一个 key 落同一个分区」(key 哈希)。 - 在
orders-01-0上再 dump 一次.index,找出 offset 与 position 的比值,验证它约等于「每批大小」。 - 断言题:活动段的
.index是 10485760 字节,而滚动完的段是几百字节。解释这两个数字分别从哪来。
自测
- 稀疏索引间隔的配置是 4096 字节,为什么实测是约 16364 字节?批次在哪一步介入了?
- 为什么
.index文件可以是 10 MB 而磁盘占用几乎为 0?log.index.size.max.bytes在这里扮演什么角色? - 「Kafka 扛顺序写不扛随机读」——随机读具体慢在哪一步?稀疏索引为什么不解决这个问题?
- 段文件名(如
00000000000000002165.log)里的数字是什么?为什么删旧数据是按段删而不是按条删? magic: 2和baseSequence在存储层各解决什么问题?如果magic是 0(旧格式)会少掉哪些能力?