后台任务队列

任务队列把“宣布执行”和“实际执行”解耦 — 代价是每个任务可能被执行不止一次。at-least-once 是分布式投递唯一诚实的契约, 不是 bug。围绕这一事实,展开幂等、重试、死信、卡死检测、顺序与粒度的取舍。项目里的 BullMQ 实现见 agent/background-jobs。

解耦的代价:任务必然重复执行

队列是生产者—持久化—消费者三段:

  • Producer——入队即返回,不等执行,延迟毫秒级。
  • 持久化层——Redis 或数据库,任务存到被消费完为止。Redis 是内存存储,durability 取决于 AOF / RDB 配置;裸跑一次 flush 就把任务丢光。
  • Worker——独立进程拉取执行,并发可配。

状态机:waiting → active → completed / failed。失败且有重试预算转 delayed,退避后回 waiting;预算耗尽停在 failed

为什么必然重复,根因只有一个:worker 要先写副作用(你的 DB),再 ACK 完成(Redis)。这是横跨两个系统的 dual-write,凑不进同一个事务。两步之间存在一个崩溃窗口——副作用已落库、ACK 还没发出,进程就挂了;恢复后队列只知道这个任务没 ACK,于是重投。

所以“恰好一次”(exactly-once)投递不存在,唯一诚实的契约是 at-least-once:每个任务至少跑一次,可能更多。工程目标不是消灭重复,是让重复不产生后果——

at-least-once 投递 + 幂等处理 = effectively-once(效果上恰好一次)。

下面每一条,都是这个根因的某个侧面。

什么时候值得上队列

适用——

  • 持久性:进程崩溃任务不能丢。setTimeout / void promise / Promise.allSettled 随进程退出一起蒸发。
  • 延迟约束:计费、发邮件、压缩摘要这类耗时操作,不该卡住用户的 HTTP 请求。
  • 重试:第三方 API 偶发超时、provider 限流;裸 await 一次失败就丢。
  • 并发控制:上游有速率上限(provider 每秒 N 请求、SMTP 每分钟 M 封),不节流就被拒绝服务。
  • 可观测性:状态分布、堆积、失败率、延迟分位——没有仪表板就是盲跑。

不适用——

  • 结果必须同步可见:表单校验、SSR 路由所需数据。
  • 可丢且无副作用:缓存刷新,下次访问重做即可。
  • 任务只耗几毫秒:入队和拉取的开销超过任务本身。
  • 单进程 / Desktop:内联 fire-and-forget 就够,不必上 Redis 和独立 Worker(项目里对应 createInlineJobQueue())。

setTimeout 不是任务系统。 一个是进程内异步,一个是跨进程崩溃的持久化执行。判别只有一条:进程退出后,任务还在不在。

幂等:提交去重 ≠ 执行幂等

worker 在 ACK 前崩溃,任务回到 waiting 被另一个 worker 重跑。每个任务都要假设至少跑两次。 设计任务的第一个问题不是“怎么跑”,是“重复跑会怎样”:

  • 扣费重复 → 重复扣款 → 客户投诉(不可接受)
  • 邮件重复 → 用户收两封 → 体验下降(可补救)
  • 生成 PDF 重复 → 浪费 CPU、结果一致(可接受)

这里有两件不同的事,常被混为一谈:

提交去重——队列用 jobId 拒绝重复入队,防的是“同一请求被提交两次”。它的窗口有限(job 完成清理后 jobId 即失效),不解决崩溃重跑

执行幂等——对抗崩溃重跑的真兜底。核心是一个幂等键,加上一条铁律:幂等记录和副作用必须写进同一个 DB 事务

任务幂等键scope
扣费credit:{taskId}:{messageId}一条 message 只扣一次
发邮件email:{userId}:{templateId}:{eventId}同事件同模板一次
资源索引resource:{resourceId}:v{version}版本变了才重建

落地就是 DB unique constraint + ON CONFLICT DO NOTHING / upsert:第二次执行撞唯一键直接跳过。关键在于唯一键的写入和副作用同事务提交——否则“已生效、未记录”之间仍有窗口,等于没幂等。

payload 传引用,不传快照

// WRONG — 入队到执行之间数据可能变
jobQueue.enqueue("send-email", { user: currentUser, template: tmpl });

// OK — worker 执行时再从 DB 取最新
jobQueue.enqueue("send-email", { userId: currentUser.id, templateId: tmpl.id });
  • 时效:入队后用户可能改了邮箱、模板可能被禁用。传快照就是拿着入队那一刻的过时数据执行。
  • 可序列化:闭包、文件句柄、Date 实例、Class 实例跨 Redis 序列化时,要么报错,要么静默丢掉 prototype 方法。

例外:不可重构的纯文本快照(审计日志里的 userAgent 字符串)可以直接存——前提是你清楚那是快照,不是引用。

不是所有错误都值得重试

错误类型重试理由
Network / 5xx瞬时故障,退避后大概率成功
Rate limit(429)指数退避 + jitter
Validation / 4xx输入有问题,重试只会放大脏数据
Business rule reject业务拒绝,重试不会改变结果

不可重试的错误,立即转抛 UnrecoverableError(BullMQ 内置)跳过剩余重试,或在确认无副作用后 swallow 并正常 ACK。一律 throw,等于让脏数据耗尽重试预算——真正该重试的网络抖动反而抢不到机会。

重试预算(attempts)耗尽后,任务既不该消失,也不该堆在 failed 里等人肉打捞——见下一节。

重试耗尽之后:死信队列

退避重试到 attempts 上限仍失败的任务,叫毒丸消息(poison message):它不是瞬时故障,重投只会反复占用 worker。正确做法是把它移进独立的死信队列(DLQ)——隔离出主流程,保留完整 payload、错误栈、重试次数,供排查和手动重放

DLQ 要配两样东西:

  • 告警:任务进 DLQ 即触发(Webhook / Slack),不靠人轮询。
  • 重放入口:修掉根因后,能把 DLQ 里的任务重新投回原队列,而不是手敲数据库。

failed 当垃圾桶、靠人盯,是队列上线后最常见的运维黑洞。

卡死的任务:租约与 stalled

worker 拉到任务后崩溃、死循环、或被 OOM kill,任务会永远卡在 active——除非有机制发现它。机制是租约(lock):worker 必须周期性续约证明自己还活着;续约超时,任务被判 stalled,自动回到 waiting 重投(BullMQ 的 stalledInterval / maxStalledCount 即此)。

这是“为什么会重复执行”的另一半:不只是崩在 ACK 前,还包括假死被判 stalled 后重投。两条都回到根因,也都靠幂等兜底。

因此可观测性的第一个指标是:active 里有没有任务停留远超正常执行时长——那就是卡死的前兆。

入队顺序 ≠ 执行顺序

并发 worker 加上不同退避时长,入队序就不等于执行序:

入队:A, B, C
执行:B(快路径),A(失败重试一次),C

队列默认不保证顺序。真要顺序,按代价从低到高:

  • per-key 串行:按实体分区(per-user / per-resource 各一条队列),同 key 串行、不同 key 并行——既保顺序又不牺牲整体吞吐,首选。
  • 依赖编排:B 入队前先确认 A 已 completed(BullMQ Flows)。
  • 全局单 worker(并发 = 1)+ FIFO:顺序最强,但吞吐归零,只用于天然低频的串行流。

别假设“入队晚就一定执行晚”。

任务粒度:一个 job 一个业务单元

  • 过细——一个用户的 10 篇文章拆 10 个 job:编排开销超过节省的工作量。
  • 过粗——1000 篇文章塞进一个 job:失败时整体重做、无法部分 ACK、worker 内存陡增。

判别标准是失败时重做的代价能不能吞下。一个 job 对应一个有意义的业务单元(一次用户操作、一次 chat round 的扣费)。

没有 dashboard 就是黑盒

投产前先回答一个问题:worker 卡住时,从哪里看到? 答不上来就是黑盒。dashboard 至少要有:

  • 状态分布——waiting / active / delayed / failed / DLQ 各多少。
  • 卡死监控——active 停留超时(见 stalled)。
  • 失败详情——错误信息、payload、重试次数。
  • 告警——失败率或堆积超阈值 → Webhook / Slack。
  • 延迟分位——入队到完成的 P95 / P99。

项目实现

回到具体接线:

  • JobQueue 接口——同一 contract,Server 端 BullMQ + Redis,Desktop 端内联 fire-and-forget。
  • 三队列拓扑——critical(扣费,高并发快速消费) / llm(压缩摘要、记忆抽取,低并发控成本) / indexing(资源索引,独立伸缩)。
  • 提交去重——BullMQ jobId(如 credit:{taskId}:{messageId})拒绝重复入队。
  • 执行幂等——ledger 表 unique referenceId + ON CONFLICT DO NOTHING,与余额扣减同事务提交。
  • 可观测性——Bull-board 仪表板(/admin/queues)、指标 API、失败 Webhook。

详细架构、文件接线、Worker 进程模型见 → 后台任务队列架构

写新任务前的自检

  • 重复执行会怎样?幂等键就位了吗?幂等记录和副作用同事务吗?
  • payload 传的是 ID 还是对象快照?入队到执行之间这份数据会变吗?
  • 错误分类清楚了吗?可重试和不可重试在代码里分开了吗?
  • 重试耗尽进 DLQ 了吗?有告警和重放入口吗?
  • 依赖前一个任务完成吗?是显式编排,还是隐式假设了入队顺序?
  • 任务卡死时怎么发现?dashboard 和 stalled 监控就位了吗?
这页有帮助吗?