Apache Kafka Compact Topic

FreeGuideOnline 最新 2026-07-12

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 线程池周期性触发。具体步骤:

  1. 分段选择:Log Cleaner 选择脏数据率(dirty ratio)最高的分段文件进行处理。脏数据率指分段中可被清理的记录占比。
  2. 构建偏移量映射:扫描整个分区的最新偏移量,确定每个键的最后出现位置。
  3. 重新复制:将分段中那些在映射里对应的消息保留,其余丢弃,生成新的清洁分段。
  4. 原子替换:用新分段替换旧分段。

因为压缩保留的是键的最新值,消费者依然可以完整回溯到所有键的最终状态,而不会中途丢失数据。


配置 Compact Topic

关键配置参数

创建主题或修改主题配置时,以下参数控制压缩行为:

参数 说明 默认值
cleanup.policy 设置清理策略,可组合 deletecompact 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 协议仍需要遍历偏移量,可能会有少量跳过开销。
  • 避免组合策略陷阱:如果同时使用 deletecompact 策略(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 可给下游消费者更长时间处理墓碑消息。


快速上手小结

  1. 创建 compact topic,指定 cleanup.policy=compact
  2. 发送带有逻辑键的消息;删除数据时发送 value 为 null 的墓碑消息。
  3. 消费者按偏移量顺序读取,处理每个键的最新值或 null
  4. 根据键的基数、墓碑保留需求调整 delete.retention.ms 等参数。

有了 Compact Topic,你可以在 Kafka 中自然地构建持久化的键值快照流,为状态管理、数据同步等场景提供简洁可靠的基础设施。