Apache Kafka Compact Topic
Apache Kafka Compact Topic 入门与实战指南
Apache Kafka 的 Compact Topic(日志压缩主题)是一种特殊的主题类型,它不按时间或大小删除旧消息,而是通过保留每个消息键(Key)的最新值来清理历史数据。本教程将从零开始解释 Compact Topic 的原理、配置、典型场景与注意事项,帮助你快速掌握这一核心机制。
什么是 Compact Topic
Kafka 主题的删除策略
Kafka 普通主题默认采用基于时间或大小的删除清理策略(delete)。当消息保留超时(retention.ms)或分区日志总大小超过上限(retention.bytes)时,旧分段文件会被直接删除,无论消息的键是什么。
Compact Topic 的定义
Compact Topic 使用压缩清理策略(compact)。压缩过程会扫描分区日志,对于同一个键的所有消息,只保留最新的一条(即最大偏移量的那条)。如果某条消息的值是 null(称为墓碑消息),且后续没有该键的新消息,那么在足够长的清理周期后,该键的所有记录都会被移除。
简单理解:Compact Topic 像一个分布式的键值数据库的 changelog,永远只保留每个键的最新状态。
日志压缩的工作原理
传统日志与压缩日志对比
- 普通日志:按顺序保留所有写入的消息,适合时间驱动的处理(如监控、点击流)。
- 压缩日志:后台线程定期将分区的分段文件重写,移除旧值的键,只保留最新值或墓碑消息。
清理过程(Log Compaction)
压缩不是实时进行,而是由 Log Cleaner 线程池周期性触发。具体步骤:
- 分段选择:Log Cleaner 选择脏数据率(dirty ratio)最高的分段文件进行处理。脏数据率指分段中可被清理的记录占比。
- 构建偏移量映射:扫描整个分区的最新偏移量,确定每个键的最后出现位置。
- 重新复制:将分段中那些在映射里对应的消息保留,其余丢弃,生成新的清洁分段。
- 原子替换:用新分段替换旧分段。
因为压缩保留的是键的最新值,消费者依然可以完整回溯到所有键的最终状态,而不会中途丢失数据。
配置 Compact Topic
关键配置参数
创建主题或修改主题配置时,以下参数控制压缩行为:
| 参数 | 说明 | 默认值 |
|---|---|---|
cleanup.policy |
设置清理策略,可组合 delete 和 compact |
delete |
min.cleanable.dirty.ratio |
分段文件中脏数据比例达到该值时,才会被清理 | 0.5 |
segment.ms |
日志分段文件的最大时间跨度,控制分段轮转速度 | 7 days |
delete.retention.ms |
墓碑消息被彻底删除前的保留时间 | 24 hours |
max.compaction.lag.ms |
消息从写入到可以被压缩的最大延迟(保证新消息不被过早清理) | 无限制 |
创建 Compaction 主题示例
使用 kafka-topics.sh 创建主题时指定压缩策略:
bin/kafka-topics.sh --create \
--bootstrap-server localhost:9092 \
--topic user-profiles \
--partitions 3 \
--replication-factor 2 \
--config cleanup.policy=compact \
--config delete.retention.ms=86400000
也可以为已有主题动态修改:
bin/kafka-configs.sh --bootstrap-server localhost:9092 \
--entity-type topics --entity-name user-profiles \
--alter --add-config cleanup.policy=compact
注意:修改策略后,只有新写入的数据会按
compact策略处理;旧数据不会被立即压缩,但后续清理周期会逐步处理。
使用场景
键值存储与状态恢复
在 Kafka Streams 或 ksqlDB 中,状态存储的 changelog 主题通常采用 compact 策略。应用程序重启时,可以从头读取 compact topic,快速重建本地内存中的最新键值状态,而无需重放全部变更历史。
CDC(变更数据捕获)事件流
Debezium 等 CDC 工具将数据库表的变更事件写入 Kafka。对于需要快照的场景,可以将 CDC 数据导入 compact topic,实现对每个主键当前行的最新快照。下游消费者直接读取紧凑日志即可获得最新全量数据。
保证最终一致性
在事件驱动架构中,Compact Topic 作为实体状态的唯一可信源。即使生产者发送多条更新,消费者只需关心每个键的最新值;结合墓碑消息,可以明确表示数据的删除,让下游系统实现最终一致。
注意事项与最佳实践
消息键的设计
- 必须有键:没有键的消息无法被压缩,它们永远不会被清理,最终可能导致磁盘无限增长。务必为所有发送到 compact topic 的消息指定一个逻辑上的唯一键。
- 键的基数:超高基数的键(如 UUID)使压缩效果变差,因为每个键几乎都只有一条消息,清理比率很低。需要在“保留更新历史”和“磁盘空间”之间权衡。
墓碑消息(Tombstone)
- 发送一个 value 为
null且键存在的记录,即标记该键的删除。 - 墓碑消息会在
delete.retention.ms后从日志中彻底移除,之后该键不再有任何记录。 - 消费者必须能正确处理
null值:例如在重建状态时,看到某个键的值为null就应该将其删除。
性能考虑
- 生产端:由于压缩过程需要维护键索引,broker 端会有额外的 CPU 和 I/O 消耗。如果每个消息的键变化频繁且值很大,压缩会对性能产生较大影响。
- 消费端:从头读取 compact topic 时,消费者会跳过中间已清理的消息,但 Kakfa 协议仍需要遍历偏移量,可能会有少量跳过开销。
- 避免组合策略陷阱:如果同时使用
delete和compact策略(cleanup.policy=compact,delete),则日志会先按时间/大小条件删除分段,再在剩余分段中做压缩。这种组合适合既要保留键最新值,又要按时间淘汰过时键的场景。
监控与调优
关注 broker 指标:
kafka.log:type=LogCleaner,name=max-clean-time-sec:最大清理耗时,过高说明日志压力大。kafka.log:type=LogCleanerManager,name=max-dirty-percent:当前最脏分段的比例,可用于判断压缩是否滞后。
合理调整 min.cleanable.dirty.ratio 可以平衡压缩频率与磁盘占用;增加 delete.retention.ms 可给下游消费者更长时间处理墓碑消息。
快速上手小结
- 创建 compact topic,指定
cleanup.policy=compact。 - 发送带有逻辑键的消息;删除数据时发送 value 为
null的墓碑消息。 - 消费者按偏移量顺序读取,处理每个键的最新值或
null。 - 根据键的基数、墓碑保留需求调整
delete.retention.ms等参数。
有了 Compact Topic,你可以在 Kafka 中自然地构建持久化的键值快照流,为状态管理、数据同步等场景提供简洁可靠的基础设施。