日志清理与压缩
定义与作用
Kafka 的数据不是无限增长的。日志清理(Log Cleanup)机制负责删除过期数据,释放磁盘空间。Kafka 提供两种清理策略:基于时间的删除(Delete) 和 基于 Key 的压缩(Compact)。
- Delete:适合"流式数据"——超过保留时间后整条丢弃
- Compact:适合"状态数据"——只保留每个 Key 的最新值
核心原理
Delete 策略
删除判断逻辑:
- 检查 Segment 文件中最后一条消息的时间戳
- 如果
now - maxTimestamp > retention.ms→ 删除整个 Segment - 活跃 Segment(当前正在写入的)永远不会被删除
Compact 策略
Compact 的特点:
- 保留每个 Key 的最新值,删除旧值
- 不影响 Consumer 的 Offset 消费(Consumer 按 Offset 读取,不会被压缩影响)
- 适合保存"当前快照"(如用户配置表、最新交易状态)
清理线程架构
Log Cleaner 线程池:
- 每个 log.dir 有 1 个清理线程
- 周期扫描每个 Partition
- Delete: 检查 Segment 时间戳 → 删除过期 Segment
- Compact:
1. 选择"最脏"的 log(dirty ratio = 可清理消息数 / 总消息数)
2. 清理旧 Key → 合并为新 Segment
3. 替换旧 Segment
完整示例
示例一:配置基于时间的删除
# 创建保留 1 小时的数据流 Topic
bin/kafka-topics.sh --create --topic raw-events \
--partitions 6 --replication-factor 3 \
--config retention.ms=3600000 \
--config segment.ms=600000 \
--bootstrap-server localhost:9092
参数说明:
| 参数 | 值 | 含义 |
|---|---|---|
retention.ms | 3600000 (1小时) | 消息保留 1 小时 |
segment.ms | 600000 (10分钟) | 每 10 分钟切新 Segment |
效果:每个 Segment 10 分钟。1 小时后,最早的那个 Segment 被删除。
示例二:配置 Log Compact 存储最新状态
# 创建 Compact Topic(保存用户最新登录信息)
bin/kafka-topics.sh --create --topic user-sessions \
--partitions 3 --replication-factor 3 \
--config cleanup.policy=compact \
--config min.cleanable.dirty.ratio=0.5 \
--config segment.ms=3600000 \
--bootstrap-server localhost:9092
# 发送数据
bin/kafka-console-producer.sh --topic user-sessions \
--property parse.key=true --property key.separator=: \
--bootstrap-server localhost:9092
> user1:login-2026-06-01
> user2:login-2026-06-01
> user1:login-2026-06-13
操作前后对比:
| Key | 初始状态(3 条消息) | Compact 后 | Consumer 可见 |
|---|---|---|---|
| user1 | v1, v3 | v3(保留最新) | v1(如果 Offset 在那条), v3 |
| user2 | v2 | v2(唯一值) | v2 |
注意:Consumer 仍可读到被压缩的旧值(如果从那条 Offset 消费),Compact 只释放磁盘空间,不改变 Consumer 行为。
易错场景
易错 1:retention.bytes 设置为 -1(无限制)
场景:生产环境只配置了 retention.ms,没有配置 retention.bytes。
后果:如果写入速率极高,7 天的数据可能写满整个磁盘。
正确做法:始终设置 retention.bytes 作为磁盘保护:
retention.bytes=107374182400 # 100GB,即使时间未到也触发删除
易错 2:Compact Topic 不能保证"立即"回收
场景:期望 Compact Topic 中的数据实时反映最新值。
事实:Compact 是异步后台进程,不会在写入时立即压缩。消息可能在压缩前被 Consumer 消费到旧值。
规则:Compact 适用于最终一致性场景(如配置同步),不适合强一致性场景(如实时库存)。
面试高频考点
Q:Compact Topic 中,如果同一个 Key 写入 100 次,Consumer 从最早开始消费会看到多少条?
A:100 条。Compact 不影响 Consumer 的读取——Consumer 按 Offset 遍历所有消息。Compact 只删除标记为"已清理"的旧消息的磁盘副本,但这些消息对应的 Offset 仍然存在(以"墓碑"或直接跳过的方式)。
如果 Consumer 需要只看到每个 Key 的最新值,需要业务层去重或使用 KTable(Kafka Streams 提供)。