李锋镝的博客

  • 首页
  • 时间轴
  • 说说
  • 每日心情
  • Now
  • 系列文章
  • 论坛
  • 左邻右舍
    • 左邻右舍
    • 博友圈
  • 留言
    • 留言
    • 走心评论
  • 关于
    • 关于我
    • 网站地图
    • 网站统计
    • 另一个网站
    • 我的导航站
    • 赞助
  • 🚇开往
Destiny
自是人生长恨水长东
  1. 首页
  2. 后端
  3. 正文

Kafka Log Compaction(日志压缩)详解

2026年8月26日 约 1,713 字6 分钟 6 0 0

一、核心概念

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.policydelete / compact / compact,deletedelete
min.cleanable.dirty.ratio触发压缩的脏比例阈值0.5
min.compaction.lag.ms消息在 dirty 中至少停留多久才能被压缩0
max.compaction.lag.ms消息在 dirty 中最多停留多久必须被压缩Long.MAX
delete.retention.mstombstone 保留时长86400000 (1 天)
segment.ms / segment.bytessegment 滚动条件(active segment 不参与压缩,所以这个值影响新数据"变冷"的速度)7 天 / 1GB
log.cleaner.threadsbroker 级:Cleaner 线程数1
log.cleaner.dedupe.buffer.sizebroker 级: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 Store
KTable 语义上就是"某个 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 从消息队列演进到流处理平台的关键机制之一。

除非注明,否则均为李锋镝的博客原创文章,转载必须以链接形式标明本文链接

本文链接:https://www.lifengdi.com/hou-duan/4928

本作品采用 知识共享署名-非商业性使用-相同方式共享 4.0 国际许可协议 进行许可
标签: Kafka Log Compaction MQ 日志
最后更新:2026年8月26日
相关文章
  • Aeron详解2026年8月13日
  • RocketMQ入门级教程2022年8月16日
  • MQ消费端遇到瓶颈该怎么办?2021年3月19日
  • Kafka 为什么要抛弃 Zookeeper?2025年10月11日
  • 关闭apache httpclient4.5 DEBUG日志2019年7月9日

李锋镝

既然选择了远方,便只顾风雨兼程。

打赏 点赞
< 上一篇
1234567891112131415161718192021222324252627282930313233343536373839404142434446474849505152535455575859606162636465666769727476777879808182858687909293949596979899
取消回复

文章评论

还没有评论,快来抢沙发吧~

晚上,我在宿舍和几个兄弟在聊天,忽然一哥们冲进屋内。他大喊:“我突然有一种热烈的想学习的冲动!”我说:“想学就学啊,没人拦你!”这哥们走到桌子前坐下。端起了一杯水说:“不行,这太冲动了,喝口水冷静一下。”我:……

听点儿音乐吧 朋友~
文章目录
最新 热点 随机
最新 热点 随机
Kafka Log Compaction(日志压缩)详解 WordPress缓存插件WP Fastest Cache、WP Rocket 、FlyingPress对比 每日早报 · 2026年8月20日 · 早上好 Kratos+ v1.1.18版本更新说明 每日早报 · 2026年8月19日 每日早报 · 2026年8月18日
给主题增加了Now、每日心情、年度回顾、岁月同一天、随机漫步等功能Kratos+ v1.1.16版本更新说明AI时代,个人技术博客的出路在哪里?增加了两套复古皮肤-牛皮纸、千禧网页写了一个订阅每日新闻的WP插件Kratos+ v1.1.18版本更新说明
深入剖析 ZGC 和 G1 垃圾回收器的区别 SpringBoot框架自动配置之spring.factories和AutoConfiguration.imports 架构师究竟比高级开发厉害在哪? 纪念中国人民抗日战争胜利76周年 Spring WebFlux底层原理深度剖析-从响应式流到事件循环的全链路拆解 为什么 Spring 不建议使用 @Autowired?@Resource 才是王道
最近评论
李锋镝 发布于 19 小时前(08月25日) 基本上大家都是这样,闲着没事就捣鼓捣鼓
不凡 发布于 19 小时前(08月25日) 嗐,闲时玩玩,没那么“极” :8:
李锋镝 发布于 19 小时前(08月25日) 哈哈哈~说明你也是极客
不凡 发布于 20 小时前(08月25日) 我最早用Rocket Cache,对个别主题不支持,然后换成Fastest Cache,设置虽然不多...
李锋镝 发布于 2 天前(08月24日) 都是AI的功劳~
标签聚合
JAVA ElasticSearch Claude SpringBoot 数据库 日常 MySQL MQ JVM Spring 多线程 WordPress 架构 SQL Redis K8s AI IDEA AI编程 分布式
友情链接
  • 彬红茶日记
  • 若梦博客
  • 老张博客
  • Serendipity
  • Mr.Sun的博客
  • sssr7844的博客
  • 志文工作室
  • 皮皮社
  • 韩情脉脉
  • 懋和道人
  • 搬砖日记
  • 知向前端
  • 蜗牛工作室
  • 瓦匠个人小站
  • 林羽凡
  • 哥斯拉
  • 临窗旋墨
  • 韩小韩博客
  • Honesty
  • 九仞之行

COPYRIGHT © 2026 lifengdi.com. ALL RIGHTS RESERVED.

正在博友圈履约中

域名年龄

Theme Kratos-plus By Dylan Li

津ICP备2024022503号-3

京公网安备11011502039375号