死信队列与重试策略
消息处理不可能永远成功。下游数据库偶尔抖动、第三方 API 临时超时、消息体本身格式错误,这些都会让消费失败。失败后怎么办,是重试还是丢弃,重试几次,重试间隔多久,这些问题不解决好,要么丢数据,要么无限重试拖垮系统。这篇讲清楚 Cloudflare Queues 的重试机制、死信队列以及延迟投递。
消息失败的两种原因
先分清失败类型,不同类型策略完全不同
| 失败类型 | 特征 | 例子 | 应对策略 |
|---|---|---|---|
| 瞬时失败 | 重试可能成功,下游本身没问题 | 网络抖动、数据库短暂连不上 | 自动重试 |
| 永久失败 | 重试多少次都失败 | 消息格式错误、引用不存在的资源 | 进死信队列人工排查 |
区分这两类是重试策略的前提。瞬时失败多试几次就好,永久失败重试只会浪费资源,必须挪走。
ack retry 和显式失败
消费者处理消息后有三种收尾方式
| 方法 | 作用 | 消息去向 |
|---|---|---|
| msg.ack() | 标记成功 | 从队列删除 |
| msg.retry() | 主动要求重试 | 重新投递 |
| 不调用任何方法 | 默认视为失败 | 按重试策略处理 |
代码里显式控制
export default {
async queue(batch: MessageBatch, env: Env): Promise<void> {
for (const msg of batch.messages) {
try {
const body = msg.body as { email: string };
if (!body.email) {
// 永久失败,格式错误,重试无用
// 不 ack 也不 retry,让它走死信流程
console.error(`消息格式错误 ${JSON.stringify(body)}`);
continue;
}
await sendMail(body.email);
msg.ack();
} catch (e) {
// 瞬时失败,请求重试
msg.retry();
}
}
},
};
retry 不带参数表示按默认退避重试,也可以传选项控制重试行为,后面会讲。
重试次数与退避
队列绑定里有个关键参数 max_retries,控制一条消息最多重试几次
[[queues.consumers]]
queue = "demo-queue"
max_batch_size = 10
max_retries = 3
dead_letter_queue = "demo-dlq"
max_retries 默认 3。超过次数后消息不再重试,进入死信队列。
重试间隔不是立刻就试,Queues 会自动做指数退避,避免下游刚恢复就被打爆
| 重试轮次 | 大致间隔 | 说明 |
|---|---|---|
| 第 1 次 | 几秒 | 快速重试,应对瞬时抖动 |
| 第 2 次 | 几十秒 | 拉长间隔 |
| 第 3 次 | 几分钟 | 等下游充分恢复 |
具体间隔由 Queues 内部调度,不可手动指定。如果需要自定义延迟,用后面讲的延迟投递。
msg.retry 可以传选项指定额外延迟
msg.retry({ delaySeconds: 60 });
这条消息会在 60 秒后才重新投递,适合下游明确需要冷却时间的场景,比如限流恢复。
死信队列
死信队列(Dead Letter Queue,DLQ)是存放重试耗尽消息的特殊队列。配置 dead_letter_queue 参数后,超过 max_retries 的消息会被原样投递到这里。
[[queues.consumers]]
queue = "demo-queue"
max_retries = 3
dead_letter_queue = "demo-dlq"
demo-dlq 要先用 wrangler(Cloudflare 的命令行工具)创建好
npx wrangler queues create demo-dlq
死信队列本身也是个普通队列,可以再绑一个消费者去处理。常见做法是绑一个发告警的 Worker
export interface Env {
ALERT: Queue<unknown>;
}
export default {
async queue(batch: MessageBatch, env: Env): Promise<void> {
for (const msg of batch.messages) {
// 死信消息带原始内容和失败次数
await env.ALERT.send({
text: `死信告警 消息 ${JSON.stringify(msg.body)} 重试 ${msg.attempts} 次仍失败`,
});
msg.ack();
}
},
};
msg.attempts 是这条消息已经被尝试的次数,可以用来判断严重程度。
死信队列要不要配
| 场景 | 是否配 DLQ | 原因 |
|---|---|---|
| 关键业务数据 | 必须配 | 不能丢,要人工兜底 |
| 日志统计 | 可不配 | 丢几条不影响大局 |
| 推送通知 | 建议配 | 失败的存档,方便补发 |
不配 DLQ 时,重试耗尽的消息直接被丢弃,不可恢复。
延迟投递
有时候不是消息处理失败,而是业务上要求消息延后处理。比如下单后 30 分钟未支付自动取消,注册后 24 小时发提醒邮件。这种场景用延迟投递。
生产者 send 时可以指定延迟
// 30 分钟后投递
await env.DEMO_QUEUE.send(
{ orderId: 123, action: "check_payment" },
{ delaySeconds: 30 * 60 }
);
delaySeconds 最大 12 小时,最小 1 秒。延迟消息会在指定时间后才进队列被消费。
延迟投递和重试退避的区别
| 机制 | 触发方 | 用途 | 间隔控制 |
|---|---|---|---|
| 延迟投递 | 生产者 | 定时任务、延后处理 | 精确指定 |
| 重试退避 | 消费者 | 失败后自动恢复 | Queues 内部 |
实际业务里两者经常配合。比如支付检查消息,生产者用延迟投递定在 30 分钟后,消费者处理时如果支付还没完成,用 msg.retry 加 delaySeconds 再延 10 分钟,形成轮询直到支付完成或超时。
完整失败处理流程
把前面几节串起来,一条消息从进入到结束的完整路径
投递到队列
-> 消费者处理
-> 成功 ack 删除
-> 失败 retry
-> 退避后重投
-> 重试次数 < max_retries 继续处理
-> 重试次数 >= max_retries 进死信队列
-> 死信消费者告警或补发
对应配置
[[queues.consumers]]
queue = "demo-queue"
max_batch_size = 10
max_batch_timeout = 5
max_concurrency = 2
max_retries = 3
dead_letter_queue = "demo-dlq"
这套配置覆盖了正常消费、失败重试、超限兜底三个层面,是生产环境的标准模板。
小结
失败处理三件套,retry 控制瞬时失败的重试,max_retries 控制重试上限,dead_letter_queue 兜底永久失败。延迟投递是另一条线,解决的是定时和延后需求,和重试退避不要混淆。关键业务一定配死信队列,别让消息默默丢失。
上一篇 批处理与并发控制