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

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

重复到达——重试造成的重复与幂等生产者

在 TT Lab 中继续学习

一句话总结

生产者没有收到响应时,无从得知消息是否已经提交,所以会重新发送, 如果原来的请求其实成功了,日志里就会留下两份。幂等生产者借助 broker 分配的 生产者 ID 和记录序号过滤掉重新发送,只保留一份。在 4.x 中 默认是开启的,但它无法覆盖跨会话的重新发送。

为什么需要它

设计文档中的“消息传递语义” 一节,准确地写出了这个问题。生产者在发布过程中遇到网络错误时,无法得知这个错误 是在消息提交之前还是之后发生的—— 就像往带有自动生成主键的表里 INSERT 时连接断开一样。0.11 之前的 生产者除了重新发送别无选择,所以是 at-least-once。如果原来的 请求其实成功了,重新发送就会往日志里把同一条消息再写一遍。

从 0.11 起出现了幂等投递选项。broker 给每个生产者分配一个 ID,生产者在 每条消息上附带一个序号(sequence)发送,broker 会把相同 ID、相同序号的消息 过滤掉。同一时期还加入了事务,使其可以原子地写入多个分区。本课程 只讲到幂等——事务是在它之上的一层。

工作原理

acks 决定要等待什么。 生产者配置文档 中的 acks 条目——0 完全不等待服务器的确认,所以无法保证对方已收到, 也不会发生重试(偏移量始终是 -1)。1 是领导者写入自己的日志后就应答, 所以如果在追随者复制之前领导者宕机,就会丢失。all(= -1)会等待整个 ISR 都确认,只要 ISR 中有一个还活着就不会丢失,是最强的 保证。默认值是 all,要开启幂等,也必须是 all。

重试默认几乎是无限的。 retries 的默认值是 2147483647,文档建议 不要动这个值,而是用 delivery.timeout.ms(默认 120000)来控制重试的 总时长。delivery.timeout.ms 必须大于等于 request.timeout.ms(默认 30000)+ linger.ms。request.timeout.ms 条目里有这样一句话—— 这个值必须大于 broker 的 replica.lag.time.max.ms(默认 30000),才能减少 由不必要的重试造成的消息重复的可能性。这是文档明确把重复写成 重试结果的地方。

幂等是有条件的。 enable.idempotence 条目——开启之后,每条消息只会 在流中写入一份;关闭的话,因 broker 故障等引起的重试可能会写入重复。 要开启,必须满足 max.in.flight.requests.per.connection 小于等于 5、retries 大于 0、 acks 为 all。如果存在冲突的配置,并且没有明确开启幂等,幂等就会悄悄地被关闭。 如果明确开启了 却又有冲突,则是 ConfigException。所以,一旦因为“性能”而设置了 acks=1,防重复的保护 就消失了,却没有任何警告。

멱등 프로듀서의 배치 헤더 (kafka-dump-log.sh)
  producerId: 1  producerEpoch: 0  baseSequence: 0  lastSequence: 1   ← ID 와 순번
멱등을 끈 배치
  producerId: -1 producerEpoch: -1 baseSequence: -1 lastSequence: -1  ← 걸러 낼 재료가 없다

broker 通过这些头部来识别重新发送。在实测中(4.3.1,给本地 broker 加上 600ms 的延迟, 并设置 request.timeout.ms=300、retries=2),关闭幂等的生产者把同一条 记录留下了三份,开启幂等的生产者只留下一份——两个生产者 最后都报告“失败”。只是没有收到响应,broker 上其实是有的。

幂等不覆盖的东西。 ID 和序号属于生产者会话。进程 重新启动就会拿到新的 ID,序号也从 0 开始,所以应用在收到失败报告之后 再次 send,在 broker 看来就是新的记录。transactional.id 条目把这称为 “跨越多个生产者会话的可靠性”,并把它划归事务的范畴。 如果没有那一层,答案就是在消费者一侧用业务键(订单号)来过滤。

在现场相遇的样子

支付服务收到“同一笔支付被记录了两次”的反馈时,第一步是检查生产者的 配置。如果有 acks=1 或 max.in.flight=10 之类的值,幂等就是被悄悄关闭的。 第二步是检查日志中的批次头部——如果同一笔支付出现在 producerId 不同的两个批次里, 那就是应用层的重新发送,这无法通过生产者配置来 阻止。

反过来,“说是发了却没有”就是 acks=0 的样子。生产者在放进套接字缓冲区的 那一刻就报告成功,之后发生什么都不知道。对于对延迟敏感的日志收集, 这是可以接受的交易,而对于支付则不行。

下一项实验要做什么

用转储比较默认生产者和关闭了幂等的生产者的批次头部,打开慢速 网络,重现重试留下三份、而幂等使其变成一份的情形, 然后分两次单独运行生产者,看到跨会话的重新发送无法被拦住。