重试造成重复到达——用幂等生产者防止重复
本实验在 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 个分区的主题
payments,并用默认配置放入两行(pay-1、pay-2)。 - 用
kafka-dump-log.sh转储payments-0的日志段,保存到/root/kafka/dump-default.txt。批次的producerId不能是 -1,baseSequence必须是 0(默认生产者是幂等的)。 - 用
enable.idempotence=false再放入一行(pay-3),再次转储并保存到/root/kafka/dump-noidem.txt。新批次必须是producerId: -1。 - 打开
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。日志里必须留下两份以上同一条记录。 - 在相同条件下,往主题
payments-idem中开启幂等(enable.idempotence=true)放入,并把 stderr 保存到/root/kafka/idem-producer.log。重试会发生,但日志里必须只留下一份。 - 在 slow 已关闭的状态下,往主题
payments-app放入pay-77,方法是分两次单独运行生产者(保持幂等开启)。把转储保存到/root/kafka/two-sessions.txt——两个批次的producerId必须互不相同。 - 在
/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-dump-log.sh --files /var/lib/kafka/<토픽>-0/00000000000000000000.log --print-data-log(占位符为主题)。批次行中会出现producerId和baseSequence,记录行中会出现payload。 - 记录数:
kafka-get-offsets.sh --bootstrap-server localhost:9092 --topic <토픽>(占位符为主题)的最后一个数字(日志末尾偏移量)。 - 客户端配置通过
--command-property 키=값(占位符依次为键、值)给出(在 4.3 中--producer-property已被弃用)。5,000 个字符的一行用head -c 5000 /dev/zero | tr '\0' x > big.txt; echo >> big.txt来生成。 - 生产者配置文档:
delivery.timeout.ms必须大于等于request.timeout.ms + linger.ms,要开启幂等,必须满足acks=all、retries>0、max.in.flight.requests.per.connection<=5。如果在明确开启幂等的同时给出冲突的值,就会是 ConfigException。 - 常见错误 1:开着 slow 做完第 4 步却不关掉。小请求不受影响,很难察觉,但处理大记录的下一步会变慢。
- 常见错误 2:在第 6 步中往同一个生产者里放入两行。那属于同一个会话,序号是连续的,会成为一个批次。必须是“分两次单独运行”,会话才是两个。
用默认生产者发送两条
创建只有 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 上数这三个数字并做比较。