Skip to content

第 14 章 Log Compaction 与 Tombstone:Kafka 的「最终态快照」

目标读者:知道 Kafka 是「日志型」存储,但不理解为什么 __consumer_offsets 不会无限增长、不知道怎么用 Compaction 做 KV 快照、看到「Tombstone」一脸茫然的同学。

学完你会:能说出 delete / compact / compact,delete 三种保留策略的工程取舍;理解 Compaction 怎么从「无限追加日志」中变出「最终态快照」;会用 null value 删除 Key;能解释 KStream / KTable 的 changelog 为什么必须配 compact


0. 导读:从「日志」变出「KV 快照」的魔法

Kafka 的核心抽象是「追加日志」(Append-Only Log):消息只能往末尾追加,永远不会原地修改。这给了 Kafka 极致的写入性能(顺序写 + zero-copy),但也带来一个问题:

如果同一个 Key 被反复更新(比如商品库存),日志会无限增长,存储吃不消,下游回放也得过滤掉一堆历史无效值。

Log Compaction 解决的就是这个问题:

同一个 Key 只保留最新的 Value,旧的物理删除。

打个生活化比方:

  • 普通 Topic(cleanup.policy=delete) 像「记账流水本」——每笔交易都顺序写一行,过期日期到了整本撕掉。适合事件流(订单事件、点击日志)。
  • Compact Topic(cleanup.policy=compact) 像「配方册」——同一道菜的最新配方覆盖旧配方,只保留每道菜最新一版。适合状态快照(用户最新地址、商品最新价格、消费组最新 offset)。

第二种是 KV Store 的「日志化版本」,是 Kafka 跨界做「分布式状态存储」的关键。__consumer_offsets__transaction_state__cluster_metadata、KStream 的 changelog topic 全部都用 Compaction。


1. 三种保留策略(cleanup.policy)

properties
# Topic 级配置
cleanup.policy=delete         # 默认;按时间/大小删除整段 Segment
cleanup.policy=compact        # 同 Key 只保留最新 Value
cleanup.policy=compact,delete # 同 Key 只保留最新 + 整段也按时间删除

1.1 delete(默认)

经典「Topic」语义:消息按 retention.ms(默认 7 天)/ retention.bytes(默认无限)保留,过期就整段 Segment删除。

log.dirs/orders-0/
├── 00000000000000000000.log   ← 7 天前的消息,过期被删
├── 00000000000123456789.log   ← active segment(持续追加中)
├── ...

特点:

  • 顺序消费:offset 永远递增,消费者从某个 offset 拉到末尾就行。
  • 删除粒度:整段 Segment 一起删(不会单独删某条消息),所以回放历史有 7 天窗口。
  • 适合:事件流、点击日志、CDC binlog 透传。

1.2 compact

「同 Key 最新值」语义:Compaction 后台线程定期扫描分区日志,同一个 Key 的旧 Value 物理删除,只保留最新一条。

原始日志:
offset  key  value
   100  k1   v1_old
   101  k2   v2_old
   102  k1   v1_mid
   103  k3   v3
   104  k1   v1_new
   105  k2   v2_new
   106  null v6  ← 没 key 的消息,compact 不动它(一直保留?看版本)
                  实际上 compact 会跳过 null-key

Compaction 后:
offset  key  value
   103  k3   v3
   104  k1   v1_new
   105  k2   v2_new

注意:

  • offset 是不连续的(100 → 103 → 104 → 105,其中 100/101/102 被删了)。
  • 保留的 offset 是「该 Key 最新一条」的原 offset——不会重新编号。
  • Active Segment(正在写入的最末尾 Segment)不参与压缩——避免热段抖动。

特点:

  • 保证「无限期保留每个 Key 的最新 Value」——不会因为时间删掉。
  • 回放时:消费者从头扫到尾就能拿到「最终态快照」,相当于一次 KV dump。
  • 适合:用户资料、商品价格、__consumer_offsets、Streams changelog、CDC 表最终态。

1.3 compact,delete

两种策略叠加:先按 compact 压缩留最新 Value,再按 retention.ms 老化删除整段。

适合:「想保留每个 Key 的最新值,但 Key 本身也有时效」。例如:

  • 用户最近 30 天的最新行为(30 天后整体过期)
  • 实时风控特征:每个用户最新的得分,但保留 7 天

2. Compaction 的「最终一致性」语义

注意:Compaction 不是实时去重,它是最终一致的:

  • 你刚写完 (k1, v_new),再写 (k1, v_newer),磁盘上两条都在
  • Compaction 线程周期性扫描,满足触发条件才合并,可能要等几分钟、几小时。
  • 在压缩前,新启动的消费者从头读会看到 v_newv_newer 两次。

所以你的应用必须能接受同一个 Key 被多次看到——只要保证「最后一次看到的是最新的」就行。换句话说:

应用要按「重放每条 (key, value) 都做 PUT」的方式编写,不要做「INSERT 失败就报错」

KStream / Flink 的 KTable 就是这种「重放等价于 PUT」的语义。


3. Tombstone:用 null value 实现「删除」

3.1 什么是 Tombstone?

发一条 (key=k1, value=null) 的消息——这条消息的语义是「删除 k1」。它叫 Tombstone(墓碑标记)。

写入时间线:

1. 应用 produce(k1, "alice")        → 插入
2. 应用 produce(k1, "bob")          → 更新
3. 应用 produce(k1, null)           → 删除标记 (tombstone)
4. (Compaction 触发)                → k1 的所有旧值删除,只保留 tombstone
5. (delete.retention.ms 过期)       → tombstone 也物理删除

3.2 为什么不直接物理删?

如果只删 k1 的所有记录、不写任何标记:

  • 已经在消费 / 即将启动消费的 read_committed 消费者读不到「删除」事件——它只会看到 k1 突然消失,无法区分「k1 没存在过」「k1 被删了」。
  • KStream 的 KTable 状态里 k1 还留着旧值,因为它没收到删除信号

所以 Compaction 设计成:先保留 tombstone,让所有消费者都有时间看到「这条 Key 被删除了」,然后才物理清理。

3.3 delete.retention.ms:tombstone 保留多久

properties
delete.retention.ms=86400000   # 默认 24h

意思是:tombstone 至少保留 24h,让下游消费者有时间看到。超过 24h 后,下次 Compaction 会把 tombstone 也物理清掉。

📌 生产经验:如果消费者下线超过 delete.retention.ms 又重启,可能永久错过某些 Key 的删除事件。要么调大这个值,要么消费者自己定期对账。

3.4 Compaction 的物理布局

Segment 文件示意(简化):
.../my-topic-0/
├── 00000000000000000000.log     ← 已压缩段
├── 00000000000000000000.index
├── 00000000000000004500.log     ← 已压缩段(offset 不连续)
├── 00000000000000004500.index
├── 00000000000000010000.log     ← active segment,未参与压缩
└── 00000000000000010000.index

每条日志条目:

RecordBatch:
  offset    : long
  timestamp : long
  key       : bytes
  value     : bytes  (null = tombstone)
  headers   : ...

Compaction 后:旧 (k, v_old) 整条记录消失,offset 跳号;新值 (k, v_new) 保留原 offset。


4. 触发条件:Compaction 什么时候跑

Cleaner 线程不是「时刻在跑」,而是有触发阈值。

4.1 关键参数

参数默认值含义
min.cleanable.dirty.ratio0.5dirty 部分占整个 partition 的比例 ≥ 该值才触发
min.compaction.lag.ms0消息距今至少 N ms 才能被压缩(不动太新的)
max.compaction.lag.ms长整数最大消息距今最多多久必须被压缩(再忙也得轮到)
segment.ms7 天active segment 多久滚动一次(滚动后才能参与压缩)
segment.bytes1 GB同上,按大小滚动
delete.retention.ms24 htombstone 至少保留多久
log.cleaner.threads1Broker 全局 cleaner 线程数(默认 1,I/O 重的可以加)
log.cleaner.dedupe.buffer.size128 MB每个线程构建「key → 最新 offset」哈希表的内存

4.2 「dirty ratio」直觉

分区日志:
[已压缩区  clean]  [可压缩区  dirty]  [active segment 不动]
       ^                  ^
       |                  |
       压缩起点(cleanerCheckpoint)   active 起点(activeBaseOffset)

dirty ratio = dirty 大小 / (clean + dirty)

dirty ratio < 0.5:cleaner 不动它(攒一会儿再压)。 dirty ratio ≥ 0.5:触发一次压缩任务,把 dirty 区合并进 clean 区。

📌 调优:流量稳定的小 topic(比如 __consumer_offsets)默认 0.5 就好;写入超大的 KStream changelog 可以调到 0.3 让压缩更勤;大对象 + 写入慢可以调到 0.7 减少 cleaner 抢 IO。

4.3 min.compaction.lag.ms

防止「刚写就被压缩掉」——某些场景应用需要时间消费到原始消息(比如 audit log):

properties
min.compaction.lag.ms=600000   # 至少 10 分钟内的消息不压

4.4 max.compaction.lag.ms(重要!)

再忙也得有个上限」——保证 GDPR 删除请求等场景在合理时间内生效:

properties
max.compaction.lag.ms=86400000   # 24h 内必须压缩一轮

没有这个的话,dirty ratio 一直不到 0.5 的小 topic 可能永远不压缩,tombstone 永远存活,违反「右forgotten」要求。


5. Cleaner 线程模型

Broker 启动时拉起一组 log.cleaner.threads 个 Cleaner 线程,循环做:

loop forever:
    1) 选一个 dirty ratio 最高的 (topic, partition)
    2) 扫描 dirty 区,构建「key → 最新 offset」OffsetMap(哈希表)
    3) 重新读 clean + dirty,每条 record:
        - 如果 record.offset == OffsetMap[key] → 保留
        - 否则 → 跳过(被新值覆盖)
       写入新的 segment 文件
    4) 用新 segment 替换旧 segment
    5) 推进 cleanerCheckpoint
    6) 写 .cleaned / .swap 后缀文件,原子切换

线程数选型

  • 默认 1,对小集群足够。
  • 大集群(每 Broker 1000+ 分区)可以加到 4~8。
  • 不要无脑调大——cleaner 抢 disk IO,会拖慢实时读写。

观察 cleaner 工作:

bash
tail -f /var/log/kafka/log-cleaner.log

输出示例:

[2026-04-17 10:23:45] Beginning cleaning of log __consumer_offsets-23
[2026-04-17 10:23:45] Building offset map for log __consumer_offsets-23 for 1 segments in offset range [4096, 8192)
[2026-04-17 10:23:46] Offset map for log __consumer_offsets-23 complete.
[2026-04-17 10:23:46] Cleaning segment 4096 in log __consumer_offsets-23 (largest timestamp ...)
[2026-04-17 10:23:47] Swapping in cleaned segment ...
[2026-04-17 10:23:47] Cleaner: __consumer_offsets-23 retained 423/1024 records (size 4MB → 2MB)

6. 应用场景详解

6.1 状态快照:KStream 的 KTable changelog

KTable 是「键值表」抽象:

scala
KTable<UserId, Profile> profiles = builder.table("user_profiles");

每个 KTable 在底层有一个 changelog topic(默认命名 <application_id>-<store_name>-changelog),用 cleanup.policy=compact

  • 每次本地状态更新(RocksDB),先写一条 (key, new_value) 到 changelog
  • 任务故障转移到新 Broker → 新机器从 changelog 重放,恢复 RocksDB 状态
  • 因为是 compact topic,重放只看最新值,不会回到中间态

6.2 CDC 数据:保留主键最终态

Debezium / Flink CDC 把 MySQL 变更同步到 Kafka,常用 compact topic 保存「每行的最终态」:

源 MySQL:
  user_id=1: name=alice (insert)
  user_id=1: name=alice2 (update)
  user_id=1: (delete)

Kafka topic (compact):
  → 写: (1, {name:alice})
  → 写: (1, {name:alice2})
  → 写: (1, null)            ← tombstone
  → Compaction 后:
        所有 user_id=1 的旧记录消失,仅留 tombstone(且 24h 后也清掉)

下游消费者重启时只看到 tombstone(或啥也没有),知道「user 1 已删除」。

6.3 __consumer_offsets

每个消费组的每个 (group, topic, partition) → offset,按 (group, topic, partition) 当 Key,每次 commit 一条新值。Compact 后只保留最新 offset。

Key: (group=A, topic=T, partition=0)
旧值: offset=100  (10 分钟前)
新值: offset=200  (5 分钟前)
新值: offset=300  (1 分钟前)
→ Compaction 后只剩 offset=300

6.4 __transaction_state

每个 transactional.id 的状态变迁也是用 compact 存:每次状态变化追加一条,旧状态被新状态覆盖。

6.5 __cluster_metadata(KRaft 模式)

整个集群的元数据(topic、partition、ACL)以「事件溯源」形式追加到 __cluster_metadata,Compaction 保留每个元数据对象的最新版本。


7. 与 KV Store 的概念对应

KV 概念Kafka Compact Topic
PUT(key, value)produce(key, value)
DELETE(key)produce(key, null)
GET(key)从头消费整个 topic 取 key 最新值(笨)
全表扫描 / Snapshot从头消费整个 topic(自动得到「快照」)
变更订阅订阅 topic,每次新写都收到通知

📌 重大启示:Compact Topic 本质上是一个「带变更通知的 KV Store」。它不能像 RocksDB 那样按 Key 随机查,但提供了 KV 没有的「完整变更历史」。所以 Streams 把 RocksDB(用于查询)+ Compact Topic(用于持久化和恢复)配合使用。


8. 配套实战

  • init.sh:建一个 cleanup.policy=compact 的 topic
  • compaction_demo.py:写入大量同 Key 不同 Value,触发 compaction,前后对比
  • tombstone_delete.py:用 null value 删除 Key 的演示

9. 与其他 MQ 的对比

维度Kafka CompactionRabbitMQRocketMQPulsar
是否支持✅ 原生❌(消息只按时间过期)✅ Topic Compaction
用途状态快照、changelog、__consumer_offsets队列即时消费同 RabbitMQ类 Kafka
删除标记null value (tombstone)N/AN/Anull value
Key 的作用必须有 key 才能压缩routing key(仅路由用)message tag(仅过滤)key 用于 partition + compaction

📌 本质差异:Kafka 把「日志型存储」和「KV Store」用一个文件格式统一起来——这是 RabbitMQ / RocketMQ 都没做的设计。Pulsar 学了这一点。


10. 小结

  1. 三种 cleanup.policydelete(默认,按时间/大小删段)、compact(同 Key 最新值)、compact,delete(叠加)。
  2. Compaction = 把「无限追加日志」变成「最终态 KV 快照」,保留每个 Key 的最新 Value。
  3. Tombstone(null value)= 删除标记,先保留 delete.retention.ms(默认 24h)让消费者读到,然后才物理清理。
  4. 触发条件:min.cleanable.dirty.ratio(默认 0.5)、min/max.compaction.lag.mssegment.ms,缺一不可正确触发。
  5. active segment 永远不压缩,所以你刚写的消息至少要等 segment 滚动后才被压缩。
  6. 应用场景:KStream changelog、CDC 主键最终态、__consumer_offsets__transaction_state__cluster_metadata
  7. 性能上 Cleaner 抢 IO,需要监控 log-cleaner.log 与 JMX kafka.log:type=LogCleaner

11. 面试高频题

Q1. cleanup.policy=compact 和 delete 的区别?什么场景下选 compact?

考察点:保留策略选型。

标准答案

  • delete(默认):按时间(retention.ms,默认 7 天)或大小(retention.bytes)整段 Segment 删除,offset 连续。适合事件流(订单事件、点击日志)。
  • compact:同一个 Key 只保留最新 Value,其他物理删除,offset 不连续,永不按时间删除。适合「状态快照」语义:用户资料、商品价格、__consumer_offsets、KStream changelog、CDC 主键最终态。
  • compact,delete:两者叠加,保留每 Key 最新值的同时,还按时间老化整段。

加分项:指出 active segment 不参与 compact、compact 不影响 producer 写入吞吐。


Q2. Tombstone 是什么?为什么需要它?

考察点:删除语义。

标准答案

  • Tombstone 是「value=null」的一条特殊消息,表示「这个 Key 被删除」。
  • 不能直接物理删除是因为消费者需要显式接收到「删除」事件才能更新自己的本地状态(KTable 删 key、CDC 下游删除行)。
  • 所以 Compaction 先保留 tombstone 至少 delete.retention.ms(默认 24h),让所有消费者有时间消费到,然后才物理清理。

易错点:以为「null value 立刻被删」——实际 tombstone 也要走 Compaction 流程才物理清理。


Q3. Compaction 是怎么触发的?关键参数有哪些?

考察点:触发条件。

标准答案

  • min.cleanable.dirty.ratio(默认 0.5):dirty 区占整个 log 比例 ≥ 该值才触发。
  • min.compaction.lag.ms(默认 0):消息至少多久之后才允许压缩(防止刚写就被压)。
  • max.compaction.lag.ms(默认无限):消息至多多久必须压缩一轮(兜底)。
  • segment.ms / segment.bytes:active segment 多久 / 多大滚动一次(滚动后才能参与压缩)。
  • delete.retention.ms(默认 24h):tombstone 至少保留多久。
  • log.cleaner.threads(默认 1):Broker 全局 cleaner 线程数。

加分项:提到 cleaner 用一个 OffsetMap 哈希表(log.cleaner.dedupe.buffer.size)做去重,内存不足时压缩效果差。


Q4. 为什么 __consumer_offsets 用 compact 策略?

考察点:内部 Topic 的设计哲学。

标准答案

  • 每次 commit() 都向 __consumer_offsets 追加一条 (group, topic, partition) → offset 的消息。
  • 一个活跃消费者每秒可能 commit 几次,N 个消费组 × M 个分区 × 365 天 = 海量数据,普通 delete 策略会把磁盘吃爆。
  • 用 compact:相同 (group, topic, partition) 只保留最新 offset,存量恒定。
  • 同时永不按时间过期,避免「停机一周后 offset 丢失」(实际靠 offsets.retention.minutes 控制无活跃 offset 的清理)。

Q5. compact topic 上的消费者会重复看到同一个 Key 吗?为什么?

考察点:Compaction 的「最终一致」语义。

标准答案

  • 。Compaction 是后台异步的,新消息写入到被压缩之间存在时间窗口;active segment 永远不压;dirty ratio 没到阈值时也不压。
  • 所以消费者在压缩前从头读,会看到同一个 Key 的多个旧 Value
  • 业务必须能容忍——一般写成「每条 (key, value) 当作 PUT」,最后一次 PUT 就是最新值。
  • KStream / Flink 的 KTable 抽象天然适配这种语义。

加分项:指出这正是「一切都是 last-write-wins 的最终一致 KV」语义。


🎯 小练习:建一个 compact topic,写 1000 条 (k1..k10, v1..v100)(即每个 key 写 100 次),用 kafka-dump-log.sh --files xxx.log --print-data-log 看物理布局;调小 min.cleanable.dirty.ratio、调小 segment.ms 让 segment 频繁滚动,等几秒再 dump,对比 record 数量减少了多少——你会亲眼看到 Compaction 在干活


12. 附录:常见踩坑与排查

12.1 「我建了 compact topic,怎么没看到 Compaction?」

最常见原因是 active segment 永远不被压缩

定位方法:

  1. segment.ms 是不是设了很大(默认 7 天)。低流量 topic 几个月都不会滚一次 active segment。

  2. segment.bytes 是不是设了很大(默认 1 GB)。

  3. 临时调小:

    bash
    kafka-configs.sh --bootstrap-server 127.0.0.1:9092 \
        --alter --entity-type topics --entity-name my-topic \
        --add-config segment.ms=10000

    或者强制滚动:写一条消息后等过 segment.ms。

12.2 「Cleaner 报错 OOM」

Cleaner 用 log.cleaner.dedupe.buffer.size(默认 128 MB)构建 OffsetMap。如果分区里不同 Key 太多,hash 表撑爆。

解决:

  • 调大 log.cleaner.dedupe.buffer.size(生产建议 ≥ 256 MB)
  • 调大 log.cleaner.threads,把内存分散到多个线程

12.3 「__consumer_offsets 增长不停」

正常情况 compact 把存量保持稳定。如果一直增长:

  1. cleanup.policy 是不是被误改成 delete

    bash
    kafka-configs.sh --bootstrap-server 127.0.0.1:9092 \
        --describe --entity-type topics --entity-name __consumer_offsets
  2. log-cleaner.log 有没有报错。

  3. 看 dirty ratio:

    JMX: kafka.log:type=LogCleaner,name=max-dirty-percent

12.4 「Tombstone 没被消费者收到」

通常是消费者离线超过 delete.retention.ms(默认 24h)后才回来,那段时间 tombstone 已经被物理清理。

解决:

  • 调大 delete.retention.ms(业务能接受存储多保留就调大)
  • 消费者引入「定期对账」逻辑,定期全量扫一遍 compact topic 重建本地 KV state

12.5 「为什么 offset 跳号?我的程序读不下去了」

Compact topic 的 offset 本来就是不连续的——这是正常现象,消费者代码不能假定 offset 连续。

如果你的代码用 offset 做主键或者期望「下一个 offset = current + 1」,会读不下去。

正确写法:用 consumer.position() 拿到下次要读的 offset,不要自己 current + 1


13. 附录:Compact 与流处理结合实战

13.1 用 Compact Topic 实现「最新订单状态表」

python
# Producer: 订单状态变化时发到 compact topic
producer.produce(
    "orders.status.compact",
    key=order_id.encode(),
    value=json.dumps({"status": "PAID", "ts": now}).encode(),
)

# 删除订单:发 tombstone
producer.produce("orders.status.compact", key=order_id.encode(), value=None)

# 消费者: 启动时全量重建本地 KV
consumer.subscribe(["orders.status.compact"])
consumer.seek_to_beginning(consumer.assignment())
state = {}
while not at_end():
    msg = consumer.poll(1.0)
    if msg is None: continue
    if msg.value() is None:
        state.pop(msg.key(), None)        # tombstone → 删除
    else:
        state[msg.key()] = json.loads(msg.value())
# 重建完成,state 就是「订单当前状态」的快照

这就是 Streams KTable 的本质 —— 一切「全量状态 + 增量变更」的场景都可以这样做。

13.2 与 Debezium CDC 的搭配

Debezium 把 MySQL 的 binlog 转成 Kafka 消息,常用 compact topic:

  • Key:MySQL 主键(如 users.id
  • Value:行的最新内容(INSERT / UPDATE 都写一条);DELETE 时写一条 tombstone

下游(ClickHouse / Elasticsearch / Search Index)从 compact topic 重建索引:

  • 全量重建:从头消费整个 topic,自动得到「所有行的最终态快照」
  • 增量同步:从上次 offset 继续,只处理新变更

这是「CDC 的标配模式」,Debezium / Maxwell / Flink CDC 都这么干。


14. 附录:JMX 监控指标

观察 Compaction 健康度的关键 JMX 指标:

MBean含义
kafka.log:type=LogCleaner,name=max-dirty-percent全集群最大 dirty ratio(应 < 阈值)
kafka.log:type=LogCleanerManager,name=time-since-last-run-ms上次 cleaner 运行距今(不应过大)
kafka.log:type=LogCleaner,name=cleaner-recopy-percent「重写率」
kafka.log:type=LogCleaner,name=DeadThreadCount死掉的 cleaner 线程数(应 = 0)

报警建议:

  • max-dirty-percent > 0.8 持续 5min → 报警,可能 cleaner 跟不上写入
  • DeadThreadCount > 0 → 立即报警,去 log-cleaner.log 看异常栈

15. 总结:一句话记住每个概念

  • cleanup.policy=delete = 流水本,过期撕掉
  • cleanup.policy=compact = 配方册,只留最新版
  • Tombstone (null value) = 图书馆下架标签,先贴着等所有人看到,再扔掉
  • min.cleanable.dirty.ratio = 攒到一定份量才开工
  • min.compaction.lag.ms = 新书暂时不能下架
  • max.compaction.lag.ms = 再忙下架请求也得在 X 时间内处理
  • delete.retention.ms = 下架标签至少挂多久
  • active segment = 正在添新书的那个货架,不许动
  • KStream changelog = 用 compact topic 给本地状态做云端备份
  • __consumer_offsets = Kafka 自己也用 compact topic 存自己的状态

🎬 可视化演示

演示加载缓慢或样式异常?点此在新标签页打开 ↗

💻 示例代码

python
#!/usr/bin/env python3
"""
compaction_demo.py
==================
往 cleanup.policy=compact 的 topic 写入大量「同 Key 不同 Value」,
等几秒让 Cleaner 触发,然后从头消费看实际剩了多少条。

结论:物理记录数 << 写入次数(因为相同 Key 的旧值被压缩掉),但每个 Key
的最新 Value 仍然能读到。

用法:
  pip install confluent-kafka
  bash ../init.sh                         # 先建好 compact topic
  python compaction_demo.py write 10 100  # 10 个 Key,每个写 100 次
  python compaction_demo.py wait 30       # 等 30 秒让 Cleaner 工作
  python compaction_demo.py read          # 从头消费,看剩多少条
"""

import os
import sys
import time
import json
from collections import Counter

from confluent_kafka import Producer, Consumer

BOOTSTRAP = os.environ.get("KAFKA_BOOTSTRAP", "127.0.0.1:9092")
TOPIC = os.environ.get("TOPIC", "learn.14.compact")


def write(num_keys, repeat):
    p = Producer({"bootstrap.servers": BOOTSTRAP, "linger.ms": 5})
    total = num_keys * repeat
    print(f"writing {num_keys} keys × {repeat} updates = {total} records")
    t0 = time.time()
    for r in range(repeat):
        for k in range(num_keys):
            key = f"item-{k:03d}".encode()
            val = json.dumps({"version": r, "ts": time.time()}).encode()
            p.produce(TOPIC, key=key, value=val)
        if r % 10 == 0:
            p.poll(0)
    p.flush(30)
    print(f"done in {time.time()-t0:.1f}s; latest version per key = {repeat - 1}")


def wait_cleaner(secs):
    print(f"sleeping {secs}s to let Log Cleaner do its work …")
    print("(init.sh 把 segment.ms=5s, dirty ratio=0.1, max.lag=60s,"
          "所以一般 30s 内就能看到 compaction 效果)")
    for i in range(secs, 0, -1):
        sys.stdout.write(f"\r  remaining: {i:3d}s ")
        sys.stdout.flush()
        time.sleep(1)
    sys.stdout.write("\n")


def read():
    c = Consumer({
        "bootstrap.servers": BOOTSTRAP,
        "group.id": "compaction-reader-" + str(os.getpid()),
        "auto.offset.reset": "earliest",
        "enable.auto.commit": False,
    })
    c.subscribe([TOPIC])

    n = 0
    cnt = Counter()
    latest = {}
    end_t = time.time() + 5
    while time.time() < end_t:
        msg = c.poll(1.0)
        if msg is None: continue
        if msg.error(): continue
        n += 1
        end_t = time.time() + 2   # 每读到一条就续 2 秒
        k = msg.key().decode() if msg.key() else "<null>"
        cnt[k] += 1
        try:
            latest[k] = json.loads(msg.value())["version"] if msg.value() else None
        except Exception:
            pass
    c.close()

    print(f"\n=== 物理读到 {n} 条记录 ===")
    print(f"涉及 {len(cnt)} 个 Key")
    print()
    if cnt:
        avg = sum(cnt.values()) / len(cnt)
        print(f"平均每个 Key 还剩 {avg:.2f} 条记录")
        print(f"  ↑ 没压缩前每个 Key 应该有 N 条;压缩后理论 1 条;中间值表示压缩进行中")
        print()
        print("Key 现在的「最新版本号」:")
        for k in sorted(latest.keys())[:10]:
            print(f"  {k}: version={latest[k]} (occurrences={cnt[k]})")
        if len(latest) > 10:
            print(f"  ... ({len(latest)} keys total)")


def main():
    if len(sys.argv) < 2:
        print(__doc__); return
    cmd = sys.argv[1]
    if cmd == "write":
        num = int(sys.argv[2]) if len(sys.argv) > 2 else 10
        rep = int(sys.argv[3]) if len(sys.argv) > 3 else 100
        write(num, rep)
    elif cmd == "wait":
        secs = int(sys.argv[2]) if len(sys.argv) > 2 else 30
        wait_cleaner(secs)
    elif cmd == "read":
        read()
    else:
        print(__doc__)


if __name__ == "__main__":
    main()
python
#!/usr/bin/env python3
"""
tombstone_delete.py
===================
演示用 null value 实现「删除一个 Key」的完整流程:
  1) 先写一些 (key, value) 数据
  2) 再写 (key, null) 作为 Tombstone
  3) 等若干秒让 Cleaner 触发
  4) 从头消费,看到 Key 物理消失(或者剩一个 tombstone)

用法:
  pip install confluent-kafka
  bash ../init.sh
  python tombstone_delete.py demo

观察:
  - 第一次 read:Key=item-001 还有数据
  - 写 tombstone 后立刻 read:能看到 (item-001, null) + 历史 value
  - 等 Cleaner 触发:历史 value 消失,只剩 tombstone
  - 等 delete.retention.ms (init.sh 里设 10s) 过期 + Cleaner 再跑:tombstone 也消失
"""

import os
import sys
import time
import json

from confluent_kafka import Producer, Consumer

BOOTSTRAP = os.environ.get("KAFKA_BOOTSTRAP", "127.0.0.1:9092")
TOPIC = os.environ.get("TOPIC", "learn.14.compact")


def produce(records):
    """records: list of (key_str, value_str | None)"""
    p = Producer({"bootstrap.servers": BOOTSTRAP, "linger.ms": 5})
    for k, v in records:
        p.produce(TOPIC,
                  key=k.encode() if k else None,
                  value=v.encode() if v is not None else None)   # ← None 即 Tombstone
    p.flush(10)


def read_snapshot(label):
    print(f"\n--- snapshot: {label} ---")
    c = Consumer({
        "bootstrap.servers": BOOTSTRAP,
        "group.id": "ts-snap-" + str(os.getpid()) + "-" + str(time.time()),
        "auto.offset.reset": "earliest",
        "enable.auto.commit": False,
    })
    c.subscribe([TOPIC])
    items = {}
    deadline = time.time() + 4
    while time.time() < deadline:
        msg = c.poll(1.0)
        if msg is None: continue
        if msg.error(): continue
        deadline = time.time() + 1.5
        k = msg.key().decode() if msg.key() else "<null-key>"
        v = msg.value()
        items.setdefault(k, []).append({
            "offset": msg.offset(),
            "value": v.decode() if v else "TOMBSTONE",
        })
    c.close()
    for k in sorted(items.keys()):
        print(f"  Key={k!r}:")
        for r in items[k]:
            print(f"    offset={r['offset']} value={r['value']}")
    if not items:
        print("  (空)")


def wait(secs):
    print(f"\nsleeping {secs}s ...")
    for i in range(secs, 0, -1):
        sys.stdout.write(f"\r  remaining: {i:3d}s ")
        sys.stdout.flush()
        time.sleep(1)
    sys.stdout.write("\n")


def demo():
    print("STEP 1: 写入 3 个 Key 各 5 个版本")
    rs = []
    for v in range(5):
        for k in ["user-1", "user-2", "user-3"]:
            rs.append((k, json.dumps({"v": v, "name": f"{k}@v{v}"})))
    produce(rs)
    read_snapshot("写入后")

    print("\nSTEP 2: 给 user-2 发 Tombstone")
    produce([("user-2", None)])
    read_snapshot("发完 tombstone 立刻读")

    print("\nSTEP 3: 等 25 秒让 Cleaner 跑(init.sh 里 segment.ms=5s, dirty ratio=0.1)")
    wait(25)
    read_snapshot("Cleaner 跑过后(user-2 应该只剩 tombstone;user-1/3 应该只剩最新版本)")

    print("\nSTEP 4: 再等 25 秒让 delete.retention.ms (10s) 过期 + 下一轮 Cleaner")
    wait(25)
    read_snapshot("最终态(user-2 的 tombstone 也物理消失,只剩 user-1/3 的最新版本)")


def main():
    if len(sys.argv) < 2 or sys.argv[1] == "demo":
        demo()
    elif sys.argv[1] == "snapshot":
        read_snapshot(sys.argv[2] if len(sys.argv) > 2 else "now")
    else:
        print(__doc__)


if __name__ == "__main__":
    main()

compaction_demo.py ↗ · tombstone_delete.py ↗