重复到达——重试造成的重复与幂等生产者
一句话总结
生产者没有收到响应时,无从得知消息是否已经提交,所以会重新发送, 如果原来的请求其实成功了,日志里就会留下两份。幂等生产者借助 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 的样子。生产者在放进套接字缓冲区的
那一刻就报告成功,之后发生什么都不知道。对于对延迟敏感的日志收集,
这是可以接受的交易,而对于支付则不行。
下一项实验要做什么
用转储比较默认生产者和关闭了幂等的生产者的批次头部,打开慢速 网络,重现重试留下三份、而幂等使其变成一份的情形, 然后分两次单独运行生产者,看到跨会话的重新发送无法被拦住。