重试与死信:失败消息怎样恢复,怎样不拖垮消费

从库存服务故障出发,解释失败分类、退避与重试预算,对比 RocketMQ、RabbitMQ 与 Kafka 的恢复路径,展开责任转移、顺序阻塞、死信治理和受控重放。

库存服务暂时不可用,消费者返回失败,消息稍后重试。听起来合理,但如果一万条消息都失败,每条消息又立即请求库存,库存刚恢复就可能被重试流量再次压垮。与此同时,一条永远无法解析的消息还会反复占用工作线程。

重试让未完成任务继续得到执行机会。它需要判断下一次尝试有没有意义,安排间隔,限制资源占用,并保存最终未解决的任务。死信队列是恢复责任的一个存放位置,进入死信并不代表业务已经完成。

本文沿“订单请求库存预占,库存短暂故障”的设计示例说明失败路径。RocketMQ 区分 4.x Remoting 回调与 5.0 消费模型;RabbitMQ 的死信安全配置以 4.2 quorum queue 文档为具体参考;Kafka 的重试 Topic 是应用设计模式,不把它描述为所有消费者都会自动启用的核心功能。

一、消费失败至少有四种不同含义

暂时失败包括连接池短时耗尽、下游过载和可恢复的数据库冲突。参数不合法、协议版本不支持、目标业务不允许当前动作,则不是等几秒就会自行消失的问题。业务前置事实尚未到达,也和网络错误不同,需要等待事实或查权威状态。

最容易误判的是结果未知。库存请求超时,可能已经预占成功;直接换一个预占键再调用,就可能重复占用。下一次尝试应该复用同一业务身份或先查询,而不是把所有异常都当成“没有执行过”。

失败类型 例子 下一步
暂时性故障 下游短时不可用、数据库可恢复冲突 有界退避后重试
永久输入错误 消息无法解析、参数违反协议 持久隔离并告警,修复后再恢复
前置事实缺失 订单尚未同步到下游 等待、查询或按业务流程补齐
结果未知 预占请求超时但可能已经成功 同键查询或幂等重试,必要时对账

错误分类最好来自协议和业务结果,不要只按异常类名或 HTTP 状态码机械处理。某个 500 可能是部署缺陷,重试一百次都没有意义;同一个数据库异常,也可能要求重做整个事务,而不是仅重复最后一条写入。

还要区分业务拒绝与基础设施失败。库存明确不足,可以正常记录拒绝事实并通知订单,不必让原消息一直失败等待;“库存服务没有回应”则不能伪造库存不足。消费成功表示已经正确处理这个输入,不一定表示业务请求被批准。

二、重试间隔、预算和隔离决定恢复速度

固定每秒重试容易形成同步尖峰。指数退避增加连续失败之间的间隔,jitter 在允许窗口内随机分散尝试。它们降低瞬间冲击,却不增加下游容量;持续大故障期间,旧重试和新消息仍会争抢相同资源。

示例策略,不是产品默认配置:
backoff(attempt) = min(cap, base × 2^attempt)
nextAttemptAt = now + random(0, backoff(attempt))
停止条件 = 最大尝试数、最长失败年龄或业务截止条件
执行条件 = 下游可用且重试预算仍有余量

最大尝试数应该明确是否包含首次执行。一次原始执行加三次重试,是四次尝试;不同 SDK 字段的命名可能不同。失败次数、首次失败时间和下一次尝试时间分开记录,才能解释一条消息为什么还在等。

重试预算限制恢复流量占用。可以为实时消息和历史恢复分别设置并发、速率与连接池边界。否则部署恢复时,把几天的死信同时投入主链路,会让正常用户请求也超时。隔离队列并不自动隔离数据库,共享连接池和下游配额仍需要控制。

多层重试会乘法放大。MQ 再投三次,消费者 HTTP 客户端每次三次,下游服务内部又三次,最底层可能面对远超原始流量的请求。应明确哪个层拥有业务恢复责任,内部仅做有限且安全的短重试,避免各层独立无限重试。

不应把消费失败当常态限流工具。大量正常消息因为限流进入失败路径,会污染重试计数,推高死信,并形成额外投递。过载时应减小接收与执行速度、暂停相应消费单元或调节并发,让正常积压留在正常队列。RocketMQ 官方也给出这一使用边界。

三、RocketMQ 的重试需要按消费模型理解

4.x 并发监听器返回 RECONSUME_LATER,适用模式下由 SDK 和 Broker 的重试路径安排后续消费;顺序监听器使用不同状态,对当前队列的推进也不同。不能把两个回调状态混用,更不能在已经无恢复记录地异步处理之后返回成功。4.x Push API

// 4.x Remoting、集群并发消费的业务结构示例。
try {
    applyOrPersistRecovery(message);
    return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
} catch (RetryableFailure failure) {
    return ConsumeConcurrentlyStatus.RECONSUME_LATER;
}
// 永久错误由 applyOrPersistRecovery 持久隔离;接管失败仍不能确认。

这里的辅助函数有严格含义:要么完成本地业务,要么已经把恢复任务提交到可靠存储并明确转移责任。只把错误写到普通日志不算接管;日志可能滚动清理,也没有自动恢复执行者。

5.0 文档把 PushConsumer 的失败等待和 SimpleConsumer 的不可见时间机制分开说明。Push 失败后按相应策略等待再交付;SimpleConsumer 的 InvisibleDuration 是收到消息后的处理保护窗口,未确认则可能重新可见,修改它有时效条件。5.0 消费重试

不可见时间太短,正常慢任务还没结束就收到第二份;太长,执行者崩溃后恢复被拖延。延长窗口只能调整重新交付时机,不会保证旧任务停止,也不会让外部调用幂等。长任务更适合持久接管后用业务状态机管理,而不是一直占住回调并不断延长。

重试上限、顺序消息与普通消息的间隔,需要看实际消费组元数据、SDK 与 Broker 版本。广播、集群和两代客户端不能共享一份无条件承诺。业务应该有自己的失败年龄与恢复状态,不只依赖默认次数。

四、RabbitMQ 重新入队与死信转发的区别

basic.nack 或 basic.reject 的 requeue=true 表示重新入队;它并不是“过一分钟再执行”。消息可能立即重新交付,多个消费者都因为同一个故障重新入队,就形成高频循环。确认与重新入队

requeue=false 在配置 DLX 时进入死信路径;没有 DLX 则可能直接丢弃。因此“拒绝永久坏消息”要先确认它有持久隔离目的地,而不是只改一个布尔参数。AMQP 的 delivery tag 属于对应 channel,也不能拿旧连接的 tag 去确认新连接上的消息。

可以借助固定 TTL 等待队列建立退避档位,到期后转发回来,但要审查队头影响和死信转发可靠性,详见 延迟与定时消息。原消息确认与应用重新发布之间也有双写窗口:先确认后发布会漏,先发布后确认会重复。

RabbitMQ 4.2 quorum queue 支持配置 at-least-once dead-lettering,默认策略仍需要核对,不能因为用了 quorum queue 就推断所有死信转发都有确认保护。对应配置要求匹配死信策略、溢出行为及其他条件。

更可靠的转发会保留未被目标确认的消息,也因此占用源队列资源。目标长时间不可用,可能阻止源队列继续接收。安全性与容量之间有真实取舍,需要监控目标路由、确认和空间,而不是只查看消费者报错。

五、Kafka 失败后选择原地等待还是迁移重试

Kafka 的提交位点是分区恢复位置。失败记录后面的输入已经处理,直接推进位点会把失败记录跳过;不推进则恢复后可能重读已经完成的部分。最简单的做法是暂停相关分区,有限重试并保持客户端必要的轮询和组协作,不要用长时间阻塞导致反复退出消费组。

原地等待保留分区顺序,但一个坏事件可能阻塞同分区其他业务键。若业务允许独立对象继续执行,可以把失败任务可靠转移到重试 Topic 或业务任务表,再推进原输入。转移后要自己保存原 Topic、partition、offset、eventId、尝试次数和下一次时间。

失败输入转移重试责任的正确边界

Kafka 输入位点和 Kafka 重试输出,可以在满足客户端、消费组及隔离条件的 Kafka 事务中协调;输入位点与外部数据库任务表则不是天然同一事务。可以在数据库持久建立恢复记录并以唯一键去重,再推进位点,允许故障窗口重读而不重复新建任务。Kafka 交付语义

迁移重试会改变顺序。订单支付消息进入重试 Topic,发货消息沿主 Topic 继续,必须由订单状态机阻止缺失支付时的发货。重试 Topic 只分离执行流量,不会自动保留主链路的因果依赖。

还要注意重试输出有自己的保留期与消费者。写入成功后没有恢复消费者、重试数据被清理、重新交付时间没人检查,都可能让可靠转移变成持久但永远不执行的任务。转移成功只代表责任被接管,必须继续观察终态。

六、顺序消费的坏消息不能随便跳过去

普通独立通知失败,可以隔离后继续处理其他通知;账户增量扣减失败,后一笔依赖前一笔的余额状态,跳过就可能破坏账务。决定是否阻塞,应根据业务前置条件,而不是因为顺序监听器不好调优就放弃顺序。

消息进入死信后,队列可能继续推进,但业务缺口仍存在。可以让下游保存对象断点及后续待办,等待原消息修复;或者在适合完整快照的投影场景,从权威状态重建。不可省略的扣款和发券不能仅靠最新状态覆盖补上。

阻塞范围也需要设计。按商家号排序可能让一个坏订单拖住整个商家;按订单号排序能缩小范围,但不能覆盖商家账务的跨订单依赖。MQ 的队列或 message group 与业务排序单元应该匹配,相关路由与扩容边界见 顺序消费。

针对协议损坏的消息,隔离记录必须保留原始载荷及解析错误,恢复时使用修复后的解释规则。不能在失败处理里强行填默认金额、默认用户,让错误输入变成看似成功的业务结果。

七、死信队列需要负责人和恢复状态

死信是自动重试暂时停止的任务集合。应记录原始身份、来源、首次失败时间、最近错误、尝试历史、业务键和已有副作用。只有消息体,没有失败上下文,恢复者很难判断应该修代码、补数据还是查询外部结果。

从失败隔离到修复重放再对账的流程

死信处理可以建立状态:待诊断、待修复、允许重放、重放中、完成、人工终止。负责人对业务结果负责,平台只提供存储和投递。把一条消息从死信移到主队列,应该保存重放批次和操作记录,不能以“队列清空了”作为完成标准。

重放之前先查询已发生的效果。退款可能已经成功,只是本地确认失败;再次调用与修复本地记录是不同动作。原始 eventId 和业务动作键保持稳定,新增 replayBatchId 用来观察恢复批次,不应通过换业务键绕过去重。重复消费与幂等解释了这个身份边界。

死信自身也会被清理或填满。需要根据允许恢复期限配置保留、容量和归档,并告警最老未处理年龄。关键任务还应通过业务表或对账结果识别未完成义务,避免死信过期后业务系统也失去恢复依据。

人工终止要有明确业务解释。过期营销通知可以放弃,已收款订单的权益未发放则不能简单删除。终止消息处理与终止业务责任不是同一个操作,有些任务需要补偿、退款或明确的人工结案。

八、一次故障怎样有界恢复

库存服务故障十分钟,消费者观察到暂时不可用,限制接收并在允许策略内退避。已完成的预占事件不重复生效,结果未知的请求按同一预占键查询。正常业务拒绝记录终态,不进入无意义的失败重试。

少数超出尝试或年龄限制的任务持久进入隔离。库存恢复后,先检查实时流量是否稳定,再开放有限的历史恢复并发。恢复任务发现原预占已经成功,补齐本地结果;确实没有预占才按原键重新请求。协议坏消息在修复前继续隔离,不能混进恢复批次反复失败。

假设下游每秒只能稳定处理一定数量请求,恢复配额必须给实时业务留空间。单纯增加消费线程可能提高连接等待和超时率,使新失败继续产生更多重试。恢复速率应根据下游成功率、尾延迟和资源余量调整,而不是只追求死信数量尽快归零。

按批次灰度重放更容易发现问题。先恢复少量具有代表性的任务,观察业务终态和重复约束,再扩大范围;每个批次能够暂停,不影响实时消费。恢复过程出现新的参数冲突或状态倒退时,暂停对应批次,保留未完成任务供进一步诊断。

九、观察失败链,而不是只数失败回调

重要指标包括首次失败率、重试流量占比、最终成功率、最老失败年龄、死信增长、恢复完成率和业务未闭环数量。消息系统的成功确认率很高,可能只是大量任务被持久接管;如果接管任务一直没有完成,业务仍然故障。

日志应能关联原事件、业务动作、来源位置与尝试,而不要求输出完整敏感载荷。重复失败的热点输入可以采样,关键状态转移与终态保留结构化记录。把大量相同异常打印到磁盘,也会成为故障期间的额外负担。

我会把一套重试方案看成完整恢复协议:下一次由谁尝试,何时尝试,允许多少资源,什么条件停止,停止后由谁处理,以及恢复怎样证明业务已经完成。每个问题能在持久状态中找到答案,重试才不会只是失败消息绕圈。

参考资料