Kafka 日志段与索引内部机制

拆解 Kafka 分区日志的物理实现:段文件命名与滚动、记录批次的二进制布局、偏移索引与时间索引的稀疏结构、二分查找定位、零拷贝 sendfile 与页缓存、磁盘布局与 IO 模式、日志清理压缩的段级粒度、JBOD 与 RAID 的取舍、索引损坏修复与排错要点

Kafka 的高吞吐常被归结为「顺序写磁盘」,但这只是故事的一半。真正让它在同样的磁盘上跑出数量级优势的,是日志段 + 稀疏索引 + 零拷贝这套组合:写入永远是追加,读取靠一层极小的索引把随机查找压缩成一次二分加一次顺序读,而数据从磁盘到网卡全程不经过用户态。理解这三个文件的二进制布局,是诊断「为什么消费变慢了」「为什么磁盘 IO 打满了」这类问题的前提。

1. 分区目录与文件命名

1.1 目录结构

Kafka 的每个分区在 broker 的 log.dirs 下有一个独立目录,命名是 <topic>-<partition>:

/var/lib/kafka/data/
├── orders-0/
│   ├── 00000000000000000000.log        # 日志数据
│   ├── 00000000000000000000.index      # 偏移索引
│   ├── 00000000000000000000.timeindex  # 时间索引
│   ├── 00000000000000000000.snapshot   # 幂等/事务生产者快照
│   ├── 00000000000000000123.log        # 下一个段,基偏移 123
│   ├── 00000000000000000123.index
│   ├── 00000000000000000123.timeindex
│   ├── leader-epoch-checkpoint
│   └── partition.metadata
├── orders-1/
└── __consumer_offsets-0/

同一分区的所有段文件共享一个基偏移(base offset) 前缀,就是这个段第一条消息的绝对 offset,用 20 位十进制零填充。文件名即段的身份,排序文件名等价于排序时间(因为 offset 单调递增)。

leader-epoch-checkpoint 记录每个 leader 任期开始时的 offset,用于副本同步时判断「这个 offset 是哪一任 leader 写的」,是日志截断(truncation)的依据。partition.metadata 存 topic id,与 __cluster_metadata 里的 topic 定义对应。

1.2 段滚动条件

活跃段(active segment)写满后滚动出新段。触发条件有四个:

# 1. 大小:段文件达到阈值即滚动(默认 1GB)
log.segment.bytes=1073741824

# 2. 时间:段创建后多久强制滚动(默认 7 天)
log.roll.ms=604800000

# 3. 索引写满:索引文件达到 log.index.size.max.bytes 时
log.index.size.max.bytes=10485760

# 4. 偏移溢出:offset 达到 Int.MAX 附近时(索引用 4 字节存相对偏移)

log.segment.bytes 是最重要的旋钮。段太小,段数量暴涨、索引文件碎片化、文件句柄吃紧;段太大,日志清理与压缩的粒度变粗(一个段要么整体保留要么整体删除,无法只删一半)。默认 1 GB 是经验平衡点,对低流量 topic 可以调小到 100 MB,让 retention 的时间精度更准。

注意第 4 条:偏移索引用 4 字节 存相对基偏移的差值,所以单段内的 offset 跨度不能超过 2^31。这也是为什么会有「索引写满」这个滚动条件。

1.3 相对偏移的设计

为什么索引里存的是「相对偏移」而不是绝对 offset?因为绝对 offset 是 8 字节 long,而大多数段的跨度远小于 2^31,用 4 字节存相对值就能省一半空间。索引文件因此极小:默认每 4 KB 数据才记一条索引项,1 GB 的段索引只有约 10 MB。

2. 段文件的二进制布局

2.1 记录批次

从 Kafka 0.11 起,.log 文件里不再是单条消息,而是记录批次(RecordBatch)。一个批次是生产者在一次请求里发往同一分区的若干条消息,被 broker 原样落盘:

RecordBatch 布局(v2)
┌─────────────────────────────────────────────┐
│ baseOffset        int64   8B                │
│ batchLength       int32   4B                │
│ partitionLeaderEpoch int32 4B               │
│ magic             int8    1B  (=2)          │
│ crc               uint32  4B  (CRC32C)      │
│ attributes        int16   2B                │
│ lastOffsetDelta   int32   4B                │
│ firstTimestamp    int64   8B                │
│ maxTimestamp      int64   8B                │
│ producerId        int64   8B                │
│ producerEpoch     int16   2B                │
│ baseSequence      int32   4B                │
│ records           int32   4B  (条数)        │
│ ── 以下为逐条记录 ──                        │
│ record...                                   │
└─────────────────────────────────────────────┘

批次头的固定部分共 61 字节。baseOffset 是批次内第一条记录的绝对 offset,批次内后续记录的 offset 是隐式的——通过 baseOffset + i 推导。这是省空间的经典手法:N 条记录只存一个 offset。

attributes 是位标志,低 3 位编码压缩类型(0 无、1 gzip、2 snappy、3 lz4、4 zstd),第 4 位是时间戳类型(创建时间 / 日志追加时间),第 5 位是事务标志,第 6 位是控制批次标志。控制批次(transaction marker)是事务提交/回滚的标记,它们也在日志里,消费端要过滤掉。

2.2 压缩批次与偏移语义

当批次被压缩时,压缩的是 records 部分,批次头不压缩。这意味着 broker 无需解压就能读到 baseOffset、lastOffsetDelta、maxTimestamp——这正是它能高效处理压缩批次的秘密。

对消费端的影响是:一个压缩批次内的所有记录必须整体投递给消费者(除非消费者自己解压后逐条处理),因为它们共享一个批次头。这也是为什么「压缩 + 大批量」会让消费端的 max.poll.records 语义变得微妙——一次 poll 可能拿到一整个大压缩批次。

2.3 稀疏索引的两个文件

.index 与 .timeindex 都不是「每条记录一条索引」,而是每 log.index.interval.bytes(默认 4096 字节)写一条。这就是「稀疏」的含义。

.index(偏移索引)条目格式:8 字节
┌───────────────────────┬───────────────────────┐
│ relativeOffset  int32 │ position       int32  │
│ 4B                    │ 4B                    │
└───────────────────────┴───────────────────────┘
relativeOffset = 绝对 offset - 段基偏移
position       = 该批次在 .log 文件中的字节位置
.timeindex(时间索引)条目格式:12 字节
┌───────────────────────┬───────────────────────┐
│ timestamp       int64 │ relativeOffset int32  │
│ 8B                    │ 4B                    │
└───────────────────────┴───────────────────────┘

.timeindex 只在时间戳单调递增时才写条目。若生产者写入的时间戳乱序(比如用了事件时间且事件延迟),Kafka 会跳过写时间索引项,导致基于时间戳查找退化成顺序扫描——这是「按时间戳查消息很慢」的根因。

3. 索引查找与二分定位

3.1 按 offset 查找的完整路径

消费者要读 offset = 100000 的消息,broker 的处理链路是:

  1. 定位段:在段的跳跃表(ConcurrentSkipListMap)里找到「基偏移 ≤ 100000 的最大段」。
  2. 查偏移索引:把 100000 - baseOffset 作为相对偏移,在 .index 里二分查找,找到「相对偏移 ≤ 目标的最大索引项」,得到起始字节位置。
  3. 顺序扫描:从该字节位置开始,顺序读 .log 里的批次头,比对 baseOffset 与 lastOffsetDelta,直到找到包含目标 offset 的批次。
  4. 返回:从批次里解出目标记录(若压缩则解压批次)。

关键在于第 3 步的扫描量:索引间隔是 4 KB,所以最多扫 4 KB 就能命中。一次随机查找 = 一次二分(内存里的索引,通常已被页缓存)+ 一次 4 KB 顺序读。这就是稀疏索引把随机 IO 转化为近顺序 IO 的原理。

3.2 按时间戳查找

按时间戳查找用 .timeindex,逻辑类似,但有个边界陷阱:时间索引条目指向的是「该时间戳之后第一个记录的位置」,因此找到条目后仍要向前扫描确认。Kafka 的做法是:在 .timeindex 里二分找到第一个 timestamp >= 目标 的条目,取它前一条的位置开始扫描。

如果 .timeindex 因为时间戳乱序而缺少条目,查找会回退到「从段头顺序扫描」,代价陡增。生产上若依赖时间戳查询(如 kafka-consumer-groups --to-datetime 或 offsetsForTimes),必须保证生产端时间戳大致单调。

3.3 索引损坏与重建

索引文件是可重建的派生数据——它们完全能从 .log 推导出来。因此 Kafka 对索引损坏的处理是「删除并重建」:

# 若 broker 日志出现 "Found invalid index file",停止该 broker
# 删除损坏的索引文件(保留 .log!)
rm /var/lib/kafka/data/orders-0/00000000000000000123.index
rm /var/lib/kafka/data/orders-0/00000000000000000123.timeindex

# 重启 broker 后,Kafka 会在加载段时自动重建索引

重建的触发点在 LogSegment.recover():broker 启动扫描 .log,从每个批次的头部重建索引项。大段的重建会拖慢启动,因此不要在启动期间对索引目录做并发操作。

一个常见误操作是「删了 .index 顺手也删 .log」——那等于删数据。索引丢了能重建,日志丢了不可恢复。

3.4 索引间隔调优

log.index.interval.bytes 控制索引密度:

  • 调小(如 1024):索引更密,查找扫描量更小,但索引文件更大、占用更多页缓存。
  • 调大(如 16384):索引更稀,省空间,但每次查找要多扫几 KB。

默认 4096 对绝大多数场景合适。只有在「磁盘随机读性能极差、而页缓存充裕」的场景(如机械盘 + 大内存)才值得调小。相反,SSD 时代顺序读 4 KB 几乎无成本,调大反而省内存。

4. 零拷贝与页缓存

4.1 传统读写的四次拷贝

常规的「读文件 → 发网络」路径涉及四次数据拷贝与两次系统调用:

磁盘 → 内核页缓存 →(copy 1)→ 用户态缓冲区
用户态缓冲区 →(copy 2)→ socket 缓冲区
socket 缓冲区 →(copy 3)→ 网卡(DMA)
外加一次内核读 + 一次内核写 = 2 次上下文切换

对 Kafka 这种「读出来的字节原样发给消费者」的场景,用户态缓冲区完全是多余的中间站。

4.2 sendfile 与 transferTo

Kafka 用 FileChannel.transferTo(),底层走 Linux 的 sendfile 系统调用,把路径压缩成:

磁盘 → 内核页缓存 →(DMA,无 CPU 参与)→ 网卡

用户态完全不参与,拷贝次数从 4 降到 2(其中一次还是 DMA 不占 CPU),上下文切换从 2 次降到 0 次(就数据传输而言)。这就是 Kafka 用普通磁盘就能跑出高吞吐的核心原因之一。

// Kafka 内部的发送路径(简化)
public long writeTo(TransferableChannel channel, long position, int length) {
    // 直接把文件的一段交给网卡,绕过用户态
    return channel.transferFrom(this, position, length);
}

零拷贝的边界必须清楚:它只在「数据无需修改」时可用。一旦需要格式转换(比如老版本 consumer 需要 down-convert 消息格式,或 broker 要做加解密),零拷贝就失效,Kafka 会退化成普通读写路径。这也是为什么「老客户端 + 新 broker」的性能明显更差——broker 被迫解压、转换、重新压缩。SSL/TLS 加密同理,TLS 终结在 broker 上时无法零拷贝。

4.3 页缓存的角色

Kafka 的读写高度依赖操作系统页缓存(page cache),而不是 JVM 堆:

  • 写:消息写入页缓存即返回(flush 由操作系统决定),不主动 fsync(除非 log.flush.interval.messages 强制)。可靠性靠副本复制而非本地刷盘保证。
  • 读:刚写入的消息大概率还在页缓存里,消费者读取命中缓存,根本不碰磁盘——这就是「生产后立即消费」性能极高的原因。

由此推出两条重要结论:

  1. 不要给 JVM 分配过大堆。堆占用的内存无法被页缓存使用,反而挤压了缓存命中率。生产上堆给 6~8 GB 足够,剩下的内存留给页缓存。
  2. log.flush.interval.messages 不要设小。强制 fsync 会打断顺序写、把 IO 变成同步随机写,吞吐断崖式下跌。正确做法是依赖 replication.factor >= 3 与 min.insync.replicas 保证持久性。

4.4 验证零拷贝是否生效

零拷贝失效往往悄无声息——性能掉了,但没有任何报错。可以用几种方式确认:

# 1. 看 CPU 的软中断与系统态占比
#    零拷贝生效时,broker 进程的 sys CPU 占比明显偏低
pidstat -p $(pgrep -f kafka.Kafka) 1

# 2. 用 strace 观察是否走 sendfile(而非 read+write)
strace -f -e trace=sendfile,read,write -p $(pgrep -f kafka.Kafka) 2>&1 | head

# 3. 检查消费者是否触发了格式转换
#    broker 日志里出现 "Down-converted" 即说明零拷贝已失效
grep -i "down-convert" /var/log/kafka/server.log

若 strace 看到大量 read + write 成对出现,而不是 sendfile,基本可以确定零拷贝没生效。常见原因有三:消费者版本过低触发 down-conversion、broker 开启了 SSL 终结、或者配置了某种消息格式转换。

5. 磁盘布局与 IO 模式

5.1 顺序写、随机读

Kafka 的 IO 模式是「写永远顺序,读大多顺序、少量随机」:

  • 生产:追加到活跃段末尾,纯顺序写。
  • 消费:从索引定位后顺序读;每个消费者在自己的 offset 处顺序推进。
  • 副本同步:follower 从 leader 拉取,也是顺序读。

随机 IO 只出现在首次查找(读索引 + 跳转到位置)以及多消费者跨分区读时。因此 Kafka 对磁盘的顺序吞吐远比随机 IOPS 敏感——这也是它能用大容量机械盘的原因,但在多消费者并发、页缓存不足时,随机读会成为瓶颈。

5.2 JBOD 与 RAID

log.dirs 可以配置多个目录,Kafka 会把分区轮询分布到各目录(每个分区整体落在同一个盘上):

# 推荐:JBOD,多个独立盘
log.dirs=/data1/kafka,/data2/kafka,/data3/kafka

这与 RAID 有本质区别:

方式容量吞吐故障影响运维
JBOD各盘之和各盘之和单盘故障只影响其上的分区需副本兜底
RAID 0各盘之和各盘之和任一盘故障全损危险
RAID 10一半较高可容忍单盘成本高

Kafka 官方推荐 JBOD:因为 Kafka 本身有副本机制,单盘故障可以通过副本重建恢复,不必用 RAID 的冗余去换容量损失。RAID 10 的写放大与成本在 Kafka 场景下并不划算。

5.3 挂载参数与文件系统

# 推荐挂载参数(ext4/xfs)
# noatime:不更新访问时间,省一次写
# data=writeback:只保证元数据一致性,牺牲一点崩溃安全性
mount -o noatime,data=writeback /dev/nvme0n1p1 /data1

noatime 对 Kafka 意义明确:读日志不需要记录访问时间,省掉每次读带来的元数据写。XFS 在大文件顺序写与并行 IO 上通常优于 ext4,是 Kafka 的常见选择。

6. 日志清理与压缩的内部视角

6.1 清理线程的粒度

日志清理(delete 策略)以段为单位:只有非活跃段才能被删除,活跃段永远保留。这带来一个常被误解的现象:retention.ms = 1 小时 的 topic,实际数据可能保留数小时——因为活跃段没满、没滚动,就一直删不掉。这正是第 1.2 节里 log.segment.bytes 与 log.roll.ms 的意义:它们决定清理的时间精度。

6.2 压缩的墓碑与清除点

compact 策略下,清理线程在段级别做重写:读旧段,保留每个 key 的最新值,写出新段。删除一个 key 靠墓碑消息(value 为 null 的记录),墓碑会在 delete.retention.ms(默认 24 小时)后随段清理被移除。

清除点(cleaner checkpoint / first dirty offset) 记录「哪些 offset 之前已经压缩干净」。压缩线程每次只处理「dirty 区间」的段,避免重复工作。若一个 key 的更新持续发生在活跃段,它永远不会被压缩——压缩只处理已滚动的段。

这里与段设计的耦合是:段越大,压缩的延迟越高(要等整段滚完才开始压缩)。对 compact 主题,log.segment.bytes 可以调小(如 100 MB)让压缩更及时。

7. 监控与排错

7.1 关键指标

  • LogFlushRateAndTimeMs:刷盘速率与耗时。异常升高说明 IO 压力大或刷盘策略过激进。
  • RequestHandlerAvgIdlePercent:IO 线程空闲率,低于 30% 说明请求处理饱和。
  • 段数量与平均段大小:段数量异常多(每个分区上万)说明段太小或清理失效。
  • 磁盘使用率与 IO util:iostat -x 看 %util 与 await,%util 长期接近 100% 说明盘到瓶颈。
  • 页缓存命中:无直接指标,但可通过「消费延迟与磁盘读速率是否背离」间接判断。

7.2 典型故障

故障一:消费延迟飙升但磁盘读不高。 多半是页缓存命中率高、数据在内存里,瓶颈其实在网络或消费者处理速度,而非磁盘。别盲目加盘。

故障二:Found invalid index file。 索引损坏,按 3.3 节删索引重建。

故障三:磁盘空间不释放。 检查是否有未滚动的活跃段(调小 log.roll.ms)、是否有 compact 主题的段一直没压缩(检查清除点是否推进)。

故障四:broker 启动极慢。 大量段需要重建索引。减少段数量(调大 log.segment.bytes)或增加 num.recovery.threads.per.data.dir 并行加载。

# 并行加载段文件,加快启动
num.recovery.threads.per.data.dir=8

8. 常见坑清单

  • 把 .log 当临时文件删掉「重建」,数据不可恢复。
  • JVM 堆配到 30 GB,页缓存被挤压,消费性能反而下降。
  • log.flush.interval.messages 设成 1,每次都 fsync,吞吐暴跌。
  • 依赖时间戳查询却让生产者写入乱序时间戳,.timeindex 退化。
  • compact 主题段设得过大,压缩迟迟不触发,key 的历史版本堆积。
  • 多消费者并发读冷数据,页缓存不足,随机读打满机械盘。
  • log.dirs 用 RAID 10,容量减半且写放大,而副本机制本就够用。
  • 挂载不加 noatime,读日志带来额外元数据写。
  • 段设得过小,文件句柄耗尽,报 Too many open files。

9. 总结

Kafka 存储层的三个文件各司其职:.log 顺序追加存数据,.index 用 8 字节条目把随机查找压缩成「二分 + 4 KB 顺序读」,.timeindex 支持按时间定位(前提是时间戳单调)。数据从磁盘到网卡的零拷贝路径绕开用户态,页缓存则让热数据几乎不碰盘——这两点决定了「堆要小、盘要多、fsync 要少」的调优方向。

理解这些机制后,很多看似玄学的现象都能解释:为什么加内存比加盘有效、为什么老客户端更慢、为什么保留期不生效、为什么按时间查很慢。若要继续看留存策略、日志压缩与分层存储的工程实践,可以参考 Kafka 存储内核与日志压缩 ;若关注段与索引对写入路径的整体影响,可看 Kafka 集群高可用 。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「kafka」更多文章

  1. 集群升级与滚动重启实践
  2. Kafka 应用测试策略:Testcontainers 与集成测试
  3. 压缩算法选型:lz4、zstd、snappy 与 gzip