TT Lab
开始
学习 学习路径 课程

订单重复到达,又消失了一次

提交时机导致消息丢失——回退与幂等处理

在 TT Lab 中继续学习

本实验在 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. 创建只有 1 个分区的主题 shipments,并放入从 order-1 到 order-5 的五行。
  2. 用组 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 就是“丢失”的。
  3. 用新的组 audit-svc 从头读取五条,并保存到 /root/kafka/ship-audit.txt。Kafka 里全部都还在。
  4. 把 ship-svc 的偏移量倒回最前面(--reset-offsets --to-earliest --execute),并把输出保存到 /root/kafka/ship-reset.txt。
  5. 用 SHIP_FIXED=1,让 ship-svc 重新读取五条,并通过管道传给 ship-process。processed.txt 会变成六行,order-1 必须出现两次——这就是倒回的代价:重复。
  6. 创建 /root/kafka/dedup.sh。它读取标准输入中的订单,只把 /root/kafka/seen.txt 中没有的写入 /root/kafka/processed-dedup.txt,并把它记录到 seen 中(如果存在环境变量 DEDUP_SEEN、DEDUP_OUT,就使用那些路径)。即使把 shipments 从头读两遍并通过管道传入,也必须只剩五行。
  7. 不带 --from-beginning,用新组 late-svc 读取 5 秒(0 条),再用 auto.offset.reset=earliest 的新组 early-svc 读取五条。在 /root/kafka/offset-reset.txt 中写两行:late_count=0 和 early_count=5。
  8. 在 /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 的行数)。

参考

五个发货订单

创建只有 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 的行数减去不同订单的数量。评分器会从同样的文件中重新数一遍。