Skip to content

Kafka 日志保留与清理:retention 删除策略与 log compaction 压缩策略

提出问题

Kafka 的 Topic 消息能存多久?一天?一周?还是"随便设个 retention 就不用管了"?

线上 Kafka 的磁盘报警总在一夜之间蹦出来——race condition 删日志、compaction 卡住不干活、消费端 reset offset 跑到了已被清理的数据上。面试官问"Kafka 消息什么时候真正被删掉",最常见的回答是"时间到了就删",但实际是Segment 整体过期才能删,不是单条消息。

更隐蔽的是 log compaction:你以为它"只保留每个 key 最新一条"很省空间,但 active segment 不压缩、delete 标记不会自动消失、压缩线程跑不过消息写入速度时,磁盘只会越涨越多。

本文讲清楚两件事:retention 删除log compaction 压缩——它们怎么工作、什么时候触发、踩过哪些坑。

分析问题

Segment:Kafka 的最小清理单位

Kafka 的消息存在每个 Partition 下的日志目录里,不是一条一条文件,而是**分段(Segment)**存的。

/tmp/kafka-logs/mytopic-0/
├── 00000000000000000000.log        # Segment 1 数据文件
├── 00000000000000000000.index      # Segment 1 偏移量索引
├── 00000000000000000000.timeindex  # Segment 1 时间戳索引
├── 00000000000000000343.log        # Segment 2
├── 00000000000000000343.index
├── 00000000000000000343.timeindex
└── ...

每个 Segment 文件名的数字是第一条消息的 offset。Segment 滚动条件:

  • 大小达到 log.segment.bytes(默认 1GB)
  • 时间超过 log.roll.hours(默认 7 天)
  • 或者主动调用 flush() 或日志切换

关键理解:清理操作的单位是整个 Segment,不是单条消息。

Retention 删除策略:按时间或按大小

删除策略由 log.cleanup.policy=delete 控制(默认就是 delete)。

时间维度

properties
# 保留最近 7 天的数据,默认 168 小时
log.retention.hours=168

Kafka 会检查每个 Segment 的最近修改时间(mtime)或文件里的最大时间戳。当 Segment 的所有消息的最后一条消息的时间超过 retention 阈值,整个 Segment 被标记为 delete。

痛点是:过期不等于马上删除。后台有一个 LogCleaner 线程池,定期扫描所有 Segment,把过期的标记为"可删除"(文件名后缀加 .delete),然后真正 unlink 掉。这个扫描间隔由 log.retention.check.interval.ms(默认 5 分钟)控制。所以消息可能多存活 5 分钟

空间维度

properties
# 每个 Partition 最大保留 500GB
log.retention.bytes=536870912000

这个值是按 Partition 算的。如果 Topic 有 20 个 Partition,总数据量最多 20 × 500GB = 10TB。Retention 时间和大小可以同时设置,任一条件触达就删除

删除时机陷阱

"Segment 整体过期" 而不是"该 Segment 里的消息全部过期,但还没到 Segment 滚动时间"——如果某个 Segment 里第一条消息的写入时间已经超过 7 天,但最后一条消息才过了 6 天,那么整个 Segment 不会被删,因为这个 Segment 的 mtime 来自最后一条消息。

            Segment 1(200MB,跨越 2 天)
┌─────────────────────────────────────────────────────┐
│ 消息 1 (7 天前) 消息 2 ... 消息 N (6 天前)           │
└─────────────────────────────────────────────────────┘
               ↑ 因为这个,Segment 不会被删

解决方案:

  • 减小 log.segment.bytes(默认 1GB 偏大,密集型 Topic 可以降到 256MB),让 Segment 滚动更快
  • 或者用 log.roll.ms 强制按时间滚动,比如每小时滚动一次

Log Compaction 压缩策略:保留 key 的最新状态

cleanup.policy=compact 的场景完全不同——它不是为了"删旧消息",而是为了保留每个 key 的最新一次值

典型场景:Kafka 作为 changelog状态存储。比如有张用户表,每次更新都发一条消息到 Kafka,key 是用户 ID,value 是用户快照。你希望消费端能拿到每个用户的最新状态,但又不关心历史变更。Compaction 帮你去掉重复 key 的旧版本。

压缩过程

mermaid
graph LR
    subgraph "压缩前"
        A["k1 v1"] --> B["k2 v1"] --> C["k1 v2"] --> D["k3 v1"] --> E["k1 v3"] --> F["k2 v2"]
    end
    subgraph "压缩后"
        G["k1 v3"] --> H["k2 v2"] --> I["k3 v1"]
    end
    subgraph "压缩后"
        D --> E --> F
    end

具体流程:

  1. LogCleaner 线程选出"脏数据比"最高的 Partition(dirty ratio = 未压缩数据 / 总数据)
  2. 读取完整 offset 范围的 key 构建一个内存哈希表(key → 最新 offset)
  3. 重写 Segment,只保留每个 key 的最新 offset 对应的消息
  4. 被清理掉的旧消息所在的 Segment 文件被删除

坑点 1:Active Segment 不参与压缩

正在写入的那个 Segment(active segment)永远不会被 compaction 处理。这是故意的——直接压缩正在写的文件会导致并发问题。

后果:最新写入的消息要等当前 Segment 滚动后才可能被压缩。如果 1GB 的 Segment 7 天才能被写满,那这 7 天内的"旧 key 重复消息"就一直占着空间。

坑点 2:Delete 标记不会消失

java
// 消费者发送一个 tombstone 消息——key 为 "user_123",value 为 null
producer.send(new ProducerRecord<>("user-changelog", "user_123", null));

Consumer 收到 tombstone 后知道"这个 key 被删了",但在 compaction 完成之前,这条 tombstone 和旧数据都还在。Compaction 之后,key 相关的所有消息(包括旧数据和 tombstone)都会被清理掉。但 tombstone 本身还要保留一段时间:log.cleaner.delete.retention.ms(默认 24 小时),防止还没消费完的 Consumer 读到被删的数据。

坑点 3:Compaction 跑不过写入速度

一个 Topic 的写入速率是 50MB/s,而 compaction 的清理速率只有 20MB/s。脏数据比从 50% 涨到 80%,再涨到 95%。这时磁盘会持续增长,直到触发 retention 的 log.retention.bytes 兜底——但 retention 会无差别删除最旧的 Segment,不管有没有被 compaction 处理过,导致数据丢失。

监控指标kafka.log:type=LogCleaner,name=MaxDirtyPercent 如果持续 > 80%,说明压缩跟不上写入,需要:

  • 增加 log.cleaner.threads(默认 1)
  • 减小 log.segment.bytes,加速滚动
  • 增加 log.cleaner.io.max.bytes.per.second(默认无限,但实际受磁盘 IO 限制)
  • 考虑换更大的磁盘

删除 vs 压缩:同一个 Topic 可以同时用

properties
# 同时启用:先按 retention 删过期 Segment,再对剩余 Segment 做 compaction
log.cleanup.policy=[delete, compact]

生产环境常见的做法:Kafka 版本的"changelog"Topic设置 compact 保留状态;"事件日志"Topic设置 delete 按时间清掉。也有混合场景——比如保留最近 7 天的所有数据(delete 策略),同时对 7 天内的数据做 compaction 去重,减少存储开销。

磁盘水位监控与清理压力

清理压力来自分区数太多。每个 Partition 都有一组 Segment,LogCleaner 线程要扫描所有 Partition 的 Segment 元数据。如果单台 Broker 上有 2000 个 Partition,每个 Partition 有 10 个 Segment,就是 20000 个文件。文件句柄压力 + 扫描开销导致清理周期延长。

磁盘水位预警:在磁盘使用率到 75% 时就要开始排查清理进度,到 85% 是紧急状态。Kafka 的硬限制是 log.dirs 所在磁盘写满时 Broker 会直接挂掉。

bash
# 检查每个 Partition 的日志大小和清理状态
# 看每个 Partition 的 offset 起止范围
./kafka-log-dirs.sh --bootstrap-server localhost:9092 --describe --topic-list my-topic

# 检查 LogCleaner 是否在正常工作
# 如果 PausedCleanerCount 长期 > 0,说明压缩被暂停了
./kafka-run-class.sh kafka.admin.LogCleanerManager --bootstrap-server localhost:9092

总结

维度Retention 删除Log Compaction
目的释放磁盘空间保留每个 key 最新一条
单位整个 SegmentSegment 内按 key 去重
触发条件时间/大小超过阈值脏数据比 > dirty ratio
适用场景日志、事件流状态表、changelog
共享配置log.retention.check.interval.mslog.cleaner.threads
核心坑Segment 整体过期才删Active segment 不压缩、tombstone 保留期

Kafka 的日志清理看起来简单,实际涉及Segment 滚动机制、后台线程扫描、磁盘 IO 预算三层约束。面试官问"Kafka 消息能存多久"时,别只答"设 retention 就行",说出 Segment 整体过期、清理线程间隔、compaction 的 active segment 不压缩,才是真理解。

参考

手撕 → 框架 → 生产化,一步步把 AI Agent 工程化搞透。
粤ICP备2026104257号-1