主题
第 7 章 存储与日志格式:Kafka 在磁盘上长什么样?
目标读者:知道 Kafka 是「写到磁盘」的,但说不清楚一条消息从 Producer
send()到磁盘上的字节序列经历了什么;面试问到.log/.index/.timeindex关系时只能干笑的同学。学完你会:能拿一个 Segment 文件用
kafka-dump-log.sh解码出来逐字节读,能解释「按 Offset 找消息」走了几次 I/O,能区分 V0 / V1 / V2 三种消息格式的字节布局,知道.snapshot和leader-epoch-checkpoint的存在意义。
0. 导读:Kafka 的本质是「分布式提交日志」
很多人第一次学 Kafka 都被「消息队列」这个词带偏了。Kafka 官方对自己的定义是:
"a distributed, partitioned, replicated commit log service."
Commit Log(提交日志) 才是 Kafka 的灵魂,「消息队列」只是这层日志暴露给业务的接口。理解了存储层,你才能理解为什么:
- Kafka 可以回放(rewind)—— 因为消息一直在磁盘上躺着。
- Kafka 可以多消费者互不干扰 —— 因为每个消费者只是维护自己的 Offset 指针。
- Kafka 可以抗住超大吞吐 —— 因为顺序写日志比随机写数据库快两个数量级(详见第 8 章)。
- Kafka 可以实现 Log Compaction —— 因为日志格式自带 Key,可以定期压缩。
本章会带你亲手扒开 Kafka 的磁盘目录,看清楚每一个文件、每一个字段是干什么的。
1. 一个生活类比:图书馆的卡片目录
把一个 Kafka Topic 想象成一座图书馆:
- Topic = 整个图书馆(按主题分类)
- Partition = 一个书架(按子分类切分)
- Segment = 书架上的一层(按入库时间分段,写满了就开新一层)
.log文件 = 这一层上真正的书(消息内容).index文件 = 这一层的编号索引卡(书号 → 位置).timeindex文件 = 这一层的时间索引卡(购入日期 → 书号)leader-epoch-checkpoint= 这层书架的主管交接登记本(谁什么时候管的).snapshot= 盘点快照(用于事务/幂等的版本号)
要找一本书,先看「编号索引卡」翻到大致位置,再走到书架前顺序往下找几本就到了。这就是 Kafka 的「稀疏索引 + 顺序扫」的全部哲学。
2. 日志目录布局 log.dirs
2.1 入门视角:从一台 Broker 看起
Broker 启动时会读 server.properties 里的 log.dirs:
properties
# 单盘
log.dirs=/var/lib/kafka/data
# 多盘(提升 IO 并行度)
log.dirs=/data1/kafka,/data2/kafka,/data3/kafka每个目录被 Kafka 当作一块独立 LogDir。每个 Topic-Partition 会被分配到其中一块 LogDir 上(靠 Kafka 内部的负载均衡)。同一个 Partition 的所有 Segment 一定在同一个 LogDir。
2.2 完整目录树
假设有 learn.07.storage Topic,4 个分区,2 个内部 Topic,看一下磁盘上长什么样:
/var/lib/kafka/data/ ← log.dirs 之一
├── meta.properties ← Broker 元数据(cluster.id, broker.id)
├── recovery-point-offset-checkpoint ← 已经 fsync 到的 offset
├── replication-offset-checkpoint ← 已复制给所有 ISR 的 offset
├── log-start-offset-checkpoint ← Topic 当前的 start offset
├── cleaner-offset-checkpoint ← Log Cleaner 的进度(compaction)
│
├── learn.07.storage-0/ ← Topic-Partition 0 目录
│ ├── 00000000000000000000.log ← 第 1 个 Segment 的数据
│ ├── 00000000000000000000.index ← Offset 稀疏索引
│ ├── 00000000000000000000.timeindex ← 时间稀疏索引
│ ├── 00000000000000005678.log ← 第 2 个 Segment(base offset = 5678)
│ ├── 00000000000000005678.index
│ ├── 00000000000000005678.timeindex
│ ├── 00000000000000005678.snapshot ← 幂等/事务快照
│ ├── leader-epoch-checkpoint ← Leader Epoch 历史
│ └── partition.metadata ← (KRaft 后) Topic ID + 版本号
│
├── learn.07.storage-1/
│ └── ...
├── learn.07.storage-2/
│ └── ...
├── learn.07.storage-3/
│ └── ...
│
├── __consumer_offsets-0/ ← 内部 Topic:消费组 Offset
│ └── ...
└── __transaction_state-0/ ← 内部 Topic:事务状态
└── ...2.3 文件命名的玄学
文件名 00000000000000005678.log 中的 5678 是什么?
答:这个 Segment 包含的第一条消息的 Offset,也叫 base offset。
- 文件名是 20 位 0 填充的整数 —— 字典序就是数值序,方便
ls排序。 - 同一个 Segment 的
.log/.index/.timeindex三个文件 base offset 完全一致。 - 给一个 Offset = 7000,要找它在哪个文件,就用「所有 base offset ≤ 7000 的最大那个」作为目标 Segment(这里就是 5678)。然后再用
.index在.log里精确定位。
3. Topic-Partition-Segment 三层结构
关键不变量:
- Topic 是逻辑概念,磁盘上不存在「topic 文件夹」,只有「topic-partition 文件夹」。
- Partition 内的 Segment 严格按 base offset 排序,不可重叠。
- 任意时刻每个 Partition 只有 1 个 active Segment(最新那个),其它都是「only read」状态。
- Segment 「封存」(roll)后才能被删除 / 压缩,active Segment 永远不会被清理。
4. 五大文件的作用详解
4.1 .log:真正的消息数据
.log 是一个追加写的二进制文件,里面是一串 RecordBatch(V2 格式)拼接而成。
+-------- RecordBatch 1 --------+-------- RecordBatch 2 --------+--- ... ---+
| header (61B) | records (变长) | header (61B) | records (变长) | |
+-------------------------------+-------------------------------+-----------+
↑
file offset = 0每个 RecordBatch 内部又是若干 Record(一条消息)。Producer 一次 send() 凑成一个 batch 发给 Broker,Broker 直接把整个 batch 字节追加到 .log 文件,不解开。这是 Kafka 高吞吐的核心之一(详见第 8 章「批量」)。
4.2 .index:Offset → Position 稀疏索引
结构:固定长度的 (relative_offset, position) 对,每对 8 字节(每个 4 字节,小端)。
+----------+----------+
| relOff=0 | pos=0 | ← Offset = base + 0 的消息在 .log 文件偏移 0
+----------+----------+
| relOff=23| pos=4096 | ← Offset = base + 23 在 .log 文件偏移 4096
+----------+----------+
| relOff=51| pos=8192 |
+----------+----------+
| relOff=82| pos=12288|
+----------+----------+
...两个关键概念:
- 稀疏(sparse):不是每条消息都有索引项,而是「写入累计 N 字节才追加一个索引项」。N =
index.interval.bytes,默认 4096。 - 相对 Offset:用
relative_offset = absolute_offset - base_offset节省空间。一个 Segment 上限通常 1GB,相对 offset 用 4 字节足够。
.index 的物理上限:log.index.size.max.bytes,默认 10MB。一个 Segment 写到 .log 1GB 或 .index 10MB,都会触发 roll。
4.3 .timeindex:Timestamp → Offset 稀疏索引
结构:固定长度的 (timestamp, relative_offset) 对,每对 12 字节(8+4)。
+----------------+----------+
| ts=170000_0000 | relOff=0 | ← 时间戳 1700000000.000s 起的第一条消息的相对 offset
+----------------+----------+
| ts=170000_0050 | relOff=23|
+----------------+----------+
| ts=170000_0100 | relOff=51|
+----------------+----------+作用:支持按时间消费(consumer.offsetsForTimes())。例如想消费「昨天 8 点之后的所有消息」,先用时间戳在 .timeindex 二分查到对应 Offset,再用 .index 找 .log 位置。
注意:时间戳支持两种语义,由 Topic 配置 message.timestamp.type 控制:
CreateTime(默认):Producer 端打的时间戳。LogAppendTime:Broker 收到消息时打的时间戳。
CreateTime 下时间戳可能不严格单调(Producer 时钟漂移),Kafka 会在 .timeindex 里只追加「比上一个大」的 timestamp(这叫 max timestamp),保证二分查找正确。
4.4 .snapshot:幂等/事务快照
全名:Producer ID Snapshot
作用:保存某个 Offset 之前所有活跃的 Producer ID(PID)+ Epoch + 序列号 的快照。
为什么需要:
- 幂等 Producer 靠
(PID, partition, sequence)三元组去重。Broker 必须知道每个 PID 的最新 sequence,才能判断新消息是否重复。 - 事务 Producer 靠 PID + Epoch 防止「僵尸 Producer」。
- Broker 重启时不可能从头扫日志重建状态,所以每个 Segment 滚动时都打一份快照,重启后只要从最近的快照 + 后续日志重放即可。
生成时机:Segment roll 时、min.compaction.lag.ms 触发时、leader 切换时。
4.5 leader-epoch-checkpoint:Leader 任期历史
结构:纯文本文件,每行一个 (epoch, start_offset) 对:
0
3
0 0
1 5678
2 12300第 1 行是版本号 0,第 2 行是条目数 3,后面三行依次是:
(epoch=0, start_offset=0):Epoch 0 从 offset 0 开始。(epoch=1, start_offset=5678):第一次 Leader 切换发生在 5678,新 Leader 用 epoch=1。(epoch=2, start_offset=12300):第二次切换。
为什么重要:这是 Kafka 0.11 引入的Leader Epoch 机制,用来解决「HW 截断导致的数据不一致」。第 9 章会详讲,这里只要记住「它是 Kafka 副本一致性的核心」。
4.6 其它辅助文件(了解即可)
partition.metadata:KRaft 之后增加,存 Topic ID(UUID)+ 版本号。recovery-point-offset-checkpoint(broker 级):每个分区已经 fsync 到磁盘的 offset。Broker crash 后从这里之后开始恢复。replication-offset-checkpoint(broker 级):每个分区的 HW,重启后 ISR 沟通的起点。cleaner-offset-checkpoint(broker 级):Log Cleaner 已经压缩到哪里。
5. 稀疏索引 + 二分查找:怎么从 Offset 找消息
5.1 全流程图示
假设要查 Offset = 7050 的消息,目录里 base offset 有 0, 5678, 12346。
Step 1 根据 base offset 二分 → 选中 Segment 5678
Step 2 在 .index 里二分 (relOff = 7050 - 5678 = 1372)
.index 最近的小于等于 1372 的项是 (relOff=1300, pos=131072)
Step 3 打开 .log 文件,seek 到 position 131072
Step 4 从 131072 开始顺序读 RecordBatch,直到 batch 内的某条消息 offset == 7050
Step 5 返回该 record 的 value总 I/O:1 次 .index 二分(mmap 的,几乎全 PageCache)+ 1 次 .log 顺序读(最多读 index.interval.bytes = 4KB)。
.index (按 relOff 升序、定长 8B/项)
+--------+--------+
| relOff | pos |
+--------+--------+
| 0 | 0 |
| 300 | 4096 |
| 650 | 8192 |
| 900 | 12288 |
| 1300 |131072 | ←── 找到这一项(最大的 ≤ 1372)
| 1700 |172000 |
| 2100 |221000 |
+--------+--------+
.log (从 pos=131072 开始顺序扫)
+----------------------------------------------------------------+
| ... | RecordBatch(off=1300..1370, len=...) | RecordBatch(off=1371..1450) | ...
+----------------------------------------------------------------+
↑
这里包含 relOff=1372(即 absolute=7050)5.2 为什么是「稀疏」而不是「稠密」
如果每条消息都建索引,索引文件会和数据文件一样大。Kafka 选择「默认每 4KB 数据建一个索引项」,因为:
- 现代磁盘读 4KB 几乎是「白嫖」(一个磁盘扇区 + PageCache)。
- 索引可以放进 PageCache 甚至 mmap 全驻内存,二分查找极快。
- 索引体积 ≈
.log大小 / 4096 × 8 字节 ≈ 0.2%。
index.interval.bytes 调小 → 索引更密、查找更快、但索引体积大;调大反之。默认 4096 是经过调优的最优解。
5.3 mmap 与 .index
Kafka 用 mmap() 把 .index 文件直接映射进内存(详见第 8 章 §6.4)。所有读写都通过指针操作,没有 read/write 系统调用,零拷贝直读。
6. 消息格式:V0 / V1 / V2 演进
6.1 三代格式简史
| 版本 | 出现时间 | 关键变化 |
|---|---|---|
| V0 | Kafka 0.7 ~ 0.9 | 最早格式,每条消息独立头 |
| V1 | Kafka 0.10 | 加入 timestamp 字段 |
| V2 | Kafka 0.11+ | RecordBatch 概念,header / 事务 / 幂等元数据,多条消息共享 batch 头,体积大幅缩小 |
现代生产环境只用 V2。下面所有讨论默认 V2。
6.2 V2 RecordBatch 字节布局
一个 RecordBatch 的物理结构:
RecordBatch (1 个 batch)
+-------------------------------------------------------------+
| baseOffset (int64, 8B) batch 第一条 offset |
| batchLength (int32, 4B) 后面所有字节的总长 |
| partitionLeaderEpoch (int32, 4B) |
| magic (int8, 1B) = 2 表示 V2 |
| crc (int32, 4B) 后续字节的 CRC32C |
| attributes (int16, 2B) 压缩类型/timestampType/事务/幂等 等 |
| lastOffsetDelta (int32, 4B) batch 内最大 offset - baseOffset |
| firstTimestamp (int64, 8B) batch 第一条 timestamp |
| maxTimestamp (int64, 8B) batch 内最大 timestamp |
| producerId (int64, 8B) PID(幂等/事务) |
| producerEpoch (int16, 2B) Epoch |
| baseSequence (int32, 4B) batch 第一条 sequence |
| recordsCount (int32, 4B) batch 包含的 record 数 |
| ----- header 共 61B ----- |
| records (变长,被 compression.type 压缩) : |
| Record 1 | Record 2 | ... | Record N |
+-------------------------------------------------------------+每条 Record 的内部结构(变长 zigzag 编码节省空间):
Record
+----------------------------------------+
| length (varint) | 整条 record 字节数
| attributes (int8, 1B, 保留) |
| timestampDelta (varint) | 相对 firstTimestamp
| offsetDelta (varint) | 相对 baseOffset
| keyLength (varint) |
| key (bytes) |
| valueLength (varint) |
| value (bytes) |
| headersCount (varint) |
| headers (变长): |
| key (string) + value (bytes) ... |
+----------------------------------------+6.3 V2 vs V0/V1 的优势
- 共享元数据:一个 batch 里所有消息共享
firstTimestamp/baseOffset,每条消息只存 delta(变长 1~4 字节),同样数据 V2 比 V1 小 30% 以上。 - 批量压缩:整个 records 段一起压缩(gzip/lz4/snappy/zstd),压缩率大幅提升(重复字段更多)。
- 支持 header:Producer 可以挂任意 K-V 头部(例如
traceId、spanId),不必塞进 value。 - 幂等 + 事务原生支持:
producerId/producerEpoch/baseSequence在 batch 头里,Broker 一眼就能去重。
6.4 attributes 字段位图
V2 的 attributes(2 字节 = 16 位):
| 位 | 含义 |
|---|---|
| 0-2 | 压缩类型(0=none, 1=gzip, 2=snappy, 3=lz4, 4=zstd) |
| 3 | timestampType(0=CreateTime, 1=LogAppendTime) |
| 4 | isTransactional(事务消息) |
| 5 | isControlBatch(控制类消息:commit/abort marker) |
| 6 | hasDeleteHorizonMs(compaction 用) |
| 7-15 | 保留 |
控制类 batch(isControlBatch=1)是事务的 commit/abort 标记,consumer 用 read_committed 隔离级别时会看到这个 batch 但不返回给业务。
7. kafka-dump-log.sh 实战解码
7.1 找到要看的文件
bash
ls -lh /var/lib/kafka/data/learn.07.storage-0/模拟输出:
total 12M
-rw-r--r-- 1 kafka kafka 12M Apr 17 10:23 00000000000000000000.log
-rw-r--r-- 1 kafka kafka 10M Apr 17 10:23 00000000000000000000.index
-rw-r--r-- 1 kafka kafka 10M Apr 17 10:23 00000000000000000000.timeindex
-rw-r--r-- 1 kafka kafka 13 Apr 17 10:23 leader-epoch-checkpoint
-rw-r--r-- 1 kafka kafka 43 Apr 17 10:23 partition.metadata7.2 dump .log
bash
kafka-dump-log.sh \
--files /var/lib/kafka/data/learn.07.storage-0/00000000000000000000.log \
--print-data-log模拟输出(截取前 2 个 batch,每 batch 各 5 条消息):
Dumping /var/lib/kafka/data/learn.07.storage-0/00000000000000000000.log
Starting offset: 0
baseOffset: 0 lastOffset: 4 count: 5 baseSequence: 0 lastSequence: 4 producerId: 1001 producerEpoch: 0 partitionLeaderEpoch: 0 isTransactional: false isControl: false position: 0 CreateTime: 1713344123456 size: 247 magic: 2 compresscodec: lz4 crc: 1234567890 isvalid: true
| offset: 0 CreateTime: 1713344123456 keySize: 6 valueSize: 18 sequence: 0 headerKeys: [traceId] key: ord_01 payload: {"amount":12.5}
| offset: 1 CreateTime: 1713344123457 keySize: 6 valueSize: 18 sequence: 1 headerKeys: [traceId] key: ord_02 payload: {"amount":34.0}
| offset: 2 CreateTime: 1713344123458 keySize: 6 valueSize: 19 sequence: 2 headerKeys: [traceId] key: ord_03 payload: {"amount":120.5}
| offset: 3 CreateTime: 1713344123459 keySize: 6 valueSize: 18 sequence: 3 headerKeys: [traceId] key: ord_04 payload: {"amount":99.9}
| offset: 4 CreateTime: 1713344123460 keySize: 6 valueSize: 19 sequence: 4 headerKeys: [traceId] key: ord_05 payload: {"amount":222.0}
baseOffset: 5 lastOffset: 9 count: 5 baseSequence: 5 lastSequence: 9 producerId: 1001 producerEpoch: 0 partitionLeaderEpoch: 0 isTransactional: false isControl: false position: 247 CreateTime: 1713344123500 size: 251 magic: 2 compresscodec: lz4 crc: 9876543210 isvalid: true
| offset: 5 CreateTime: 1713344123500 keySize: 6 valueSize: 18 sequence: 5 headerKeys: [traceId] key: ord_06 payload: {"amount":42.0}
...逐字段解读:
baseOffset / lastOffset / count:这个 batch 是 [baseOffset, lastOffset],共 5 条。producerId=1001 producerEpoch=0:幂等 Producer 的身份。isTransactional: false:非事务消息。position: 0:这个 batch 在.log文件里的字节起点。CreateTime: 1713344123456:Producer 端时间戳。size: 247:整个 batch 的字节数(注意这就是 §6.2 的 batchLength + 12B 头)。compresscodec: lz4:批量压缩用的算法。- 下面以
|开头的行是 batch 内每条 record 的解开后展示。
7.3 dump .index
bash
kafka-dump-log.sh \
--files /var/lib/kafka/data/learn.07.storage-0/00000000000000000000.index模拟输出:
Dumping /var/lib/kafka/data/learn.07.storage-0/00000000000000000000.index
offset: 23 position: 4096
offset: 51 position: 8192
offset: 82 position: 12288
offset: 110 position: 16384
offset: 145 position: 20480
...每行就是稀疏索引的一项 (relative_offset, position),但工具会自动还原成 absolute_offset(这里 base offset = 0 所以等于相对值)。
7.4 dump .timeindex
bash
kafka-dump-log.sh \
--files /var/lib/kafka/data/learn.07.storage-0/00000000000000000000.timeindex模拟输出:
Dumping /var/lib/kafka/data/learn.07.storage-0/00000000000000000000.timeindex
timestamp: 1713344123456 offset: 0
timestamp: 1713344124501 offset: 23
timestamp: 1713344125678 offset: 51
timestamp: 1713344126890 offset: 82
...按时间戳升序,每条对应到一个 offset。
7.5 校验 CRC
bash
kafka-dump-log.sh \
--files /var/lib/kafka/data/learn.07.storage-0/00000000000000000000.log \
--deep-iteration--deep-iteration 会逐条解开 RecordBatch 并校验每条 Record 的 CRC,发现损坏会打印 corrupt: true。生产环境磁盘异常排查必备。
8. Segment 滚动条件
active Segment 何时「封口」开新一个?以下任意一条满足都触发:
| 配置 | 默认值 | 含义 |
|---|---|---|
log.segment.bytes | 1 GiB(1073741824) | .log 文件达到此大小 |
log.roll.ms | 168 h(7 天) | 距离首条消息时间到达此时长 |
log.roll.hours | 168 h | 上面字段的 hour 版本,二选一 |
log.index.size.max.bytes | 10 MiB | .index 文件到达此大小 |
8.1 实验:把 segment 改小看一下滚动
bash
# 创建一个 segment 上限只有 1MB 的 Topic
kafka-topics.sh --bootstrap-server kafka:9092 --create \
--topic learn.07.tinyseg \
--partitions 1 --replication-factor 1 \
--config segment.bytes=1048576
# 写入 5MB 数据
for i in $(seq 1 50000); do
echo "msg $i $(date +%s%N)"
done | kafka-console-producer.sh --bootstrap-server kafka:9092 --topic learn.07.tinyseg
# 看磁盘
ls -lh /var/lib/kafka/data/learn.07.tinyseg-0/期望输出(5 个左右的 Segment):
00000000000000000000.log 1.0M
00000000000000000000.index 16K
00000000000000000000.timeindex 24K
00000000000000010234.log 1.0M
00000000000000010234.index 16K
00000000000000010234.timeindex 24K
00000000000000020678.log 1.0M
...
00000000000000041000.log 100K ← active segment(不到 1MB)8.2 滚动的代价
- 打开新 fd:每次 roll 都新建
.log/.index/.timeindex/.snapshot4 个文件。 - PageCache 抛弃:旧 segment 的 PageCache 仍在,但访问稀少会被系统回收。
log-cleaner触发:compaction 类 Topic 的 cleaner 只清理「已 roll 的 segment」,所以 active segment 永远不会被清理。
所以 segment.bytes 不能调太小,否则文件爆炸;也不能太大,否则被删 / 被 compaction 的时延太大。1GB 是经过调优的好默认。
9. 配套配置全览(Topic 级)
bash
kafka-configs.sh --bootstrap-server kafka:9092 \
--entity-type topics --entity-name learn.07.storage \
--describe生产推荐配置参考:
properties
# Segment 滚动
segment.bytes=1073741824 # 1GB
segment.ms=604800000 # 7 天
# 索引
index.interval.bytes=4096 # 每 4KB 一个索引项
segment.index.bytes=10485760 # 索引文件上限 10MB
# 保留与压缩
retention.ms=259200000 # 保留 3 天
retention.bytes=-1 # 不按大小限制
cleanup.policy=delete # 或 compact / compact,delete
# 消息体
max.message.bytes=1048576 # 单条消息 1MB
message.timestamp.type=CreateTime
message.timestamp.difference.max.ms=922337203685477580710. 与 RabbitMQ / RocketMQ / Pulsar 的对比
| 维度 | Kafka | RabbitMQ | RocketMQ | Pulsar |
|---|---|---|---|---|
| 存储模型 | 每 Partition 独立日志,Segment 滚动 | Erlang Mnesia + 持久化(每 Queue 一个文件) | 统一 CommitLog + 每 Queue 一个 ConsumeQueue 索引 | Bookkeeper Ledger(多 Bookie 分片) |
| 顺序消费 | .log 顺序读 + .index 二分定位 | Queue 内顺序消费 | CommitLog 顺序写、ConsumeQueue 顺序读 | Ledger 内顺序读 |
| 时间索引 | .timeindex | 无(用消息 TTL 整批删) | 有 ConsumeQueue 的 storeTimestamp | 有 |
| 索引方式 | 稀疏 + mmap | 内存 hash 表 | ConsumeQueue 是稠密索引(每条 20B) | Bookkeeper 内部 LSM |
| 多消费者读同一份数据 | ✅ 各自维护 Offset | ❌(每 Consumer 一个 Queue) | ✅ ConsumeQueue 抽象 | ✅ Cursor |
| 压缩粒度 | 整个 RecordBatch 一起压 | 单条消息 | 单条消息 | 整个 entry |
一句话:Kafka 把「日志的本质」做到了极致——追加写 + 稀疏索引 + mmap + 批量压缩;其它 MQ 在不同环节做了取舍。
11. 小结
存储与日志格式
├─ 日志目录布局
│ ├─ log.dirs:可配多盘
│ ├─ topic-partition 目录
│ └─ Segment:base offset 命名
├─ 五大文件
│ ├─ .log 数据(RecordBatch 拼接)
│ ├─ .index Offset → Position 稀疏索引
│ ├─ .timeindex Timestamp → Offset 稀疏索引
│ ├─ .snapshot 幂等/事务快照
│ └─ leader-epoch-checkpoint
├─ V2 RecordBatch
│ ├─ 共享 baseOffset / firstTimestamp 节省空间
│ ├─ 批量压缩 → 压缩率高
│ ├─ 支持 header
│ └─ 原生支持幂等 + 事务(PID/Epoch/Seq)
├─ 查找过程
│ └─ 文件名二分 → .index 二分 → .log 顺序扫
├─ Segment 滚动
│ ├─ segment.bytes 默认 1GB
│ ├─ segment.ms / log.roll.ms 默认 7 天
│ └─ index 满 10MB 也滚
└─ 实战工具
└─ kafka-dump-log.sh --files xxx.log --print-data-log12. 面试高频题(7 题)
Q1:Kafka 的 .log、.index、.timeindex 三个文件分别是什么?怎么配合查消息?
考察点:存储结构、稀疏索引。
答案:
.log:真正的消息数据,由若干 RecordBatch 顺序追加。.log文件名是 base offset(这个 segment 第一条消息的 offset)。.index:Offset → Position 的稀疏索引,每对 8 字节 (relative_offset, position)。默认每写入index.interval.bytes=4KB 数据追加一项。.timeindex:Timestamp → Offset 的稀疏索引,每对 12 字节 (timestamp, relative_offset)。支持consumer.offsetsForTimes()按时间消费。- 查询流程(按 Offset 找消息):
- 在 partition 目录下,根据 base offset 二分定位到 segment。
- 在该 segment 的
.index里二分找最大的「relative_offset ≤ target」,拿到position。 - 在
.log里 seek 到 position,顺序往下扫最多 4KB,定位到目标 offset 的 RecordBatch。
- 加分项:
.index是mmap进内存的,二分查找几乎是纯内存操作。- 「稀疏」是为了让索引文件小(约 0.2% of
.log),同时单次顺序扫 ≤ 4KB 几乎免费。 - 时间戳查询稍贵:先在
.timeindex二分找 offset,再走上面流程,一次查询 2 次二分 + 1 次顺序扫。
Q2:Segment 是什么?为什么要分段?
考察点:日志分段思想。
答案:
- 定义:每个 Partition 的日志被切成多个 Segment,每个 Segment 包含
.log + .index + .timeindex + .snapshot一组文件。 - 为什么要分段:
- 只追加单一巨大文件,删除老数据要么改文件中间内容(不可能 / 性能爆炸),要么 truncate 头部(破坏 mmap)。分段后删除 = 直接
unlink整个 segment,O(1) 操作。 - 索引体积可控:每个 segment 的
.index≤ 10MB,能完整 mmap。 - 故障恢复快:只需重放最后一个 segment,前面的不动。
- Compaction 友好:log cleaner 只处理 roll 后的 segment,不影响正在写的 active segment。
- 只追加单一巨大文件,删除老数据要么改文件中间内容(不可能 / 性能爆炸),要么 truncate 头部(破坏 mmap)。分段后删除 = 直接
- 滚动条件(任意满足即触发):
log.segment.bytes(默认 1GB)log.roll.ms或log.roll.hours(默认 7 天)log.index.size.max.bytes(默认 10MB)
- active segment:当前正在写的 segment,永远不会被删除/压缩,所以 retention 至少保留 1 个 segment 的时间。
- 加分项:提到 segment 文件名是 20 位 0 填充的 base offset,字典序 = 数值序,方便 ls 排序。
Q3:.snapshot 文件是干什么的?为什么需要它?
考察点:幂等/事务、Broker 启动恢复。
答案:
- 作用:保存某 offset 之前所有活跃 Producer 的 PID + Epoch + 序列号 + 事务状态快照(即 Producer State Snapshot)。
- 为什么需要:
- 幂等 Producer 用
(PID, partition, sequence)三元组去重;事务 Producer 用 PID + Epoch 防僵尸。Broker 必须维护一份「每个 PID 当前 sequence 是多少」的状态。 - Broker 重启时不可能从头扫整个分区日志重建状态(GB ~ TB 级)。
- 所以每次 segment roll 时,把当前 Producer 状态快照写到
.snapshot,重启后从「最近的 snapshot + 后续日志」回放即可。
- 幂等 Producer 用
- 类似机制:可以理解为 Broker 内部「为分区状态打的 checkpoint」,思路与 Raft 的 snapshot、Redis 的 RDB 一致。
- 生成时机:segment roll、leader 切换、Broker 优雅关闭。
- 加分项:提到 KRaft 自身的
__cluster_metadata也大量使用 snapshot 机制;Broker 启动慢时常常是 snapshot 太老 + 日志太长导致重放时间长,可调log.roll.ms让 snapshot 更频繁。
Q4:V0/V1/V2 三种消息格式的区别?V2 解决了什么问题?
考察点:Kafka 演进史、协议优化。
答案:
- V0(Kafka 0.7 ~ 0.9):每条消息独立 header。结构:
crc + magic + attributes + keyLen + key + valueLen + value。无 timestamp。 - V1(Kafka 0.10):在 V0 基础上加 8 字节 timestamp。其它一样。
- V2(Kafka 0.11+):引入 RecordBatch 概念,多条 Record 共享一个 batch 头(61B),每条 Record 用 varint 编码 delta 字段,整体压缩。结构上完全重写。
- V2 解决的问题:
- 协议体积:V0/V1 每条消息 30+ 字节 overhead,V2 共享后每条 ~ 6 字节。同样数据量空间小 30%+。
- 批量压缩:V0/V1 是「先压缩再封 batch」,每条仍有独立 header;V2 整个 records 段一起压,压缩率大幅提升(重复 header 字段被压掉)。
- 支持事务/幂等:batch 头里直接有
producerId + producerEpoch + baseSequence,Broker 一眼判断重复或事务归属。 - 支持 header:每条 Record 可挂任意 K-V 头部(traceId、tenantId 等),不必塞进 value。
- 支持控制类消息:commit/abort marker 用
isControlBatch标识,consumer 自动过滤。
- 兼容性:现代 Broker 都支持读 V0/V1,但写都是 V2。Producer/Consumer 协商最小公共版本。
- 加分项:提到 V2 消息的
lastOffsetDelta字段让 Broker 知道一个 batch 跨多少 offset,可以一次校验整个 batch 的 sequence 连续性。
Q5:Kafka 怎么按时间戳消费?时间戳从哪来?
考察点:.timeindex、message.timestamp.type。
答案:
- 时间戳来源(由 Topic 配置
message.timestamp.type决定):CreateTime(默认):Producer 端打的时间戳(new ProducerRecord(topic, key, value)不传则用System.currentTimeMillis())。LogAppendTime:Broker 收到消息时的本地时间,自动覆盖 Producer 端时间戳。
- 存储:每条消息的 timestamp 在 V2 的 RecordBatch header 里,用 delta 编码。Broker 维护
.timeindex文件,每 4KB 数据追加一个(timestamp, relative_offset)。 - 查询流程(
consumer.offsetsForTimes(Map<TopicPartition, Long>)):- 客户端发
ListOffsetsRequest(timestamp)给 Broker。 - Broker 在
.timeindex二分找 ≥ 该时间戳的最小项,得到一个 offset。 - 再用
.index校准到 RecordBatch 起点 offset,返回给客户端。 - 客户端用
consumer.seek(partition, offset)跳到该位置。
- 客户端发
- 注意点:
CreateTime下时间戳可能不严格单调(Producer 时钟漂移、不同机器不同步)。.timeindex只追加「比上一项大」的时间戳(max 语义),保证二分正确。- 如果 Producer 故意打了一个未来时间戳,会污染整个分区的
.timeindex,需要message.timestamp.difference.max.ms限制 Producer 与 Broker 时差。
- 加分项:用时间戳消费比 offset 多一次 IO(
.timeindex二分),所以高频场景仍建议保存 offset。
Q6:Kafka 怎么决定一条消息是否过期?删除是怎么发生的?
考察点:retention 机制、segment 删除。
答案:
- 两种保留维度:
retention.ms(默认 7 天):消息时间戳早于「现在 - retention.ms」即过期。retention.bytes:分区总大小超过此值开始删旧 segment。-1表示不限。
- 删除单位是 segment,不是消息:Kafka 永远不会删 segment 中间的某条消息,只会按 segment 整体删除。
- 触发流程:
- Broker 后台线程
kafka-log-retention-task(默认每 5 分钟一次,由log.retention.check.interval.ms控制)扫描所有 segment。 - 取出最大时间戳 < (now - retention.ms) 的所有非 active segment,标记为待删。
- 重命名为
.deleted后缀,等file.delete.delay.ms(默认 1 分钟)后真正 unlink。
- Broker 后台线程
- active segment 永远不删:所以 retention 实际是「保留 ≥ retention.ms 但可能更长(最多多 1 个 segment 的时间)」。
cleanup.policy=compact:另一个语义,按 Key 保留最新值(详见第 14 章 Log Compaction)。- 加分项:
retention.bytes优先级低于retention.ms:先删过期的,再按大小删。- 大消息 + 短 retention 容易出现「消息逻辑过期但 segment 还没满,仍占着磁盘」。
- 如果 disk usage 报警了,临时降
retention.ms后会立刻触发清理。
Q7:为什么 Kafka 不用 B-Tree 这种数据库的索引?
考察点:存储模型对比、设计哲学。
答案:
- 数据库的需求是「任意 Key 都能快速查找 + 范围扫」:B-Tree / B+Tree 是为了支持任意点查 + 任意范围 + 频繁更新。
- Kafka 的需求是「只按 Offset 查 + 顺序扫」:Offset 是单调递增的整数,不会更新;查询模式 99% 是「从某 Offset 开始顺序往后读 N 条」。
- 如果用 B-Tree:
- 维护代价高(每条消息插入都要分裂 / 合并节点,破坏顺序写)。
- 索引体积大(每条都要建索引)。
- 顺序扫退化(B-Tree 叶节点链表跳来跳去,磁头乱跑)。
- 稀疏索引的优势:
- 索引体积极小(≈ 0.2%),整个索引能 mmap 进内存。
- 写入完全顺序(追加到
.log末尾、追加到.index末尾)。 - 查询:1 次二分 + 1 次 ≤4KB 顺序扫,IOPS 极少。
- 本质对比:
- 数据库 = OLTP,强事务 + 任意 Key 检索,B-Tree 最优。
- Kafka = 提交日志,只追加 + 按位置读,稀疏索引 + mmap 最优。
- 加分项:提到 RocksDB 的 LSM 也是为了「优化写 + 牺牲点查」的折中;Kafka 比 LSM 还激进,连 compaction 都仅在
cleanup.policy=compact时才做(默认 delete 模式只 unlink 整个 segment)。
本章配套:
07_storage/demo.html:Segment 滚动动画 + Offset → Position 二分查找动画 + V2 RecordBatch 字节级展示。07_storage/init.sh:建learn.07.storageTopic 并写入 1 万条消息,方便观察文件分布。07_storage/code/dump_segment.sh:封装kafka-dump-log.sh的一键脚本。07_storage/code/explore_segment.py:写大量消息后,用 Python 列出 segment 文件、统计大小、解析索引项。
🔗 延伸阅读
- 第 6 章 Topic 设计 —— 分区数 → 文件数的连锁反应。
- 第 8 章 高吞吐底层原理 —— mmap、零拷贝、PageCache 都依赖这一章的存储模型。
- 第 9 章 副本机制与 ISR ——
leader-epoch-checkpoint的完整故事。 - 第 14 章 Log Compaction ——
.snapshot的另一类应用,以及怎么按 Key 压缩。
🎬 可视化演示
演示加载缓慢或样式异常?点此在新标签页打开 ↗
💻 示例代码
bash
#!/usr/bin/env bash
# dump_segment.sh — 一键封装 kafka-dump-log.sh
#
# 用法:
# bash dump_segment.sh <segment文件路径>
#
# 自动检测文件后缀,分别解码 .log / .index / .timeindex / .snapshot:
# - .log : --print-data-log + --deep-iteration(含 batch + record 详情 + CRC 校验)
# - .index : 列出 (offset, position) 稀疏索引项
# - .timeindex : 列出 (timestamp, offset) 时间索引项
# - .snapshot : 列出 ProducerId / Epoch / 序列号快照
#
# 示例:
# bash dump_segment.sh /var/lib/kafka/data/learn.07.storage-0/00000000000000000000.log
# bash dump_segment.sh /var/lib/kafka/data/learn.07.storage-0/00000000000000000000.index
# bash dump_segment.sh /var/lib/kafka/data/learn.07.storage-0/00000000000000000000.timeindex
set -euo pipefail
if [ $# -lt 1 ]; then
echo "用法: $0 <segment文件>"
echo
echo "示例:"
echo " $0 /var/lib/kafka/data/learn.07.storage-0/00000000000000000000.log"
echo " $0 /var/lib/kafka/data/learn.07.storage-0/00000000000000000000.index"
echo " $0 /var/lib/kafka/data/learn.07.storage-0/00000000000000000000.timeindex"
exit 1
fi
FILE="$1"
if [ ! -f "$FILE" ]; then
echo "❌ 文件不存在: $FILE" >&2
exit 1
fi
# 检查 kafka-dump-log.sh 是否在 PATH
if ! command -v kafka-dump-log.sh >/dev/null 2>&1; then
echo "❌ 找不到 kafka-dump-log.sh,请确认 Kafka bin 目录已加入 PATH" >&2
echo " 常见位置: /opt/kafka/bin/kafka-dump-log.sh" >&2
exit 1
fi
EXT="${FILE##*.}"
echo "=============================================================="
echo " kafka-dump-log.sh 解码工具"
echo " 文件: $FILE"
echo " 类型: .$EXT"
echo " 大小: $(ls -lh "$FILE" | awk '{print $5}')"
echo "=============================================================="
echo
case "$EXT" in
log)
echo "[模式] .log 数据文件"
echo "[选项] --print-data-log --deep-iteration(解码每条 record + 校验 CRC)"
echo
kafka-dump-log.sh \
--files "$FILE" \
--print-data-log \
--deep-iteration
;;
index)
echo "[模式] .index 偏移量稀疏索引"
echo "[说明] 每行:offset → file position(默认每 4KB 一项)"
echo
kafka-dump-log.sh --files "$FILE"
;;
timeindex)
echo "[模式] .timeindex 时间稀疏索引"
echo "[说明] 每行:timestamp → offset"
echo
kafka-dump-log.sh --files "$FILE"
;;
snapshot)
echo "[模式] .snapshot 幂等/事务 Producer 状态快照"
echo "[说明] ProducerId / Epoch / 序列号 / 事务状态"
echo
kafka-dump-log.sh --files "$FILE"
;;
*)
echo "⚠️ 未知后缀 .$EXT,将以默认方式 dump"
echo
kafka-dump-log.sh --files "$FILE"
;;
esac
echo
echo "=============================================================="
echo " 完成。"
echo "=============================================================="python
#!/usr/bin/env python3
"""
explore_segment.py
==================
探索 Kafka Topic-Partition 目录里所有 segment 文件,打印:
- 每个 segment 的 base offset / .log 大小 / .index 项数 / .timeindex 项数 / .snapshot 是否存在
- leader-epoch-checkpoint 内容
- 所有文件的总占用
- 演示「按 offset 找 segment」的二分定位
不依赖 Kafka 服务端,仅靠文件名 + 文件大小解析。
用法:
python3 explore_segment.py --part-dir /var/lib/kafka/data/learn.07.storage-0
python3 explore_segment.py --part-dir /var/lib/kafka/data/learn.07.storage-0 --lookup 7050
也可以与 init.sh 配合:
bash init.sh
python3 code/explore_segment.py \
--part-dir /var/lib/kafka/data/learn.07.storage-0 --lookup 5000
"""
from __future__ import annotations
import argparse
import os
import struct
import sys
from dataclasses import dataclass
from pathlib import Path
from typing import List, Optional
# ---------------------------------------------------------------------------
# 常量
# ---------------------------------------------------------------------------
# .index 每项 8 字节: relativeOffset (int32) + position (int32)
INDEX_ENTRY_SIZE = 8
# .timeindex 每项 12 字节: timestamp (int64) + relativeOffset (int32)
TIMEINDEX_ENTRY_SIZE = 12
# ---------------------------------------------------------------------------
# Segment 描述
# ---------------------------------------------------------------------------
@dataclass
class Segment:
base_offset: int
log_path: Path
index_path: Optional[Path]
timeindex_path: Optional[Path]
snapshot_path: Optional[Path]
@property
def log_size(self) -> int:
return self.log_path.stat().st_size if self.log_path.exists() else 0
@property
def index_entries(self) -> int:
if self.index_path and self.index_path.exists():
return self.index_path.stat().st_size // INDEX_ENTRY_SIZE
return 0
@property
def timeindex_entries(self) -> int:
if self.timeindex_path and self.timeindex_path.exists():
return self.timeindex_path.stat().st_size // TIMEINDEX_ENTRY_SIZE
return 0
@property
def has_snapshot(self) -> bool:
return self.snapshot_path is not None and self.snapshot_path.exists()
# ---------------------------------------------------------------------------
# 扫描目录
# ---------------------------------------------------------------------------
def scan_segments(part_dir: Path) -> List[Segment]:
if not part_dir.exists():
raise FileNotFoundError(f"分区目录不存在: {part_dir}")
log_files = sorted(part_dir.glob("*.log"))
segs: List[Segment] = []
for log in log_files:
base_str = log.stem # 文件名(去后缀)
try:
base = int(base_str)
except ValueError:
continue # 不是 segment 命名(例如 leader-epoch-checkpoint)
index = part_dir / f"{base_str}.index"
timeindex = part_dir / f"{base_str}.timeindex"
snapshot = part_dir / f"{base_str}.snapshot"
segs.append(Segment(
base_offset=base,
log_path=log,
index_path=index if index.exists() else None,
timeindex_path=timeindex if timeindex.exists() else None,
snapshot_path=snapshot if snapshot.exists() else None,
))
return segs
# ---------------------------------------------------------------------------
# 打印 Segment 概要
# ---------------------------------------------------------------------------
def print_summary(part_dir: Path, segs: List[Segment]) -> None:
print(f"\n📁 {part_dir}")
print(f" 共发现 {len(segs)} 个 segment\n")
if not segs:
return
print(f" {'idx':<4} {'base offset':>13} {'.log':>10} {'.index':>8} "
f"{'.timeindex':>10} {'snapshot':>9} {'state':>8}")
print(" " + "-" * 78)
for i, s in enumerate(segs):
is_active = (i == len(segs) - 1)
state = "ACTIVE" if is_active else "SEALED"
print(
f" {i:<4} {s.base_offset:>13} "
f"{human(s.log_size):>10} "
f"{s.index_entries:>8} "
f"{s.timeindex_entries:>10} "
f"{'✓' if s.has_snapshot else '·':>9} "
f"{state:>8}"
)
total_log = sum(s.log_size for s in segs)
total_idx = sum((s.index_path.stat().st_size if s.index_path else 0) for s in segs)
total_tidx = sum((s.timeindex_path.stat().st_size if s.timeindex_path else 0) for s in segs)
print()
print(f" 📊 .log 总大小 : {human(total_log)}")
print(f" 📊 .index 总大小 : {human(total_idx)} ({total_idx/total_log*100:.3f}% of .log)" if total_log else "")
print(f" 📊 .timeindex 总大小: {human(total_tidx)}")
# ---------------------------------------------------------------------------
# 解析 leader-epoch-checkpoint
# ---------------------------------------------------------------------------
def print_leader_epoch(part_dir: Path) -> None:
f = part_dir / "leader-epoch-checkpoint"
if not f.exists():
return
print(f"\n📜 leader-epoch-checkpoint ({f})")
lines = f.read_text().strip().splitlines()
if len(lines) < 2:
print(" (文件为空或格式异常)")
return
version = lines[0]
count = lines[1]
print(f" version = {version}, entries = {count}")
print(f" {'epoch':>6} {'start_offset':>14}")
for line in lines[2:]:
parts = line.split()
if len(parts) == 2:
print(f" {parts[0]:>6} {parts[1]:>14}")
# ---------------------------------------------------------------------------
# 解析 .index 文件,打印前 N 项
# ---------------------------------------------------------------------------
def print_index_head(seg: Segment, n: int = 10) -> None:
if not seg.index_path or not seg.index_path.exists():
return
print(f"\n🔎 .index 头部({seg.index_path.name},前 {n} 项)")
print(f" {'idx':>4} {'rel_offset':>12} {'abs_offset':>12} {'position':>10}")
with open(seg.index_path, "rb") as f:
data = f.read(n * INDEX_ENTRY_SIZE)
for i in range(0, len(data), INDEX_ENTRY_SIZE):
rel, pos = struct.unpack(">II", data[i:i + INDEX_ENTRY_SIZE])
if rel == 0 and pos == 0 and i > 0:
# 已经到尾部填零段
break
print(f" {i//INDEX_ENTRY_SIZE:>4} {rel:>12} {seg.base_offset + rel:>12} {pos:>10}")
def print_timeindex_head(seg: Segment, n: int = 10) -> None:
if not seg.timeindex_path or not seg.timeindex_path.exists():
return
print(f"\n⏱ .timeindex 头部({seg.timeindex_path.name},前 {n} 项)")
print(f" {'idx':>4} {'timestamp_ms':>16} {'rel_offset':>12} {'abs_offset':>12}")
with open(seg.timeindex_path, "rb") as f:
data = f.read(n * TIMEINDEX_ENTRY_SIZE)
for i in range(0, len(data), TIMEINDEX_ENTRY_SIZE):
ts, rel = struct.unpack(">QI", data[i:i + TIMEINDEX_ENTRY_SIZE])
if ts == 0 and rel == 0 and i > 0:
break
print(f" {i//TIMEINDEX_ENTRY_SIZE:>4} {ts:>16} {rel:>12} {seg.base_offset + rel:>12}")
# ---------------------------------------------------------------------------
# 「按 offset 找 segment」演示
# ---------------------------------------------------------------------------
def lookup_offset(segs: List[Segment], target: int) -> None:
print(f"\n🔍 查找 offset = {target}")
if not segs:
print(" 无 segment 可查")
return
# 二分:找最大的 base_offset ≤ target
lo, hi, found = 0, len(segs) - 1, -1
steps = []
while lo <= hi:
mid = (lo + hi) // 2
steps.append((mid, segs[mid].base_offset))
if segs[mid].base_offset <= target:
found = mid
lo = mid + 1
else:
hi = mid - 1
print(f" 二分步骤(共 {len(steps)} 步): {steps}")
if found == -1:
print(f" ❌ 找不到合适的 segment(目标 offset 太小)")
return
seg = segs[found]
print(f" ✅ 命中 segment[{found}] base_offset={seg.base_offset}")
print(f" 文件: {seg.log_path.name}")
print(f" 下一步:在该 .index 里二分 relative_offset = {target - seg.base_offset}")
print(f" 再到 .log 顺序扫 ≤ index.interval.bytes(默认 4KB)即可命中")
# ---------------------------------------------------------------------------
# 工具
# ---------------------------------------------------------------------------
def human(n: int) -> str:
for unit in ["B", "KB", "MB", "GB"]:
if n < 1024:
return f"{n:.1f}{unit}" if unit != "B" else f"{n}{unit}"
n /= 1024
return f"{n:.1f}TB"
# ---------------------------------------------------------------------------
# 主入口
# ---------------------------------------------------------------------------
def main() -> None:
parser = argparse.ArgumentParser(description="探索 Kafka Topic-Partition 目录的 segment 结构")
parser.add_argument("--part-dir", type=str, required=True,
help="分区目录,如 /var/lib/kafka/data/learn.07.storage-0")
parser.add_argument("--lookup", type=int, default=None,
help="查找某 offset 落在哪个 segment 上(演示二分过程)")
parser.add_argument("--head", type=int, default=10,
help="每个文件打印前 N 项索引,默认 10")
args = parser.parse_args()
part_dir = Path(args.part_dir)
try:
segs = scan_segments(part_dir)
except FileNotFoundError as e:
print(f"❌ {e}", file=sys.stderr)
sys.exit(1)
print_summary(part_dir, segs)
print_leader_epoch(part_dir)
if segs:
# 默认展示第一个 segment 的索引头部
print_index_head(segs[0], args.head)
print_timeindex_head(segs[0], args.head)
if args.lookup is not None:
lookup_offset(segs, args.lookup)
print()
if __name__ == "__main__":
main()