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

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

重试造成重复到达——用幂等生产者防止重复

在 TT Lab 中继续学习

本实验在 VM 中运行

Ubuntu VM 上以 KRaft 单节点方式运行着 Apache Kafka 4.3.1 (localhost:9092)。辅助工具 kafka-lab-slow on|off 只会把发往 broker 的 4KB 以上的 请求延迟 600ms,从而引发生产者的超时和重试。首次启动大约需要 4 分钟。

目标

重现生产者没有收到响应而重试时,日志里留下同一条记录的多份副本的情形, 并用 kafka-dump-log.sh 确认幂等生产者把它变成一份。此外,还要 看看幂等覆盖不到的地方(跨生产者会话的重新发送)。

为什么重要

正如设计文档所说,生产者遇到网络错误时,无法得知这个错误是在消息 提交之前还是之后发生的。所以它会重新发送, 而如果原来的请求其实成功了,日志里就会留下两份——这就是 at-least-once。 幂等生产者的做法是:broker 给每个生产者分配一个 ID,并给每条记录附上序号, 过滤掉相同的序号,以此来防止这种情况。在 4.x 中它默认是开启的 (enable.idempotence 的默认值是 true),但只要动错一个配置,它就会悄悄 关闭,而且跨会话的重新发送它本来就不覆盖。亲眼看过这三点, 收到“到了两次”的反馈时,就知道该看哪里。

步骤

  1. 创建只有 1 个分区的主题 payments,并用默认配置放入两行(pay-1、pay-2)。
  2. 用 kafka-dump-log.sh 转储 payments-0 的日志段,保存到 /root/kafka/dump-default.txt。批次的 producerId 不能是 -1,baseSequence 必须是 0(默认生产者是幂等的)。
  3. 用 enable.idempotence=false 再放入一行(pay-3),再次转储并保存到 /root/kafka/dump-noidem.txt。新批次必须是 producerId: -1。
  4. 打开 kafka-lab-slow on,创建主题 payments-dup,然后把一行 5,000 个字符的内容(/root/kafka/big.txt)在关闭幂等的情况下,以 request.timeout.ms=300、delivery.timeout.ms=2000、retries=2、max.block.ms=10000 放入。把生产者的 stderr 保存到 /root/kafka/dup-producer.log,结束后执行 kafka-lab-slow off。日志里必须留下两份以上同一条记录。
  5. 在相同条件下,往主题 payments-idem 中开启幂等(enable.idempotence=true)放入,并把 stderr 保存到 /root/kafka/idem-producer.log。重试会发生,但日志里必须只留下一份。
  6. 在 slow 已关闭的状态下,往主题 payments-app 放入 pay-77,方法是分两次单独运行生产者(保持幂等开启)。把转储保存到 /root/kafka/two-sessions.txt——两个批次的 producerId 必须互不相同。
  7. 在 /root/kafka/producer-report.txt 中写五行:dup_copies=<payments-dup 의 레코드 수>、idem_copies=<payments-idem 의 레코드 수>、two_sessions_copies=<payments-app 의 레코드 수>、acks_default=all、idempotence_default=true(占位符依次为 payments-dup 的记录数、payments-idem 的记录数、payments-app 的记录数)。

参考

用默认生产者发送两条

创建只有 1 个分区的主题 payments,并用默认配置放入两行(pay-1、pay-2)。

先执行 --create --topic payments --partitions 1,再执行 printf 'pay-1\npay-2\n' | kafka-console-producer.sh ...。记录数用 kafka-get-offsets.sh 查看。

默认生产者是幂等的

用 kafka-dump-log.sh 转储 payments-0 的日志段,保存到 /root/kafka/dump-default.txt。批次的 producerId 不能是 -1,baseSequence 必须是 0(默认生产者是幂等的)。

kafka-dump-log.sh --files /var/lib/kafka/payments-0/00000000000000000000.log --print-data-log。broker 分配给生产者的 ID 和记录序号,原样写在批次头部里——这就是过滤重复的依据。

关闭幂等就没有 ID

用 enable.idempotence=false 再放入一行(pay-3),再次转储并保存到 /root/kafka/dump-noidem.txt。新批次必须是 producerId: -1。

--command-property enable.idempotence=false。关闭了幂等的生产者不会获得 ID,所以即使同一条记录再次到来,broker 也没有办法识别。

重试造成第二次到达

打开 kafka-lab-slow on,创建主题 payments-dup,然后把一行 5,000 个字符的内容(/root/kafka/big.txt)在关闭幂等的情况下,以 request.timeout.ms=300、delivery.timeout.ms=2000、retries=2、max.block.ms=10000 放入。把生产者的 stderr 保存到 /root/kafka/dup-producer.log,结束后执行 kafka-lab-slow off。日志里必须留下两份以上同一条记录。

如果请求在 300ms 内没有得到应答,生产者就会重新发送同一个批次(retries),而 broker 会把迟到的原请求和迟到的重试全都写入。生产者最后会报告失败,但日志里却有三份。用 kafka-get-offsets.sh 数一数。

幂等生产者只留下一份

在相同条件下,往主题 payments-idem 中开启幂等(enable.idempotence=true)放入,并把 stderr 保存到 /root/kafka/idem-producer.log。重试会发生,但日志里必须只留下一份。

和第 4 步一样打开 slow 再发送,只把 enable.idempotence=true 改掉。被重试的批次带着相同的序号到来,所以 broker 会把它过滤掉。生产者的日志里仍然会打印 REQUEST_TIMED_OUT——重试是发生了,只是没有重复。

会话不同,幂等就覆盖不了

在 slow 已关闭的状态下,往主题 payments-app 放入 pay-77,方法是分两次单独运行生产者(保持幂等开启)。把转储保存到 /root/kafka/two-sessions.txt——两个批次的 producerId 必须互不相同。

执行两次 printf 'pay-77\n' | kafka-console-producer.sh ...。进程不同,broker 就会分配新的生产者 ID,序号也从 0 开始,所以在 broker 看来是不同的记录。应用在失败之后重新发送就是这种样子,能阻止它的不是幂等生产者,而是像订单号这样的业务键。

用三个主题的记录数来总结

在 /root/kafka/producer-report.txt 中写五行:dup_copies=<payments-dup 의 레코드 수>、idem_copies=<payments-idem 의 레코드 수>、two_sessions_copies=<payments-app 의 레코드 수>、acks_default=all、idempotence_default=true(占位符依次为 payments-dup 的记录数、payments-idem 的记录数、payments-app 的记录数)。

三个数字是 kafka-get-offsets.sh 的最后一个字段。另外两个是生产者配置文档中的默认值——评分器会现在重新从 broker 上数这三个数字并做比较。