延迟、重试、死信与补偿:安排失败消息的下一次处理
同一批消息里,连接超时可能稍后恢复,缺少金额字段的消息却会每次失败。消费端需要为这两种结果安排不同去处:前者等待下一次有限尝试,后者退出正常工作队列,保存原始内容供修复。处理已经造成了业务变化时,还要决定是否需要执行补偿。
根据失败原因安排下一步
发送重试与消费重试发生在不同位置
生产者没有收到发送结果时,Broker 可能尚未接收,也可能已经保存但确认在回程丢失。此时重发应保留业务事件 ID,并根据产品的确认、幂等生产与路由规则处理结果未知。消费重试则发生在消息已被交付之后:业务未完成,应用还没有安全确认当前交付。RabbitMQ 的可靠性说明
把两者合并成一个 retryCount 会丢掉重要信息。一个事件可以经历两次发送尝试、三次消费交付,其中只有一次实际进入数据库事务。日志宜分别记录 eventId、publishAttempt、deliveryAttempt、处理阶段和失败类别。
| 失败类型 | 例子 | 处理方式 |
|---|---|---|
| 暂态依赖失败 | 下游限流、短时连接失败 | 有间隔、有总期限的重试;必要时暂停整个下游方向 |
| 结果未知 | 请求超时,但对方可能已扣款 | 先按稳定操作 ID 查询结果;支持幂等时再重发 |
| 契约或权限错误 | 必填字段缺失、不支持的版本、签名不合法 | 隔离并报告,不靠等待修复数据 |
| 业务拒绝 | 订单已关闭、额度不足 | 根据业务契约记录终态,通常无需基础设施重试 |
| 状态先后冲突 | version=2 到达时本地仍是 version=0 | 等待缺失事件、重读事实或进入顺序修复队列 |
| 程序缺陷 | 特定合法数据稳定触发异常 | 保存样本、停止无效重试,修复版本后有控制地重放 |
错误码只能帮助分类,最终还要看操作语义。读取超时通常容易重试;扣款超时必须先确定是否已经发生扣款。
次数、间隔与总期限
设基础等待为 d,指数退避可使用 min(maxDelay, d × 2^(attempt-1)),再加入随机抖动,避免大量失败消息同一时刻返回。总次数与最晚完成时间应共同约束:剩余业务期限不足以完成下一次调用时,继续排队没有意义。
例如订单通知允许总共尝试 3 次,间隔分别为 5 秒和 20 秒,每次 HTTP 请求最多 2 秒。队列等待也计入业务期限;消费端拿到消息后应重新检查 deadline,而非只看已经尝试几次。连接池耗尽时增加重试次数,可能把一次故障扩成持续的资源争用。
nack/reject + requeue=true 可以使消息重新进入待交付状态,通常没有自动产生合理的退避。RabbitMQ 4.3 的 quorum queue 提供原生 delayed retry,策略可选择哪些返回行为参与延迟,并配置线性延迟上下限;采用它时应按对应版本的 acquired/delivery 计数语义配置上限。Quorum delayed retry 与投递限制
下面使用固定 TTL 的延迟队列展示一次消息怎样被重新路由,便于观察两次重试的具体去向;生产系统可以按期限、吞吐和运维要求选择原生延迟重试、分级延迟队列或数据库任务调度。
用真实队列运行有限重试与毒消息分流
环境、账号和拓扑
下载完整工程,解压进入 message-event-lab。以下命令面向 Linux shell、Docker Engine 与 Compose v2+,由有 Docker 权限的普通宿主用户执行。RabbitMQ 4.3.5 使用 vhost lab 中的 app 角色;PostgreSQL 18.6 使用 app 与 lab schema。管理初始化由 Compose 完成,应用不使用数据库超级用户。
compose.yaml 只启动一节点 RabbitMQ 和 PostgreSQL,端口不映射到宿主。它用于本机机制实验,消息副本仍只有一份。Java 客户端 5.33.0、pgJDBC 42.7.13、Maven 3.9.12;Maven 使用 Java 25 编译到 Java 17 字节码,实验运行 Java 25。
me12.retry.work
├── healthy → 写业务结果 → ACK
├── transient / attempt 1、2
│ └── 发布到 me12.retry.delay → publisher confirm → ACK 原交付
│ └── TTL 到期 → 默认交换机 → 回到 work
│ └── attempt 3 → 写业务结果 → ACK
└── poison → reject(requeue=false)
└── me12.retry.dead 交换机 → me12.retry.dlq延迟队列的 x-message-ttl=500 单位为毫秒,过期后通过 DLX 路由回工作队列。TTL 控制过期资格,Broker 调度和排队决定实际再次交付时间,不能承诺 500 毫秒准点执行。逐消息 TTL 还会受到队头过期处理影响,混放很长与很短的延迟任务容易造成观察上的队头阻塞。队列与消息 TTL
test "$(id -u)" -ne 0 || { echo '请使用普通宿主用户'; exit 1; }
mkdir -p .m2
docker compose up -d --wait --wait-timeout 120
docker run --rm --user "$(id -u):$(id -g)" \
-e MAVEN_CONFIG=/maven-cache \
-v "$PWD:/workspace" -v "$PWD/.m2:/maven-cache" -w /workspace \
maven:3.9.12-eclipse-temurin-25 \
mvn -B -ntp '-Dmaven.repo.local=/maven-cache' '-Duser.home=/tmp' \
-pl rabbit -am clean package dependency:copy-dependencies
run_retry() {
docker run --rm --user "$(id -u):$(id -g)" --network me12-lab_default \
-v "$PWD:/workspace:ro" -w /workspace eclipse-temurin:25.0.4_7-jdk \
java '-Dorg.slf4j.simpleLogger.defaultLogLevel=warn' \
-cp 'rabbit/target/classes:rabbit/target/dependency/*' example.RetryLab "$@"
}
run_retry retry构建成功后,关键结果为:
retry: healthyAttempts=1 transientAttempts=[1, 2, 3] businessRows=2 poisonDeadLettered=true程序实际读取每次重投的 attempt header,断言 transient 的尝试序列严格为 1、2、3。healthy 和 transient 各形成一条 PostgreSQL 结果;poison 不写业务表,进入 DLQ,并带有 reason=rejected 的 x-death 历史。代码最后确认并取走实验 DLQ 样本,使再次执行从已知数据开始;生产消费端不能把“记录日志后丢弃”当作死信处理完成。
完整入口与确认顺序
源码位置是 rabbit/src/main/java/example/RetryLab.java:
package example;
import com.rabbitmq.client.AMQP;
import com.rabbitmq.client.Channel;
import com.rabbitmq.client.ConnectionFactory;
import java.nio.charset.StandardCharsets;
import java.sql.DriverManager;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.concurrent.TimeUnit;
public final class RetryLab {
static final String WORK="me12.retry.work",DELAY="me12.retry.delay",DLQ="me12.retry.dlq";
static void check(boolean ok,String why){if(!ok)throw new IllegalStateException(why);}
static java.sql.Connection database() throws Exception {
return DriverManager.getConnection("jdbc:postgresql://postgres:5432/lab","app","app-lab-only");
}
static void publish(Channel ch,String queue,String id,String body,int attempt) throws Exception {
ch.basicPublish("",queue,true,new AMQP.BasicProperties.Builder().messageId(id).deliveryMode(2)
.headers(Map.of("attempt",attempt)).build(),body.getBytes(StandardCharsets.UTF_8));
ch.waitForConfirmsOrDie(5000);
}
static com.rabbitmq.client.GetResponse receive(Channel ch,String queue,int seconds) throws Exception {
long end=System.nanoTime()+TimeUnit.SECONDS.toNanos(seconds);
while(System.nanoTime()<end){var m=ch.basicGet(queue,false);if(m!=null)return m;Thread.sleep(20);}
throw new IllegalStateException("no message from "+queue);
}
static void successful(String event) throws Exception {
try(var c=database();var p=c.prepareStatement("INSERT INTO retry_effect VALUES (?) ON CONFLICT DO NOTHING")) {
p.setString(1,event);p.executeUpdate();
}
}
public static void main(String[] args) throws Exception {
check(args.length==1,"mode: retry|hold|replay|compensate");
if(args[0].equals("compensate")){compensate();return;}
check(List.of("retry","hold","replay").contains(args[0]),"unknown mode");
var f=new ConnectionFactory();f.setHost("rabbit");f.setVirtualHost("lab");f.setUsername("app");
f.setPassword("app-lab-only");f.setAutomaticRecoveryEnabled(false);
try(var connection=f.newConnection();var ch=connection.createChannel()) {
ch.confirmSelect();
if(args[0].equals("replay")) {
String dead="me12.retry.replay.dead",work="me12.retry.replay.work";
for(String q:List.of(dead,work)){ch.queueDeclare(q,true,false,false,Map.of("x-queue-type","quorum"));ch.queuePurge(q);}
try(var c=database();var s=c.createStatement()) {
s.execute("CREATE TABLE IF NOT EXISTS retry_replay_audit(event_id text PRIMARY KEY,reason text NOT NULL,state text NOT NULL)");
s.execute("INSERT INTO retry_replay_audit VALUES('repair-42','verified order amount=10','APPROVED') ON CONFLICT(event_id) DO UPDATE SET state='APPROVED'");
}
publish(ch,dead,"repair-42","amount=invalid",3);
var original=receive(ch,dead,5);
check(new String(original.getBody(),StandardCharsets.UTF_8).equals("amount=invalid"),"unexpected repair input");
publish(ch,work,original.getProps().getMessageId(),"amount=10",1);
ch.basicAck(original.getEnvelope().getDeliveryTag(),false);
var repaired=receive(ch,work,5);
check("repair-42".equals(repaired.getProps().getMessageId())&&new String(repaired.getBody(),StandardCharsets.UTF_8).equals("amount=10"),"repair changed ID or wrong payload");
try(var c=database();var s=c.createStatement()) {
check(s.executeUpdate("UPDATE retry_replay_audit SET state='VERIFIED' WHERE event_id='repair-42' AND state='APPROVED'")==1,"approval absent");
}
ch.basicAck(repaired.getEnvelope().getDeliveryTag(),false);
check(ch.basicGet(dead,true)==null&&ch.basicGet(work,true)==null,"replay queues not drained");
System.out.println("replay: approved=1 stableEventId=true repairedPayload=true sourceRemovedAfterConfirm=true verified=1");
return;
}
if(args[0].equals("hold")) {
String source="me12.retry.hold",target="me12.retry.hold.target",exchange="me12.retry.hold.dlx";
ch.exchangeDeclare(exchange,"direct",true);
ch.queueDeclare(source,true,false,false,Map.of("x-queue-type","quorum"));
ch.queueDeclare(target,true,false,false,Map.of("x-queue-type","quorum"));
ch.queuePurge(source);ch.queuePurge(target);
ch.queueUnbind(target,exchange,"failed"); // The source's policy routes here after rejection.
publish(ch,source,"held-event","payload",1);
var original=receive(ch,source,5);ch.basicReject(original.getEnvelope().getDeliveryTag(),false);
long end=System.nanoTime()+TimeUnit.SECONDS.toNanos(2);
while(System.nanoTime()<end){check(ch.basicGet(target,true)==null,"unbound target received message");Thread.sleep(50);}
System.out.println("hold: targetUnbound=true targetEmpty=true");
ch.queueBind(target,exchange,"failed");
var recovered=receive(ch,target,240);
check("held-event".equals(recovered.getProps().getMessageId()),"wrong dead-lettered event");
check("payload".equals(new String(recovered.getBody(),StandardCharsets.UTF_8)),"body changed");
ch.basicAck(recovered.getEnvelope().getDeliveryTag(),false);
System.out.println("hold: bindingRestored=true recoveredOriginalEvent=true noApplicationRepublish=true");
return;
}
ch.exchangeDeclare("me12.retry.dead","direct",true);
ch.queueDeclare(DLQ,true,false,false,Map.of("x-queue-type","quorum"));
ch.queueBind(DLQ,"me12.retry.dead","failed");
ch.queueDeclare(WORK,true,false,false,Map.of("x-queue-type","quorum","x-dead-letter-exchange","me12.retry.dead","x-dead-letter-routing-key","failed"));
ch.queueDeclare(DELAY,true,false,false,Map.of("x-queue-type","quorum","x-message-ttl",500,"x-dead-letter-exchange","","x-dead-letter-routing-key",WORK));
ch.queuePurge(WORK);ch.queuePurge(DELAY);ch.queuePurge(DLQ);
try(var c=database();var s=c.createStatement()){
s.execute("CREATE TABLE IF NOT EXISTS retry_effect(event_id text PRIMARY KEY)");
s.execute("DELETE FROM retry_effect WHERE event_id IN ('healthy','transient','poison')");
}
publish(ch,WORK,"healthy","healthy",1);publish(ch,WORK,"transient","transient",1);publish(ch,WORK,"poison","poison",1);
var attempts=new ArrayList<Integer>();int success=0,poison=0;
long deadline=System.nanoTime()+TimeUnit.SECONDS.toNanos(20);
while((success<2||poison<1)&&System.nanoTime()<deadline) {
var message=receive(ch,WORK,5);String id=message.getProps().getMessageId();
int attempt=((Number)message.getProps().getHeaders().get("attempt")).intValue();
long tag=message.getEnvelope().getDeliveryTag();
if(id.equals("poison")){ch.basicReject(tag,false);poison++;continue;}
if(id.equals("transient")){
attempts.add(attempt);
if(attempt<3){publish(ch,DELAY,id,"transient",attempt+1);ch.basicAck(tag,false);continue;}
check(attempt==3,"retry budget exceeded");
}
successful(id);ch.basicAck(tag,false);success++;
}
check(success==2&&poison==1&&attempts.equals(List.of(1,2,3)),"retry path differs");
var dead=receive(ch,DLQ,5);check("poison".equals(dead.getProps().getMessageId()),"wrong dead letter");
var death=(List<?>)dead.getProps().getHeaders().get("x-death");
check(death!=null&&death.stream().anyMatch(x->"rejected".equals(((Map<?,?>)x).get("reason").toString())),"rejection history missing");
ch.basicAck(dead.getEnvelope().getDeliveryTag(),false);
try(var c=database();var s=c.createStatement();var rows=s.executeQuery("SELECT count(*) FROM retry_effect")){
check(rows.next()&&rows.getInt(1)==2,"unexpected business effects");
}
System.out.println("retry: healthyAttempts=1 transientAttempts="+attempts+" businessRows=2 poisonDeadLettered=true");
}
}
static void compensate() throws Exception {
try(var c=database();var s=c.createStatement()) {
s.execute("CREATE TABLE IF NOT EXISTS retry_reservation(order_id text PRIMARY KEY,state text NOT NULL)");
s.execute("CREATE TABLE IF NOT EXISTS retry_stock(sku text PRIMARY KEY,available int NOT NULL)");
s.execute("INSERT INTO retry_stock VALUES('book',9) ON CONFLICT(sku) DO UPDATE SET available=9");
s.execute("INSERT INTO retry_reservation VALUES('order-42','RESERVED') ON CONFLICT(order_id) DO UPDATE SET state='RESERVED'");
}
int applied=0;
for(int attempt=0;attempt<2;attempt++) {
try(var c=database();var s=c.createStatement()) {
c.setAutoCommit(false);
try {
int changed=s.executeUpdate("UPDATE retry_reservation SET state='RELEASED' WHERE order_id='order-42' AND state='RESERVED'");
if(changed==1){check(s.executeUpdate("UPDATE retry_stock SET available=available+1 WHERE sku='book'")==1,"stock missing");applied++;}
c.commit();
}catch (Exception | Error failure) {
try { c.rollback(); }
catch (java.sql.SQLException rollbackFailure) { failure.addSuppressed(rollbackFailure); }
throw failure;
}
}
}
try(var c=database();var s=c.createStatement();var r=s.executeQuery("SELECT available FROM retry_stock WHERE sku='book'")){
check(r.next()&&r.getInt(1)==10&&applied==1,"compensation repeated effect");
}
System.out.println("compensation: attempts=2 applied=1 reservation=RELEASED available=10");
}
}publish 等待 Broker confirm 后才确认原消息,因此进程若在这两个动作之间退出,会留下“延迟队列已有副本、原工作消息重新交付”的重复窗口。应用应继续使用同一个 event ID,并在业务事务中去重;交换两步会产生已确认原消息但重试副本尚未保存的丢失窗口。发布确认与消费确认
此处的发布目标均由程序先声明,实验期间不删除拓扑。实际系统还应处理 mandatory return:publisher confirm 为 ACK 时,消息也可能没有路由到任何队列。声明、路由与 confirm 的正反实验见 RabbitMQ 运行链。
retry_effect 的唯一键只验证重复结果记录的抑制;涉及余额等实际业务变化时,去重记录必须与业务修改一起提交,完整事务路径见 投递、幂等与顺序。
死信保留、目标恢复与受控重放
DLX 是交换机,DLQ 是最终接收队列
消息被拒绝且不重新入队、达到消息 TTL、达到某些长度限制或超过 quorum 投递限制时,可以按配置进行 dead lettering。DLX 负责重新路由,DLQ 才存放最终消息;只配置一个交换机名字而没有匹配绑定,不能形成有效的死信接收路径。Dead Letter Exchanges
x-death 保存消息在死信路径上的队列、原因和次数等信息,它适合诊断历史。应用自行重新 publish 的新消息,是否保留这些 header 取决于具体代码;不要把业务 attempt、Broker 的 acquired count 与 x-death count 当成同一个计数器。
默认死信路径与至少一次转交
前面的 TTL 和毒消息例子使用默认死信策略。默认转交并不提供目标不可用时的保留保证。对必须保存的失败业务消息,可在源 quorum queue 上启用 at-least-once dead lettering:内部消费者保留源消息,转发到目标队列并等待 publisher confirms,再移除源端记录。
以下管理命令只作用于精确命名的实验队列:
docker compose exec -T rabbit rabbitmqctl set_policy -p lab me12-hold \
'^me12[.]retry[.]hold$' \
'{"dead-letter-exchange":"me12.retry.hold.dlx","dead-letter-routing-key":"failed","dead-letter-strategy":"at-least-once","overflow":"reject-publish","max-length":100}' \
--apply-to quorum_queues
run_retry hold必须同时使用 overflow=reject-publish。源队列达到容量限制时会拒绝新发布,避免待转交消息无限积累。at-least-once 有额外资源成本,也允许目标收到重复;它不使消费者业务变成一次执行。
hold 先解绑目标队列,发布一条消息到源队列,再拒绝这次交付。前两秒确认目标没有收到消息,然后恢复绑定,等待原事件自动转交:
hold: targetUnbound=true targetEmpty=true
hold: bindingRestored=true recoveredOriginalEvent=true noApplicationRepublish=true两行之间可能等待约三分钟。RabbitMQ 4.3.5 的内部工作进程使用 publisher confirm 超时驱动周期重投,绑定恢复不直接等于立即交付。代码为恢复阶段保留 240 秒窗口,期间应用没有再次 publish。具体调度机制可查 4.3.5 死信转发工作进程。
等待期间可从另一终端查询:
docker compose exec -T rabbit rabbitmqctl list_queues -p lab \
name messages messages_ready messages_unacknowledged policy被内部死信消费者保留的记录,可能表现为源队列 messages=1,而 ready=0、unacknowledged=0;不要只用面向正常消费者的两个数量相加推断源已完全清空。恢复成功要以目标实际收到相同 ID 和正文为准。
若 240 秒仍未恢复,检查策略是否命中源队列、绑定键是否为 failed、目标是否拒绝发布和 Broker 日志中的 dead-letter 警告。尚有待转交消息时切换回 at-most-once、删去 DLX 或采用 drop-head,可能使这些保留消息被删除,不能用这些动作消除告警。
重放需要知道修了什么
run_retry replay此模式准备一条内容为 amount=invalid 的隔离消息,写入明确的批准记录,按已核对的实验金额修复为 amount=10。发布修复消息并获得 confirm 后才 ACK 原死信,接收端检查 ID 保持 repair-42、正文已修复,再把审计状态从 APPROVED 改为 VERIFIED。
replay: approved=1 stableEventId=true repairedPayload=true sourceRemovedAfterConfirm=true verified=1这里的金额由实验预设,真实系统必须从可信业务记录核对;重放工具不应猜测支付金额、身份或权限。程序保留相同事件 ID 是为了让重复修复仍可去重。如果业务上确实发生了另一件事,应生成新的事件,并用 causationId 或补偿引用关联原事件,而非修改历史事实后假装它从未存在。
批准表、Broker publish 与原死信 ACK 没有共同事务。重放程序仍可能在 confirm 后崩溃,因此生产消费者要保留幂等处理;审计记录也应包含审批人、失败原因、修复版本、批次、速率和最终业务结果。先重放一个样本,确认结果,再逐步扩大批次。
定时触发与补偿怎样进入业务状态
到期消息只是一次处理机会
“订单创建三十分钟后检查是否未支付”可以通过延迟消息触发。消费者收到消息时重新读取订单:已支付则结束;仍待支付且确实超过截止时间时,用带状态条件的 UPDATE 关闭订单。队列交付迟到、重复交付和人工操作都需要容纳在状态判断里。
产品提供的延迟粒度、最大时长、存储恢复方式与吞吐限制各有不同。RocketMQ 将定时/延迟消息作为独立消息类型,并规定相应的投递时间行为;不能直接把 RabbitMQ 的 TTL 参数搬到 RocketMQ。RocketMQ 定时消息
如果到期处理是业务强要求,还需要持久的待处理任务或按截止时间扫描的补漏任务。仅在 JVM 中使用定时器,进程退出后未执行任务会一起消失。长周期任务通常适合保存在数据库中,接近到期时再投入消息队列,数据库记录用于恢复和查询。
补偿是新的受条件约束的业务操作
补偿释放库存时,应只释放曾成功预占且尚未释放的那一笔。下列实验初始可用库存为 9,订单预占状态为 RESERVED;连续执行两次释放:
run_retry compensatecompensation: attempts=2 applied=1 reservation=RELEASED available=10事务中的第一条 UPDATE 使用 WHERE state='RESERVED' 争取释放资格,只有影响一行时才增加库存;两条修改共同 commit。第二次执行已经找不到符合条件的预占,因此不再增加库存。PostgreSQL 事务
退款、撤销预占、关闭订单分别有业务限制。已经发货的订单可能不能简单退回初始状态,退款又可能失败或结果未知。应记录补偿操作 ID、当前状态和外部结果,允许重试查询及人工处理;不要把任意 SQL 的反向操作称为安全补偿。
运行中该观察哪些变化
| 现象 | 判断与后续操作 |
|---|---|
| 重试率升高但成功率不变 | 按错误类型统计;永久失败退出重试,暂态失败检查下游恢复 |
| 工作队列持续流动,订单长时间未完成 | 用 event ID 查尝试历史、deadline 和最终业务状态,不能只看 ready 数量 |
| DLQ 新增突然增多 | 对照最近契约或代码变更,暂停批量重放,先检查样本 |
| 死信源队列占用空间却没有 ready 消息 | 检查内部转交是否在等目标确认,修复 DLX 路由或目标容量 |
| 补偿操作多次执行 | 查询条件更新影响行数和外部操作 ID,确认是重复调用还是重复生效 |
| 到期任务整体延后 | 比较计划时刻、实际交付、开始处理与业务完成时间,定位排队或下游瓶颈 |
结束时 docker compose stop 保留队列和数据库。仅当确认 me12-lab 中全部数据均为可丢弃实验数据时执行 docker compose down --volumes;它会删除本项目消息与数据库卷。若继续保留项目而只停止特殊策略,在 hold 完成且没有待转交消息后执行:
docker compose exec -T rabbit rabbitmqctl clear_policy -p lab me12-hold权威资料与规范地址
| 查阅内容 | 官方地址 |
|---|---|
| 重投、结果未知与可靠性 | https://www.rabbitmq.com/docs/reliability |
| Quorum 投递限制、延迟重试、至少一次死信 | https://www.rabbitmq.com/docs/quorum-queues |
| 消息与队列 TTL | https://www.rabbitmq.com/docs/ttl |
| 发布确认和消费确认 | https://www.rabbitmq.com/docs/confirms |
| 死信原因、交换机与历史字段 | https://www.rabbitmq.com/docs/dlx |
| 周期性死信恢复实现 | https://raw.githubusercontent.com/rabbitmq/rabbitmq-server/v4.3.5/deps/rabbit/src/rabbit_fifo_dlx_worker.erl |
| RocketMQ 定时与延迟消息 | https://rocketmq.apache.org/docs/featureBehavior/02delaymessage/ |
| 数据库事务提交与回滚 | https://www.postgresql.org/docs/18/tutorial-transactions.html |
