消息丢失——提交时机与至多一次、至少一次语义
一句话总结
消费者保存位置的时机决定了语义。如果在处理之前保存,那么处理过程中
崩溃的消息就不会再来(at-most-once,“丢失”);如果在处理之后保存,
那么在保存之前崩溃的消息就会再来(at-least-once,“到了两次”)。选择后者
并让处理程序具备幂等性才是答案,而倒回位置(--reset-offsets)是找回丢失内容
的工具。
为什么需要它
设计文档的“消息传递语义” 一节,用两段话就讲完了消费者一侧。如果消费者读取消息、保存位置之后再 处理,那么在保存之后、处理之前崩溃时,接手的进程会从已保存的位置 开始,所以在它之前的消息就没有被处理——at-most-once。如果读取、处理之后再 保存,那么在处理之后、保存之前崩溃时,接手的进程会再次收到已经处理过的消息 ——at-least-once。文档还补充说:在很多情况下,消息带有主 键,更新是幂等的(即使收到同一条消息两次,也只是覆盖同一条 记录而已)。
这两段话就是课程标题里的两起事故。“消失的那一次”是处理之前保存的结果, 而把它修好之后,就变成了“到了两次”。并不是二选一,而是选择 “到了两次”,并让处理两次也没关系,这才是设计。
工作原理
自动提交与处理无关,照常运行。 消费者配置文档
中的 enable.auto.commit 默认是 true,每隔 auto.commit.interval.ms(默认 5000),
偏移量就会在后台被提交。像控制台消费者这样依赖自动提交的
处理程序,提交的是“读过的”而不是“处理过的”。如果消费者交出三条
并提交了偏移量 3,而处理程序在第二条上崩溃,那么下一个消费者就会从
3 开始读取。第二条和第三条仍然好好地留在 Kafka 里,但这个组再也
不会经过它们了。
shipments P0: order-1 order-2 order-3 order-4 order-5
ship-svc 읽음 3건 → 커밋 3 → 처리기 order-2 에서 크래시
처리됨: order-1 사라짐: order-2, order-3
丢失的不是消息,而是位置。 消费并不是删除,所以用另一个组
从头读取,五条都在。而且可以用运维文档
中的 --reset-offsets 把组的位置倒回去——有 --to-earliest、
--to-latest、--to-offset、--shift-by、--to-datetime 等场景,
必须加上 --execute 才会真正改变(没有它就只显示计划),而且消费者
实例必须处于停止状态。倒回之后,丢失的两条会回来,但已经
处理过的 order-1 也会再来一遍。倒回的代价就是重复。
幂等消费者。 就是在处理程序一侧,实现文档所说的“有主键,所以更新是幂等的情形”。 把已处理的订单号和结果保存在同一个地方,如果传来的 订单已经存在,就跳过。如果把状态只放在内存里,进程一挂它也就 跟着没了,下一次运行又会产生重复。无论是文件还是 DB,都必须和处理结果一起 留存下来。
新的组从哪里开始。 auto.offset.reset 条目——决定当组里没有已提交的
偏移量,或者那个偏移量已经不存在(数据被删除)时,该怎么办。earliest 是回到最前面,latest(默认)是回到最后面,
by_duration:<ISO8601> 是从现在回退那么长的时间,none 则抛出异常。
因为默认是 latest,所以新建一个组、不做任何配置就接上去,会跳过此前所有的消息
——这正是“新服务没有看到旧订单”这类反馈的原因。控制台
工具的 --from-beginning 就是把它改成 earliest 的快捷方式,如果组里
已经有偏移量,则会被忽略。
提交会消失。 broker 配置文档
中的 offsets.retention.minutes(默认 10080,即 7 天)——如果组一直是空的,或者停止订阅主题,
超过这个时间,已提交的偏移量就会被丢弃。之后消费者
回来时,就成了没有偏移量的状态,会适用 auto.offset.reset。停了一周的
批处理消费者回来之后,如果以 latest 开始,就会把这期间的数据全部跳过。
在现场相遇的样子
收到“订单丢了”的反馈时,顺序是这样的:对组执行 describe,查看已提交的
偏移量;再用另一个组(或者不用组)读取那个区间,确认消息是否存在;
如果存在,就是位置的问题——在消费者停止的状态下,用 --reset-offsets --to-offset 倒回并重新处理。重新处理会不会产生重复,取决于处理程序
是否幂等,如果不是,就得在倒回之前先把它修好。
反过来,“同一个订单被处理了两次”这类反馈,通常是再平衡或重启之后 at-least-once 的正常行为。如果把它当作 bug,并把提交提前,下一条反馈 就变成“丢了”。这两条反馈是同一个旋钮的两端。
下一项实验要做什么
重现发货处理程序在第二个订单上崩溃、导致两条消息丢失的情形,用另一个
组确认它们仍然留在 Kafka 里,然后倒回组并重新处理,
看到产生重复,再用幂等处理程序把它拦住。最后用新的组比较 auto.offset.reset
的两个取值。