Learn
Kafka/12-storage-compaction

存储与日志压实

"Kafka 为什么快"的答案不在网络层,而在存储层:顺序写、PageCache、零拷贝、稀疏索引。理解磁盘上实际发生了什么,也能帮你回答"数据存多久、占多大、怎么清理"这些容量规划问题。

1. 磁盘目录结构

每个"分区副本"对应磁盘上一个目录,目录里是一段段 segment 文件:

/var/kafka/data/orders-0/            # topic "orders" 的 0 号分区
  00000000000000000000.log           # segment: 消息本体
  00000000000000000000.index         # offset 索引
  00000000000000000000.timeindex     # 时间戳索引
  00000000000005368769.log           # 下一个 segment, 文件名 = 起始 offset
  00000000000005368769.index
  00000000000005368769.timeindex
  leader-epoch-checkpoint            # 第 11 章的 epoch 记录
  partition.metadata
  • 只有最后一个 segment(active segment)在写入,其余都是只读的。
  • 滚动条件:log.segment.bytes(默认 1GB)或 log.roll.ms/log.roll.hours(默认 7 天),先到先滚。
  • 删除过期数据 = 直接删除整个老 segment 文件——这就是 Kafka 清理数据近乎零成本的原因。

2. 稀疏索引:按 offset 找消息

.index 文件是稀疏的:不是每条消息一个条目,而是每写入 log.index.interval.bytes(默认 4KB)才记一条:

.index (offset -> 物理位置)          .log
  相对offset  文件位置                 [msg 5368769] pos 0
  0          0                       [msg 5368770] pos 210
  35         4102          ------>   ...
  70         8230                    [msg 5368804] pos 4102
                                     ...
查找 offset=5368790:
  1. 二分定位 segment 文件 (文件名即起始 offset)
  2. 在 .index 里二分找 <= 目标的最大条目 (35 -> 4102)
  3. 从 4102 顺序扫描 .log 找到精确位置

稀疏索引把索引体积压到极小(可全部驻留内存),代价是最后一小段顺序扫描——对顺序读写为主的日志系统是完美取舍。.timeindex 同理,支撑按时间戳定位(offsetsForTimes、--to-datetime 重置)。

3. 为什么快:顺序写 + PageCache + 零拷贝

3.1 顺序写

追加写让磁盘免于寻道。机械盘顺序写可达数百 MB/s,与随机写差 3 个数量级;SSD 上顺序写同样对写放大更友好。

3.2 PageCache 而非进程内缓存

Kafka 写入只到 PageCache 就返回(依赖副本机制而非单机 fsync 保证可靠),读取也优先命中 PageCache:

传统做法: 应用堆内缓存 -> GC 压力大, 数据在堆内外复制两份
Kafka:    直接用 OS PageCache
  写: producer -> socket -> broker 写 PageCache (OS 异步刷盘)
  读: 消费不落后时, 直接从 PageCache 返回, 磁盘毫无压力

这就是"给 Kafka 机器留一半以上内存给 OS"的原因——那些内存都是缓存。

3.3 零拷贝(sendfile)

消费路径上,传统读文件发网络要 4 次拷贝、4 次上下文切换。Kafka 用 sendfile 系统调用:

传统:  磁盘 -> 内核缓冲 -> 用户空间 -> socket 缓冲 -> 网卡   (4 次拷贝)
零拷贝: 磁盘 -> 内核缓冲 ----------------------> 网卡        (2 次, 且不经过用户态)

注意:开启 SSL 后数据必须进入用户态加解密,零拷贝失效——这是 SSL 集群吞吐下降的重要原因之一。

4. 保留策略(cleanup.policy=delete)

默认按时间或大小删除老 segment:

log.retention.hours=168        # 默认 7 天 (retention.ms 优先级更高)
log.retention.bytes=-1         # 每分区大小上限, 默认不限制
log.retention.check.interval.ms=300000   # 清理线程检查周期
  • 判定单位是 segment:整个 segment 里最新消息也过期了才会删(所以实际保留时间略长于配置)。
  • topic 级覆盖:retention.ms、retention.bytes(用 kafka-configs 动态改,见第 4 章)。
  • 容量估算:每分区磁盘 ≈ 写入速率 × retention 时间 × (1 + 副本数-1)/生产端压缩比,规划时按峰值算。

5. Log Compaction(cleanup.policy=compact)

压实策略保留每个 key 的最新值,而不是按时间删除。适合"数据是状态而非事件"的场景:

压实前:
  offset: 0        1        2        3        4        5
  key:    u1       u2       u1       u3       u2       u1
  value:  addr=A   addr=X   addr=B   addr=M   null     addr=C
 
压实后:
  offset: 3        5
  key:    u3       u1
  value:  addr=M   addr=C
  (u2 最新值是 null 墓碑 -> 整个 key 被删除)

规则与要点:

  • 消息必须有 key,无 key 消息发到 compact topic 会报错。
  • value=null 是墓碑(tombstone),表示删除该 key;墓碑本身在 delete.retention.ms(默认 24h)后清除。
  • 压实只处理"干净段"之外的部分,active segment 永不压实,所以近期数据总是完整的。
  • offset 不会重排:压实后 offset 不连续是正常现象,消费逻辑不能假设 offset 连续。
  • 触发时机由 min.cleanable.dirty.ratio(默认 0.5,脏数据过半才清)与 log.cleaner.threads 控制。

典型用途:

场景说明
__consumer_offsetsKafka 自己用 compact 存消费进度
数据库变更流(CDC)每个主键只需最新镜像
Kafka Streams 的 changelog状态存储的备份,恢复时只需最新状态
配置/字典分发新消费者从头读一遍即得全量最新配置
# 创建 compact topic
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
  --create --topic user-profiles --partitions 6 --replication-factor 3 \
  --config cleanup.policy=compact \
  --config min.cleanable.dirty.ratio=0.2
# 也可以组合: 既压实又按时间删
#   --config cleanup.policy=compact,delete --config retention.ms=2592000000
⚠️常见误区
  1. 把 compact 当去重:压实是后台异步的,消费者完全可能读到同一 key 的多个历史版本,业务端仍需按"后读覆盖先读"处理。
  2. retention.ms 设太短:消费者积压超过保留时间,数据被删,LAG 永远追不回来且静默丢数据。保留时间至少要覆盖"最长可能的消费中断 + 修复时间"。
  3. 误以为 flush 了才安全:Kafka 默认不主动 fsync(log.flush 参数别乱动),可靠性来自多副本。单副本 topic 掉电就可能丢数据,与 acks 无关。
💡快速估算某 topic 的磁盘占用

用 kafka-log-dirs 工具直接查:bin/kafka-log-dirs.sh --bootstrap-server localhost:9092 --describe --topic-list orders,输出每个分区副本的字节数,比去机器上 du 方便得多。

小结

  • 分区 = 一串 segment 文件;只写 active segment,过期删除以 segment 为单位
  • 稀疏 offset/时间索引 + 二分定位 + 短距顺序扫描,索引小到可常驻内存
  • 快的三大支柱:顺序写、PageCache(内存留给 OS)、sendfile 零拷贝(SSL 会失效)
  • delete 策略按时间/大小删段;compact 策略按 key 留最新值,null 为墓碑
  • compact 不等于实时去重;retention 必须覆盖最长消费中断时间
🎯练习
  1. 把某测试 topic 的 segment.bytes 调到 1MB,写入几 MB 数据后去 log.dirs 目录观察 segment 滚动与文件命名规律。
  2. 创建 compact topic,对同一 key 写入 5 个版本再写一个 null 墓碑,等待压实后从头消费验证只剩最新值(可调低 min.cleanable.dirty.ratio 与 segment.ms 加速触发)。
  3. 用 kafka-log-dirs 查看 __consumer_offsets 的大小,解释它为什么用 compact 策略而不是 delete。