后台任务队列
任务队列把“宣布执行”和“实际执行”解耦 — 代价是每个任务可能被执行不止一次。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 监控就位了吗?