任务队列
Agent 为什么需要队列
普通 HTTP 请求的生命周期很短:接收请求、做一点同步检查、返回响应。Agent 运行不是这样。
一次 Agent run 可能会持续几十秒甚至更久。它要加载上下文、连接 MCP、调用模型、执行工具、写消息、扣费、抽取记忆、索引资源,还要处理用户中途 Stop、客户端断线和服务端重启。把这些都放进同一个请求处理函数里,会让 API 进程同时承担三件互相冲突的事:
- 对用户尽快返回响应:请求路径不能被长时间 LLM / 工具调用卡住。
- 稳定执行长任务:Agent run 不能因为 HTTP 连接断开就丢。
- 处理轮末副作用:扣费、记忆抽取、资源索引不能阻塞流式响应收尾。
所以我们的 Agent 系统需要队列。队列不是为了“让代码异步一点”,而是为了把请求接入、Agent 执行、轮末副作用拆成不同生命周期。
队列在系统里的位置
execute(taskId, user, body) 不直接跑 Agent。它只做请求路径上必须同步完成的事:
- 清理旧 abort flag。
- 做归属、额度等 preflight。
- 保存用户本轮消息。
- 重置本轮 stream buffer。
- 把
task.run入队。 - 读取 buffer,把流式响应返回给客户端。
真正的 Agent run 在 worker 里执行。run 产生的 UIMessageChunk 不直接写给某个 HTTP 连接,而是 produce 到按 taskId 键的 stream buffer;HTTP 响应、断线后的 resume,都只是这个 buffer 的读者。
这层拆分很关键:
HTTP 请求
-> preflight + 入队
-> 读 stream buffer
Worker
-> 执行 task.run
-> produce UIMessageChunk 到 stream buffer
-> finalize 后入队轮末 job
API 进程负责接入和读回,worker 负责执行。客户端连接断了,run 仍然可以继续;客户端回来时,仍然读同一份 buffer。
两类任务都需要队列
Agent 系统里有两类工作会活过发起它的请求。
第一类是 Agent run 本身。task.run 是长时、IO-bound、高并发的任务:它会等待模型、等待工具、等待外部系统。它不能长期占着 API 请求线程,也不能因为一次连接断开就消失。
第二类是 轮末 job。一轮流 drain 完后,finalizeExecution 还会触发扣费、记忆抽取、资源索引等后台工作。这些工作重要,但不应该挡住用户看到本轮回答的完成状态。
因此队列承担两个角色:
- 承载长运行的 Agent 执行:让 run 从请求路径中移出,并能跨进程执行。
- 承载轮末副作用:让扣费、记忆、索引等工作独立重试、独立监控、独立限流。
如果不用队列,常见替代方案是 void promise 或 setTimeout。它们的问题很直接:进程退出任务就丢,没有重试,没有并发控制,也没有 dashboard 告诉你任务卡在哪里。
队列带来的边界
队列解决了可恢复执行,但它不会给我们 exactly-once。
worker 执行任务时,通常要先写业务副作用,再向队列 ACK 完成。比如扣费 job 要先写 ledger / balance,再告诉 Redis 这个 job 完成。业务数据库和队列存储不是同一个事务,中间一定存在崩溃窗口:
副作用已经写入数据库
ACK 还没写回队列
worker 崩溃
队列恢复后重新投递这个 job
所以队列的诚实契约是 at-least-once:任务至少会被投递一次,也可能被重新投递。
这不是 BullMQ 的缺陷,而是后台任务系统的基本边界。研发写任何 job 都要默认它可能重跑。工程目标不是“消灭重复投递”,而是让重复投递不会造成重复生效:
at-least-once 投递 + 业务幂等 = 效果上只生效一次。
在我们的系统里,最典型的是扣费:credit.consume 使用稳定的 reference id,并由 ledger 表唯一约束兜底。重复投递时,第二次写入撞唯一键,余额不会被重复扣减。
写 Agent 后台任务的规则
第一,payload 只传稳定 ID,不传运行时对象。
BullMQ 路径会把 payload 序列化进 Redis,worker 再从 payload 重建依赖。闭包、sandbox 句柄、数据库连接、class 实例都不能跨进程序列化。正确的 payload 应该像这样:
jobQueue.enqueue("resource.index", {
taskId,
resourceId,
version,
});
worker 执行时再用这些 ID 读取权威数据。这样既能拿到最新状态,也能在失败后安全重放。
第二,区分 jobId 和业务幂等。
jobId 只能防止同一个提交重复入队,比如 credit:{taskId}:{messageId}。但 job 一旦开始执行,worker 在 ACK 前崩溃,队列仍可能重新投递。真正防止重复扣费、重复写副作用的,是业务数据库里的唯一约束和事务。
第三,副作用要有幂等键。
扣费用 taskId + messageId;资源索引用 resourceId + version;记忆抽取要能覆盖或跳过同一轮结果。凡是会改变外部状态、余额、索引、通知的 job,都要先回答:这个 job 重跑一次会怎样?
第四,不要假设顺序。
队列有并发、重试和 delayed 状态,入队顺序不等于完成顺序。如果某个任务只允许同一实体内有序,就按 taskId、resourceId 或其它业务 key 做串行、版本检查或覆盖写,而不是默认全局 FIFO。
为什么要拆不同队列
Agent run、扣费、记忆抽取、资源索引的负载形态不同,不能挤在同一条执行通道里。
agent:长时、IO-bound,高并发,主要在等模型和工具。critical:扣费等关键副作用,要快速消费,不能被长 run 饿死。llm:记忆抽取等 LLM job,要控制 provider 限流和 token 成本。indexing:资源索引可能更偏 CPU / DB 密集,需要独立伸缩。scheduling/nudge:时间触发任务,和用户即时 run 解耦。
拆队列的目的不是让拓扑好看,而是避免一种负载饿死另一种负载。大量 Agent run 堆积时,扣费仍要能消费;资源索引变慢时,不能拖垮 Agent 执行。
内联队列为什么还存在
同一个 JobQueue 接口有两条实现:
| 场景 | 实现 | 含义 |
|---|---|---|
| Server + Redis | BullMQ | payload 序列化进 Redis,由 worker 执行 |
| Desktop / 无 Redis dev | inline | fire-and-forget,在当前进程执行 |
这不是两套业务逻辑,而是同一个调用点的两种执行方式。
调用方同时提供 payload 和 execute 闭包:BullMQ 使用 payload,worker 从 ID 重建执行;inline 使用闭包,直接在本进程跑。这样 Server 有跨进程恢复能力,Desktop 不需要 Redis 也能复用同一套编排形状。
这个设计也反过来约束我们:即使 inline 路径能拿到闭包,job 的语义仍必须按 BullMQ 路径设计。也就是 payload 要可序列化,副作用要幂等,失败要能重试。
队列必须可观测
队列把工作移出了请求路径,也意味着失败不再天然暴露给用户请求。没有 dashboard 和告警,后台任务就会变成黑盒。
研发至少要能看到:
- 每个队列的
waiting、active、delayed、failed数量。 - 哪些 job 失败了,payload 和错误栈是什么。
- job 从入队到开始执行、到完成的延迟。
- 是否出现 stalled,是否有任务长时间停在
active。 - worker 是否存活,并发是否被占满。
- 重试耗尽后是否触发 Webhook / 告警。
在我们的系统里,对应 Bull-board、队列指标 API、失败 Webhook 和 worker 结构化日志。队列上线后最危险的不是“任务失败”,而是“任务失败了很久没人知道”。
写新 Agent job 前的自检
- 这个工作为什么不能留在请求路径或 run 的 inline finalize 里?
- 它是 Agent run 本身,还是轮末副作用?
- payload 是否只包含稳定 ID、版本、幂等键和必要快照?
- worker 能否只靠 payload 重建所需依赖?
- job 被重新投递时,副作用是否会重复生效?
jobId是否只被当作提交去重,而不是业务幂等?- 幂等记录和业务副作用是否在同一个事务里提交?
- 这个 job 失败后应该重试吗?哪些错误应该直接停止重试?
- 重试耗尽后在哪里可见?谁会收到告警?怎么重放?
- 它是否会被
agent队列饿死,或反过来饿死关键任务?
详细 BullMQ 接线、Worker 进程模型和文件清单见 → 后台任务队列架构。