六维教程

死信队列与重试策略

消息处理不可能永远成功。下游数据库偶尔抖动、第三方 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 兜底永久失败。延迟投递是另一条线,解决的是定时和延后需求,和重试退避不要混淆。关键业务一定配死信队列,别让消息默默丢失。

上一篇 批处理与并发控制

上一篇
批处理与并发控制
下一篇
对象存储基础与 Bucket 管理