Queues:生产者、消费者、批量处理与重试实战

This article is extracted from the chat log with AI. Please identify it with caution.

说明:本文由 Codex 根据作者提供的主题、对话素材与 Cloudflare 官方文档辅助生成,属于 AI 生成/整理内容,非作者原创。请读者自行甄别并交叉验证。

Queues 解决的是“请求到来时先可靠地收下工作,稍后再异步处理”。它适合削峰、异步发信、写入第三方服务、批处理和事件驱动管道;它不是严格 FIFO 队列,也不是能持久化多步骤状态机的 Workflow。它提供的是 at-least-once(至少一次)投递:消息不会因一次暂时故障轻易丢失,但消费者必须能接受重复投递。

本文以“订单创建后异步履约”为例,覆盖 Worker producer、Worker consumer、显式确认、部分失败、延迟重试、死信队列(DLQ)、幂等与监控。内容按 2026-08-07 的官方文档整理;配额、价格和默认值会变化,部署前应复核文末来源。

先建立正确的心智模型#

flowchart LR
  A["HTTP API / Durable Object<br/>生产者 Worker"] -->|send / sendBatch| Q[("order-events")]
  Q -->|按批推送| C["消费者 Worker<br/>queue(batch, env, ctx)"]
  C -->|ack| S["D1 / R2 / 第三方服务"]
  C -->|retry + 延迟| Q
  C -->|超过 max_retries| D[("order-events-dlq")]
  D --> R["DLQ consumer<br/>告警、诊断、修复后重放"]

队列是生产者与消费者之间的持久缓冲。生产者的 send()sendBatch() resolve,表示消息已经写入磁盘;消费者在成功完成一个 batch 后,平台会默认确认该 batch。若 handler 抛错、返回的 Promise reject,或传给 ctx.waitUntil() 的 Promise reject,则整个未显式确认的 batch 会按 retry 配置重新投递。JavaScript API

Queues、Workflows 与 Durable Objects 怎么分工#

要解决的问题首选能力原因
高吞吐异步消息、削峰、容忍重复处理Queues按批投递、自动扩缩、至少一次语义。
一个订单/审批/Agent 流程要跨数小时或数天等待、恢复和记录步骤状态Workflowsstep.do()、sleep 与事件等待是持久化流程边界。
同一资源必须串行修改、需要按 key 协调或 WebSocket 状态Durable Objects每个对象 ID 有单一协调点和持久状态。

三者经常组合:Queue consumer 收到大量订单事件后为每条创建 Workflow;或先把同一订单的消息交给对应 Durable Object 做串行协调。不要用 max_concurrency: 1 来伪造 per-order 顺序:队列消息顺序仅是 best effort,且并发上限不能变成严格 FIFO。

1. 先配置 producer、consumer 与 DLQ#

下例将同一个 Worker 同时作为生产者和消费者。实际项目也可以拆成两个 Worker;绑定配置分别放在各自的 wrangler.jsonc 即可。

{
  "$schema": "./node_modules/wrangler/config-schema.json",
  "name": "order-events",
  "main": "src/index.ts",
  "compatibility_date": "2026-08-07",

  "queues": {
    "producers": [
      {
        "queue": "order-events",
        "binding": "ORDER_EVENTS"
      }
    ],
    "consumers": [
      {
        "queue": "order-events",
        "max_batch_size": 10,
        "max_batch_timeout": 5,
        "max_retries": 3,
        "dead_letter_queue": "order-events-dlq",
        "max_concurrency": 5
      }
    ]
  }
}

这里的数值是示例,不是通用最优配置:

  • max_batch_sizemax_batch_timeout 谁先达到,谁触发投递。更大的 batch 能减少 Worker invocation 并更适合批量写第三方服务,但低流量时会增加端到端延迟。
  • max_retries 定义同一消息失败后允许重投的次数;当前默认是 3。没有配置 dead_letter_queue 时,达到上限的消息会被永久丢弃。
  • max_concurrency 限制同时运行的 consumer invocation 数,常用于保护有速率限制的数据库或第三方 API;不设置时平台会自动扩缩到当前支持的最大并发。
  • DLQ 也是普通 Queue,需要再绑定一个 consumer 来告警、诊断或在修复后重放。没有活动 consumer 的 DLQ 消息当前会保留 4 天。
  • 一个 Queue 可以有多个 producer,但只能有一个活跃 consumer,且只能选择 Worker push 或 HTTP pull 其中一种消费类型。吞吐扩展来自同一 consumer 的并发 invocation,不是把同一队列绑定给多个独立消费者。

修改配置后运行 npx wrangler types,让 Queue<T>ExportedHandler 与 bindings 获得对应类型。生产者可以是普通 Worker,也可以从 Durable Object 内部发送消息;消费者还可选择 HTTP pull 模式,但本文聚焦 Worker 的 push consumer。配置 Queues 批处理与重试 死信队列

2. Producer API:把稳定的业务事件写入队列#

对可能产生副作用的事件,在生产端生成稳定的 eventId。下面直接要求调用方携带 idempotency key;若客户端重试同一请求,必须复用该 key,而不是每次重新生成 UUID。生产端还应把 key 与订单 ID、请求摘要关联起来,拒绝“同一个 key 却代表另一份订单”的冲突复用。

type OrderEvent = {
  eventId: string;
  kind: "order.created";
  orderId: string;
  occurredAt: string;
};

interface Env {
  ORDER_EVENTS: Queue<OrderEvent>;
}

export default {
  async fetch(request, env): Promise<Response> {
    if (request.method !== "POST") {
      return new Response("method not allowed", { status: 405 });
    }

    const idempotencyKey = request.headers.get("idempotency-key");
    const body = (await request.json()) as { orderId?: unknown };

    if (!idempotencyKey || typeof body.orderId !== "string") {
      return Response.json(
        { error: "idempotency-key and orderId are required" },
        { status: 400 },
      );
    }

    const event: OrderEvent = {
      eventId: idempotencyKey,
      kind: "order.created",
      orderId: body.orderId,
      occurredAt: new Date().toISOString(),
    };

    // Promise resolve 表示消息已持久化,而不表示消费者已经处理完成。
    const result = await env.ORDER_EVENTS.send(event);

    return Response.json(
      {
        accepted: true,
        eventId: event.eventId,
        backlogCount: result.metadata.metrics.backlogCount,
      },
      { status: 202 },
    );
  },
} satisfies ExportedHandler<Env>;

Queue<T> 的泛型只负责 TypeScript 开发期类型;它不会替你验证外部 JSON,因此 HTTP handler 和 consumer 仍应做运行时校验。消息 body 必须能通过 structured clone;单条不超过 128 KB。将大 payload 放到 R2、D1 或外部存储,消息只携带 object key、主键、版本或 checksum。

可以把 send() 放进 ctx.waitUntil() 来缩短 HTTP 响应等待,但这不适合“订单已创建即必须可靠投递”的关键事件:响应可能已经返回成功,而后台入队失败会被忽略。此类场景应 await 入队,或采用数据库 outbox 再由独立任务可靠投递。

需要一次发送多条时,用 sendBatch()

type PendingOrder = {
  eventId: string; // 来自订单记录 / outbox,而不是本次循环临时生成。
  orderId: string;
};

await env.ORDER_EVENTS.sendBatch(
  pendingOrders.map((order: PendingOrder) => ({
    body: {
      eventId: order.eventId,
      kind: "order.created" as const,
      orderId: order.orderId,
      occurredAt: new Date().toISOString(),
    },
  })),
  { delaySeconds: 60 },
);

sendBatch() 最多 100 条;每条仍受 128 KB 限制,整个 batch 最大 256 KB。delaySeconds 可用于发送后延迟可见性,单条消息可在 0 到 24 小时之间延迟。批量代码应保证每条消息都拥有独立、稳定的事件 ID,不能把一次 batch 的 ID 复用到每个业务事件上。Producer API 与限额 Queues limits

3. Consumer API:默认确认与显式确认的差别#

消费者通过默认导出的 queue(batch, env, ctx) 收到消息。若 handler 正常返回,并且所有 ctx.waitUntil() Promise 也正常 resolve,平台会自动确认 batch 中所有未显式处理的消息。

对象接口含义
MessageBatch<T>queuemessages当前队列名和本批消息;消息到达顺序不保证等于写入顺序。
MessageBatch<T>ackAll()标记整批成功。正常返回时通常不必显式调用。
MessageBatch<T>retryAll({ delaySeconds })标记整批以后重投。
Message<T>idtimestampbodyattempts平台消息 ID、发送时间、业务 body 与本次处理尝试次数;attempts 从 1 开始。
Message<T>ack()显式确认这一条;即使后续 batch 失败,它也不会再被投递。
Message<T>retry({ delaySeconds })只让这一条以后重投,适合局部失败和退避。

如果不显式 ack(),一个 batch 中第 8 条抛错会让整批未确认消息都重新投递。若每条都是幂等操作,这个最简单的“全成或全批重试”模型很好用;若一批包含多次第三方 API 调用或数据库写入,则应在每条成功后显式确认,避免已经成功的前 7 条反复执行。

特别注意不要写 batch.messages.forEach(async () => …)forEach 不会等待 async callback。使用 for…of,或者在每个操作都安全独立时 await Promise.all(…)Consumer API

4. 生产级 consumer:每条成功就确认,失败就延迟重试#

下面的 consumer 故意串行处理一个 batch,便于阅读和控制对下游的压力。若每个事件完全独立且下游允许并发,可把循环改为受限并发的 worker pool;不要无上限地 Promise.all 100 条外部请求。

type OrderEvent = {
  eventId: string;
  kind: "order.created";
  orderId: string;
  occurredAt: string;
};

interface Env {
  FULFILLMENT: Fetcher;
}

export default {
  async queue(batch, env): Promise<void> {
    for (const message of batch.messages) {
      try {
        await deliverWithIdempotencyKey(message.body, env);
        message.ack();
      } catch (error) {
        console.error("order event failed", {
          messageId: message.id,
          eventId: message.body.eventId,
          attempts: message.attempts,
          error: String(error),
        });

        // 返回成功,让已 ack 的消息保持成功;只有这一条稍后重试。
        message.retry({
          delaySeconds: retryDelaySeconds(message.attempts),
        });
      }
    }
  },
} satisfies ExportedHandler<Env, OrderEvent>;

async function deliverWithIdempotencyKey(event: OrderEvent, env: Env) {
  const response = await env.FULFILLMENT.fetch("https://internal/fulfill", {
    method: "POST",
    headers: {
      "content-type": "application/json",
      "idempotency-key": event.eventId,
    },
    body: JSON.stringify(event),
  });

  if (!response.ok) {
    throw new Error("fulfillment returned " + response.status);
  }
}

function retryDelaySeconds(attempts: number) {
  // 30, 60, 120 ...;上限远低于 Queues 所允许的 24 小时。
  return Math.min(30 * 2 ** Math.min(attempts - 1, 6), 3_600);
}

这个模式的关键不是 retryDelaySeconds(),而是外部服务必须按 eventId 去重:第一次请求已经完成、但 Worker 在 ack() 前失联时,重投必须仍然安全。支付、邮件、履约 API 若支持 idempotency key,直接使用事件 ID;自有服务则应以事件 ID 做唯一键、outbox key 或幂等记录。

上述代码将所有失败都交给重试与 DLQ,便于保留坏消息做调查。另一种合理策略是:对已经确认不可恢复的 schema 错误,先把原始消息与原因安全地写入“隔离”存储,再调用 ack()。不经保存就确认坏消息,等于主动丢数据。

5. 至少一次投递如何落到业务幂等#

Queues 默认保证“至少一次”,因此重复消息不是异常分支,而是正常设计条件。不要把 message.id 当作唯一业务幂等键:它是本次队列消息的系统 ID;应由生产者生成并随业务事件携带稳定的 eventId

生产端:为业务动作生成 eventId
    ├─ Queue 重投:仍使用同一个 eventId
    ├─ 下游数据库:eventId 作为唯一键 / inbox 记录
    └─ 下游 HTTP API:Idempotency-Key = eventId

一个足够可靠的原则是:先调用“以 eventId 幂等”的下游副作用,再记录本地成功状态;若在两者之间崩溃,下次重复调用仍会得到同一业务结果。若下游完全不支持幂等键,就需要自行设计 inbox/outbox、事务写入或由 Durable Object 串行协调;仅靠 Queue 的 ack() 不能把外部副作用变成 exactly-once。

对于同一订单必须严格按版本顺序执行的事件,也不要只依赖 Queue 到达顺序。应携带业务版本号,由 Durable Object / 数据库条件更新拒绝过期事件,或在消费端按 key 建立明确的协调层。Delivery guarantees

6. 重试、延迟和 DLQ 的处理策略#

重试不是“吞掉异常”。先决定错误属于哪类,再把策略写入代码与报警:

情况建议原因
网络超时、HTTP 429、短暂 5xxmessage.retry({ delaySeconds })为下游恢复和限流留出时间。
未知处理异常记录 eventId、attempts 与错误,再 retry让配置的重试上限和 DLQ 承接人工诊断。
已知无效 payload写入隔离存储后 ack(),或送入 DLQ不要无限重试同一坏数据,更不要无记录地丢弃。
batch 级依赖整体不可用抛错或 retryAll()让整批保持一致的失败边界。

单条 ack() / retry() 的第一次调用优先:之后对同一消息的确认或重试会被忽略;单条决定也优先于 ackAll() / retryAll()。因此应把每条消息的最终决定集中在一个代码路径里,避免不同 helper 互相覆盖。

DLQ 不会自动“修好”消息。应为它配置 consumer,并至少做到三件事:

  1. 产生告警:DLQ 出现消息通常意味着契约、权限或下游可用性有问题。
  2. 保留诊断信息:事件 ID、业务 key、attempts、错误类别和代码版本;日志不要包含支付信息、token 等敏感 payload。
  3. 修复后重放:修复消费者或数据后,以原始 eventId 重新发送到主队列。不能重建新 ID,否则会绕过幂等保护。

如果业务需要默认延迟或调整保留期,可通过 Wrangler CLI 更新 queue,例如 npx wrangler queues update order-events –delivery-delay-secs 60 –message-retention-period-secs 1209600。目前 message retention 可配置到 14 天;Free plan 的 retention 有不同限制,因此将其作为部署配置核查项,不要硬编码为应用假设。批处理、重试与延迟 配置 Queues

7. 监控:先看积压和 lag,再看 retry 与 DLQ#

在生产端,send()sendBatch() 的返回值以及 Queue.metrics() 都能给出实时的 backlogCountbacklogBytesoldestMessageTimestamp。这类数值适合请求后的轻量反馈,不应替代聚合监控。

Dashboard 与 GraphQL Analytics 可按队列查看以下趋势:

  • backlog messages / bytes 持续上升:消费者吞吐不足、被限流,或发生系统性失败。
  • lag time 增大:从写入到消费的端到端时间变长,应检查 batch、consumer concurrency 和下游性能。
  • retry count 突增:下游故障、schema 变更或代码发布回归的早期信号。
  • delete outcome 中的 dlqfail:应触发告警,不应只看 Worker 5xx。
  • consumer concurrency 长期顶在上限:若下游健康,考虑提高上限或分片;若下游受限,先降低生产速率或加大退避。

Queues 当前可观测 backlog、consumer concurrency 与消息操作;聚合指标通过 Dashboard / GraphQL 查询,实时 backlog 还可经 JavaScript API 或 REST API 获取。Queues metrics

8. 上线前检查单#

  • 每条业务事件有生产端生成的稳定 eventId,且下游以它去重。
  • 消费者不会把 forEach(async …) 当成已等待的处理。
  • 需要局部成功时,每条成功后 ack();失败后明确 retry() 或隔离并确认。
  • max_batch_size、timeout、retry 和 concurrency 是按延迟、下游配额与成本做出的选择,而非照搬默认值。
  • 主队列配置了 DLQ,并为 DLQ 配置了报警和处理人;没有无主的死信。
  • 消费者不把 Queue 顺序当成业务顺序;按 key 的顺序约束交给业务版本、Durable Object 或数据库。
  • 为 backlog、lag、retry 与 DLQ 建立告警,并在压测或故障演练中验证重放安全性。

参考资料#

本文共 4769 字,创建于 Aug 7, 2026

相关标签: Cloud, DevOps, TypeScript, ByAI