主题
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=168Kafka 会检查每个 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具体流程:
LogCleaner线程选出"脏数据比"最高的 Partition(dirty ratio = 未压缩数据 / 总数据)- 读取完整 offset 范围的 key 构建一个内存哈希表(key → 最新 offset)
- 重写 Segment,只保留每个 key 的最新 offset 对应的消息
- 被清理掉的旧消息所在的 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 最新一条 |
| 单位 | 整个 Segment | Segment 内按 key 去重 |
| 触发条件 | 时间/大小超过阈值 | 脏数据比 > dirty ratio |
| 适用场景 | 日志、事件流 | 状态表、changelog |
| 共享配置 | log.retention.check.interval.ms | log.cleaner.threads |
| 核心坑 | Segment 整体过期才删 | Active segment 不压缩、tombstone 保留期 |
Kafka 的日志清理看起来简单,实际涉及Segment 滚动机制、后台线程扫描、磁盘 IO 预算三层约束。面试官问"Kafka 消息能存多久"时,别只答"设 retention 就行",说出 Segment 整体过期、清理线程间隔、compaction 的 active segment 不压缩,才是真理解。
参考
- Kafka 高性能核心设计:零拷贝、顺序写、页缓存、批处理
- Kafka 核心概念:Topic、Partition 与 Consumer Group
- Kafka 生产者与消费者原理
- Apache Kafka Documentation: Log Compaction