提交时机导致消息丢失——回退与幂等处理
本实验在 VM 中运行
Ubuntu VM 上以 KRaft 单节点方式运行着 Apache Kafka 4.3.1
(localhost:9092)。辅助工具 ship-process 模拟发货处理程序——它会把标准
输入中的订单全部读完,然后逐条写入 /root/kafka/processed.txt,
如果 SHIP_FIXED=1 没有设置,就会在 order-2 上崩溃。首次启动大约需要 4 分钟。
目标
重现消费者在处理之前先提交偏移量,崩溃之后消息“丢失”的情形,
再用另一个组确认那些消息在 Kafka 里其实还在,然后倒回组的
偏移量重新处理。倒回所产生的重复,用幂等消费者来防范,
并观察 auto.offset.reset 如何决定新组的起始位置。
为什么重要
设计文档的“消息传递语义”一节,就是这个实验的剧本。读取 → 保存位置
→ 处理,那么在处理过程中崩溃时,这条消息就不会再来了
(at-most-once)。读取 → 处理 → 保存,那么在保存之前崩溃时,它会再次
到来(at-least-once)。像控制台消费者那样依赖自动提交(enable.auto.commit 默认
true,5 秒间隔)的处理程序,更接近前一种形式。“消失了一次”
通常就是这个原因,而把它修好之后就成了“到了两次”——所以处理程序必须是幂等的。
文档把它写作“消息带有主键,因此更新是幂等的情形”。
步骤
- 创建只有 1 个分区的主题
shipments,并放入从order-1到order-5的五行。 - 用组
ship-svc从头读取三条,通过管道传给ship-process(它会崩溃)。之后对ship-svc执行 describe,并保存到/root/kafka/ship-crash.txt。CURRENT-OFFSET 是 3,而/root/kafka/processed.txt中必须只有order-1一条——order-2和order-3就是“丢失”的。 - 用新的组
audit-svc从头读取五条,并保存到/root/kafka/ship-audit.txt。Kafka 里全部都还在。 - 把
ship-svc的偏移量倒回最前面(--reset-offsets --to-earliest --execute),并把输出保存到/root/kafka/ship-reset.txt。 - 用
SHIP_FIXED=1,让ship-svc重新读取五条,并通过管道传给ship-process。processed.txt会变成六行,order-1必须出现两次——这就是倒回的代价:重复。 - 创建
/root/kafka/dedup.sh。它读取标准输入中的订单,只把/root/kafka/seen.txt中没有的写入/root/kafka/processed-dedup.txt,并把它记录到 seen 中(如果存在环境变量DEDUP_SEEN、DEDUP_OUT,就使用那些路径)。即使把shipments从头读两遍并通过管道传入,也必须只剩五行。 - 不带
--from-beginning,用新组late-svc读取 5 秒(0 条),再用auto.offset.reset=earliest的新组early-svc读取五条。在/root/kafka/offset-reset.txt中写两行:late_count=0和early_count=5。 - 在
/root/kafka/consumer-report.txt中写三行:lost_after_crash=<2단계에서 사라진 건수>、duplicates_after_reset=<5단계 뒤 processed.txt 의 중복 건수>、unique_orders=<processed-dedup.txt 의 줄 수>(占位符依次为第 2 步中丢失的条数、第 5 步之后 processed.txt 中重复的条数、processed-dedup.txt 的行数)。
参考
- 通过组读取:
kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic shipments --group ship-svc --from-beginning --max-messages 3 --timeout-ms 8000 | ship-process - 倒回:
kafka-consumer-groups.sh ... --reset-offsets --group ship-svc --topic shipments --to-earliest --execute。正如运维文档所说,消费者必须处于停止状态。如果不带--execute运行,就只显示计划。 - 新组的起始位置:消费者配置文档中的
auto.offset.reset——默认是latest,所以如果组里没有偏移量,就只读取现在之后的内容。用--command-property auto.offset.reset=earliest来修改。--from-beginning是控制台工具做同一件事的快捷方式。 - 常见错误 1:在第 2 步省略
--max-messages。如果把五条全读了,“丢失的两条”就无法重现。 - 常见错误 2:把 dedup 的状态(seen)只放在内存里。进程一挂,状态也跟着没了,下一次运行又会产生重复。无论是文件还是 DB,都必须和处理结果一起留存。
五个发货订单
创建只有 1 个分区的主题 shipments,并放入从 order-1 到 order-5 的五行。
printf 'order-1\norder-2\norder-3\norder-4\norder-5\n' | kafka-console-producer.sh ...。只有一个分区,所以顺序会完全保持。
处理之前先提交,就会丢失
用组 ship-svc 从头读取三条,通过管道传给 ship-process(它会崩溃)。之后对 ship-svc 执行 describe,并保存到 /root/kafka/ship-crash.txt。CURRENT-OFFSET 是 3,而 /root/kafka/processed.txt 中必须只有 order-1 一条——order-2 和 order-3 就是“丢失”的。
--group ship-svc --from-beginning --max-messages 3 --timeout-ms 8000 | ship-process。控制台消费者交出三条之后,提交偏移量 3 就结束了,而处理程序在第二条上崩溃了。下一次用这个组读取,就会从 3 开始。
在 Kafka 里它们仍然在
用新的组 audit-svc 从头读取五条,并保存到 /root/kafka/ship-audit.txt。Kafka 里全部都还在。
--group audit-svc --from-beginning --max-messages 5 --timeout-ms 8000 > /root/kafka/ship-audit.txt。消费并不是删除——丢失的不是消息,而是 ship-svc 的位置。
倒回组的偏移量
把 ship-svc 的偏移量倒回最前面(--reset-offsets --to-earliest --execute),并把输出保存到 /root/kafka/ship-reset.txt。
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --reset-offsets --group ship-svc --topic shipments --to-earliest --execute。这之所以能做到,是因为消费者的位置只是一个整数——设计文档把它称为违背队列契约、但必不可少的功能。
重新处理,就会到两次
用 SHIP_FIXED=1,让 ship-svc 重新读取五条,并通过管道传给 ship-process。processed.txt 会变成六行,order-1 必须出现两次——这就是倒回的代价:重复。
... --group ship-svc --max-messages 5 --timeout-ms 8000 | SHIP_FIXED=1 ship-process。如果组里已经有偏移量,--from-beginning 会被忽略——从倒回后的位置(0)开始读取。丢失的两条会回来,但已经处理过的 order-1 也会再来一遍。
幂等消费者
创建 /root/kafka/dedup.sh。它读取标准输入中的订单,只把 /root/kafka/seen.txt 中没有的写入 /root/kafka/processed-dedup.txt,并把它记录到 seen 中(如果存在环境变量 DEDUP_SEEN、DEDUP_OUT,就使用那些路径)。即使把 shipments 从头读两遍并通过管道传入,也必须只剩五行。
用 grep -qxF "$line" "$SEEN" 确认是否见过,只有没见过时才写入两个文件。不使用组,用 --from-beginning --max-messages 5 读两遍并通过管道传入。评分器会通过环境变量给出临时路径,并送入混有重复的输入来测试。
新的组从哪里开始
不带 --from-beginning,用新组 late-svc 读取 5 秒(0 条),再用 auto.offset.reset=earliest 的新组 early-svc 读取五条。在 /root/kafka/offset-reset.txt 中写两行:late_count=0 和 early_count=5。
--group late-svc --timeout-ms 5000 什么都读不到就结束了(默认 latest)。--group early-svc --command-property auto.offset.reset=earliest --max-messages 5 --timeout-ms 8000 会读取五条。用 wc -l 数一数,再写入文件。
数一数丢失的和来了两次的
在 /root/kafka/consumer-report.txt 中写三行:lost_after_crash=<2단계에서 사라진 건수>、duplicates_after_reset=<5단계 뒤 processed.txt 의 중복 건수>、unique_orders=<processed-dedup.txt 의 줄 수>(占位符依次为第 2 步中丢失的条数、第 5 步之后 processed.txt 中重复的条数、processed-dedup.txt 的行数)。
丢失的条数,是 ship-crash.txt 中的 CURRENT-OFFSET 减去当时已处理的条数(1);重复的条数,是 processed.txt 的行数减去不同订单的数量。评分器会从同样的文件中重新数一遍。