主题
第 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_new和v_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.ratio | 0.5 | dirty 部分占整个 partition 的比例 ≥ 该值才触发 |
min.compaction.lag.ms | 0 | 消息距今至少 N ms 才能被压缩(不动太新的) |
max.compaction.lag.ms | 长整数最大 | 消息距今最多多久必须被压缩(再忙也得轮到) |
segment.ms | 7 天 | active segment 多久滚动一次(滚动后才能参与压缩) |
segment.bytes | 1 GB | 同上,按大小滚动 |
delete.retention.ms | 24 h | tombstone 至少保留多久 |
log.cleaner.threads | 1 | Broker 全局 cleaner 线程数(默认 1,I/O 重的可以加) |
log.cleaner.dedupe.buffer.size | 128 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=3006.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的 topiccompaction_demo.py:写入大量同 Key 不同 Value,触发 compaction,前后对比tombstone_delete.py:用 null value 删除 Key 的演示
9. 与其他 MQ 的对比
| 维度 | Kafka Compaction | RabbitMQ | RocketMQ | Pulsar |
|---|---|---|---|---|
| 是否支持 | ✅ 原生 | ❌ | ❌(消息只按时间过期) | ✅ Topic Compaction |
| 用途 | 状态快照、changelog、__consumer_offsets | 队列即时消费 | 同 RabbitMQ | 类 Kafka |
| 删除标记 | null value (tombstone) | N/A | N/A | null value |
| Key 的作用 | 必须有 key 才能压缩 | routing key(仅路由用) | message tag(仅过滤) | key 用于 partition + compaction |
📌 本质差异:Kafka 把「日志型存储」和「KV Store」用一个文件格式统一起来——这是 RabbitMQ / RocketMQ 都没做的设计。Pulsar 学了这一点。
10. 小结
- 三种
cleanup.policy:delete(默认,按时间/大小删段)、compact(同 Key 最新值)、compact,delete(叠加)。 - Compaction = 把「无限追加日志」变成「最终态 KV 快照」,保留每个 Key 的最新 Value。
- Tombstone(null value)= 删除标记,先保留
delete.retention.ms(默认 24h)让消费者读到,然后才物理清理。 - 触发条件:
min.cleanable.dirty.ratio(默认 0.5)、min/max.compaction.lag.ms、segment.ms,缺一不可正确触发。 - active segment 永远不压缩,所以你刚写的消息至少要等 segment 滚动后才被压缩。
- 应用场景:KStream changelog、CDC 主键最终态、
__consumer_offsets、__transaction_state、__cluster_metadata。 - 性能上 Cleaner 抢 IO,需要监控
log-cleaner.log与 JMXkafka.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」语义。
🎯 小练习:建一个
compacttopic,写 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 永远不被压缩。
定位方法:
看
segment.ms是不是设了很大(默认 7 天)。低流量 topic 几个月都不会滚一次 active segment。看
segment.bytes是不是设了很大(默认 1 GB)。临时调小:
bashkafka-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 把存量保持稳定。如果一直增长:
看
cleanup.policy是不是被误改成delete:bashkafka-configs.sh --bootstrap-server 127.0.0.1:9092 \ --describe --entity-type topics --entity-name __consumer_offsets看
log-cleaner.log有没有报错。看 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()