后台任务

了解 Reaction 异步命令、异步日志和访问统计三条队列的职责、重试方式与排查入口。

Saavo 使用 Cloudflare Queues 处理不适合阻塞当前请求的工作。项目内置三条用途明确的队列:业务事件产生的异步 Command、异步日志,以及第一方访问统计事件。

Reaction 并不是一条“万能异步队列”。Event 表示已经发生的业务事实,Command 表示系统接下来要执行的动作;每个 Command 自己决定同步执行还是异步执行。

警告

不要将自定义事件插入这些预定义的队列中,它们都有自己的功能和处理逻辑,如果想要业务上实现异步,可以尝试使用 Reaction 的 Event / Command 机制。又或者自己新建一条异步业务处理队列。

内置队列

队列Binding默认批次与并发用途
async-policy-taskASYNC_POLICY_TASK_QUEUE每批 1 条,并发 1执行 Reaction 异步 Command
async-loggerASYNC_LOGGER_QUEUE每批最多 50 条,并发 1把日志写入 KV
analytics-eventsANALYTICS_QUEUE每批最多 8 条,并发 1把第一方统计事件写入 ANALYTICS_DB

三条消费者的 max_retries 当前都是 3,平台配置允许首次投递后最多再投递 3 次,也就是最多 4 次消费尝试。但消费者可以提前确认并停止重试:Reaction 消费者按最多 4 次投递处理,统计和日志消费者在第 3 次处理失败时就会主动确认消息并记录错误或告警。

wrangler.jsonc 同时配置了队列名、Binding 和 *_QUEUE_NAME。入口根据实际队列名精确选择消费者;名称不一致时,消息不会自动改由其他队列处理。

Reaction 命令如何执行

业务代码通过 reaction.processor.emit() 发出 Event。处理器会先持久化 Event 和 Command 执行记录,再按顺序执行:

  • mode: 'sync':在当前请求、Webhook 或定时入口中立即尝试执行;
  • mode: 'async':把 Event 执行 ID 发送到 async-policy-task,由消费者继续执行。

当前同步 Command 主要用于授予或收回角色和权益;异步 Command 包括发送邮件、发送通知、创建站内通知、同步支付客户邮箱、处理到期权益周期和清理过期数据。

同步只说明“当前调用链会执行它”,不代表失败一定会变成异常。Command 返回的正常失败会写入执行记录,emit() 仍可能返回 accepted。如果调用方必须确保某项 Command 已成功,除了检查 Event 是否被接收,还要检查返回的 Command 状态,或读取后台执行记录。

管理后台提供了两个排查入口:

/dashboard/reaction/events
/dashboard/reaction/commands

重试、去重与失败处理

Event 使用稳定的 idempotencyKey 去重。相同 Event 类型和幂等键再次发出时,处理器返回 duplicate,不会创建第二套执行记录。

每个 Command 还有自己的 maxAttempts,默认通常为 3。它控制的是 Command 处理函数的业务重试;Cloudflare 的 max_retries 控制的是队列消息重新投递。两者不是同一层重试,排查时不要混为一谈。

队列采用至少一次投递,消息可能重复到达。Reaction 通过持久化状态、稳定执行 ID 和租约减少重复执行,但涉及发邮件、调用第三方接口等外部副作用时,Command 仍应使用稳定业务键实现幂等。

当前 async-policy-task 没有配置死信队列。队列消息达到最终投递仍失败时,消费者会确认消息并发出高优先级告警,后续需要人工检查执行记录。如果外部副作用已经完成、但结果状态无法持久化,消费者也会停止重投并告警,以免重复执行该副作用。

入队成功不等于任务完成

emit() 返回 accepted 表示 Event 和 Command 已被接受并开始分发;异步 Command 此时可能还在排队。返回 abandoned 则表示 Event 在多次发出尝试后仍未被可靠接受,调用方不能把它当作成功。

在业务代码中触发事件

业务侧应使用 Reaction,不要直接向 ASYNC_POLICY_TASK_QUEUE 发送自定义 JSON:

import { reaction } from '@/core/reaction';

const result = await reaction.processor.emit(
    workerCtx,
    SomeEvent.create({
        idempotencyKey: 'some-event:stable-business-id',
        payload: { /* 业务数据 */ },
    }),
);

emit() 可能返回以下结果:

结果含义
accepted新 Event 已持久化并开始执行或入队
duplicate相同幂等事件已经存在,返回原执行记录
ignoredEvent 没有生成任何 Command
abandoned发出过程最终失败,已记录日志和告警

调用方是否需要检查这些结果,取决于业务完成边界。注册授权、支付权益等关键流程不能只确认函数已经返回;需要确认相应同步 Command 的状态。纯通知类异步工作则可以在可靠入队后结束当前请求。

HTTP 路由使用 resolveFetchWorkerCtx(c) 获取上下文。Cron 和队列入口使用各自运行时已经创建的 WorkerCtx。应用代码不要直接调用 createWorkerCtx。

选择合适的处理方式

  • 业务事实及其后续动作:定义 Reaction Event / Command;
  • 页面访问和产品埋点:使用访问统计接口,由 analytics-events 处理;
  • 应用日志:使用现有 logger,由 async-logger 处理;
  • 与现有业务无关的长时间计算或通用作业:单独评估运行平台,不要直接塞进 async-policy-task。

config/deploy.ts 中的 logger.enableAsyncLoggerQueue 默认开启。关闭后,日志不再进入异步日志队列;消费者收到遗留消息时会确认并丢弃。访问统计整体关闭后,统计消费者也会确认并丢弃已经排队的事件。切换这些开关前,应先评估是否允许丢弃积压消息。

不要在没有测量数据的情况下提高批次或并发。async-policy-task 保持单条、单并发,是为了按执行记录的顺序和租约安全推进业务 Command。

上线检查

  • 三条队列的名称、Binding 和 *_QUEUE_NAME 完全一致。
  • 生产者和消费者已部署到同一套环境配置。
  • 注册、测试购买和订阅结束后,Reaction 执行记录符合预期。
  • 新 Event 使用稳定的业务幂等键。
  • 外部副作用可以安全重试,不依赖“队列只投递一次”。
  • 关键调用方会识别 abandoned,并按业务要求检查同步 Command 状态。
  • 日志或统计开关关闭前,已确认可以丢弃队列中的遗留消息。
  • 日志与告警渠道能收到最终消费失败和结果持久化失败通知。

常见问题

接下来