一、核心概念
Log Compaction(日志压缩)是 Kafka 提供的一种基于 key 的消息保留策略。它的核心保证是:对于同一个 partition 内的相同 key,Kafka 至少保留最新的一条消息,而旧版本的消息会被清理掉。
这与传统的按时间或大小删除(delete 策略)是完全不同的思路:
cleanup.policy=delete:按时间或总大小删除整段旧数据,不关心内容cleanup.policy=compact:按 key 去重,只留每个 key 的最新值cleanup.policy=compact,delete:两者结合
二、为什么需要 Log Compaction
考虑一个场景:你要用 Kafka 记录用户的最新地址变更。
如果用普通的 delete 策略:
- 保留 7 天:7 天前注册的用户地址就丢了
- 永久保留:磁盘会无限增长
而实际上你只关心每个用户最新的地址,历史版本没用。Log Compaction 正是为这类"变更日志 / changelog"场景设计的:Topic 表现得像一张只保留最终状态的 KV 表,同时又保留了 Kafka 顺序日志的所有优点(有序、可重放、多消费者)。
三、工作原理
Kafka 的每个 partition 的 log 会被划分成若干个 segment 文件,压缩过程只作用在非 active segment(正在写入的那个 segment 不动)。
从结构上看,log 被分成两部分:
- Clean 部分:已经被压缩过的旧数据,每个 key 只剩最新一条
- Dirty 部分:新写入、还未被压缩的部分,可能存在重复 key

从上图可以看到几个关键点:
- Key A 有三次更新(offset 0, 2, 5),压缩后只保留最新的
A → 3(offset 5) - Key B 有两次更新(offset 1, 4),压缩后只保留
B → 2(offset 4) - Key C 的最新值是
null(tombstone,墓碑消息),表示"删除这个 key" - 消息的 offset 不会重排:压缩前 offset 5 是
A → 3,压缩后依然是 offset 5,只是中间有"空洞",消费者按顺序读取时会跳过
四、压缩过程细节
压缩由 broker 后台的 Log Cleaner 线程池完成,配置项为 log.cleaner.threads(默认 1)。它的处理流程大致是:
第一步:挑选要压缩的 log
遍历所有 cleanup.policy=compact 的 partition,计算每个 log 的 dirty ratio(脏数据比例 = dirty 部分字节数 / 总字节数),选择比例最高、且超过 min.cleanable.dirty.ratio(默认 0.5)的 log 优先处理。
第二步:构建 OffsetMap
Cleaner 从 dirty 部分从头扫描一遍,构建一个 key → 最新 offset 的哈希表(OffsetMap,通常是 skimpy hash map,只存 key 的 hash 而不是原始 key,以节省内存)。占用内存由 log.cleaner.dedupe.buffer.size 控制(默认 128MB)。
第三步:重写 segment
再次从头扫描 log(这次包括 clean 部分和 dirty 部分),对每条消息:
- 查 OffsetMap,如果这条消息的 offset 等于该 key 记录的最新 offset → 保留
- 否则 → 丢弃
保留的消息被写入新的 segment 文件,原 segment 被替换。
第四步:处理 tombstone
tombstone(value 为 null 的消息)在压缩时会被保留一段时间,让消费者有机会读到"这个 key 被删除了"的信号。这段时间由 delete.retention.ms 控制(默认 24 小时)。超过后 tombstone 本身也会在下次压缩时被清掉。
五、关键配置参数
| 参数 | 作用 | 默认值 |
|---|---|---|
cleanup.policy | delete / compact / compact,delete | delete |
min.cleanable.dirty.ratio | 触发压缩的脏比例阈值 | 0.5 |
min.compaction.lag.ms | 消息在 dirty 中至少停留多久才能被压缩 | 0 |
max.compaction.lag.ms | 消息在 dirty 中最多停留多久必须被压缩 | Long.MAX |
delete.retention.ms | tombstone 保留时长 | 86400000 (1 天) |
segment.ms / segment.bytes | segment 滚动条件(active segment 不参与压缩,所以这个值影响新数据"变冷"的速度) | 7 天 / 1GB |
log.cleaner.threads | broker 级:Cleaner 线程数 | 1 |
log.cleaner.dedupe.buffer.size | broker 级:OffsetMap 内存 | 128MB |
六、典型使用场景
1. Kafka 内部:__consumer_offsets topic
这是最经典的例子。每个消费者组消费到哪个 offset,用 (group, topic, partition) 作为 key 写入这个 topic。消费者不断提交 offset 就是不断产生新版本,但你只关心每个 (group, topic, partition) 最新的那个 offset 值。用 log compaction 完美契合这个模型。
2. 数据库变更捕获(CDC)
用 Debezium 之类的工具把 MySQL/PostgreSQL 的变更同步到 Kafka,以主键作为 key。下游即使从头重放,也能重建出源表的最新状态(而不是每一行的全部历史)。
3. Kafka Streams 的 KTable / State StoreKTable 语义上就是"某个 key 的最新值",其底层的 changelog topic 用的就是 log compaction。任务重启时从这个 topic 恢复状态。
4. 配置分发、字典数据、用户 profile
写入以业务实体 ID 为 key,Kafka 变成一个可订阅、可重放的 KV 存储。
七、需要注意的地方
必须有 key。 没有 key(key 为 null)的消息不会被压缩,会永远堆积在 clean 部分(除非同时启用 compact,delete)。
Active segment 不会被压缩。 最新写入的那个 segment 是"热"的,不参与压缩。如果 topic 写入很慢,重复 key 可能长时间无法压缩掉;可以通过调小 segment.ms 强制滚动。
offset 是不连续的。 消费者不能假设 offset 是紧密相邻的,poll() 得到的相邻两条消息 offset 可能差很大。这对绝大多数消费者是透明的,但如果你自己做 offset 计算要注意。
Compaction 不保证任意时刻只有一个版本。 "至少保留最新值"是承诺,"只保留最新值"不是。dirty 部分中同一个 key 可能存在多个版本,消费者从头消费时会读到这些历史值,直到追到 clean/dirty 边界之后才是唯一版本。
Tombstone 有过期。 如果消费者停机超过 delete.retention.ms,可能错过删除信号,恢复后看到的 KTable 里会缺失那个 key 的"被删除"事件——它就直接消失了。
压缩不是免费的。 Cleaner 线程要读写磁盘、维护 OffsetMap 内存。key 空间非常大(比如 UUID + 高吞吐)时 OffsetMap 可能不够用,导致单轮压缩不能覆盖整个 dirty 部分,需要多轮才能收敛。
一句话总结:Log Compaction 让 Kafka topic 从"事件流"变成一张"最终状态表 + 变更日志"的合体,既能像 KV 存储那样查到每个 key 的最新值,又保留了顺序日志的全部好处。这是 Kafka 从消息队列演进到流处理平台的关键机制之一。
除非注明,否则均为李锋镝的博客原创文章,转载必须以链接形式标明本文链接

文章评论