后台任务队列
基于 BullMQ 的后台任务处理——JobQueue 接口、六队列拓扑(含 agent 运行本身)、Worker 架构、双路径分发(BullMQ vs 内联)、 幂等扣费、可观测性仪表板
比一次请求活得久的,都交给队列
有两类工作会活过发起它的那个请求,于是都经同一个 JobQueue 接口、在请求路径之外执行——Server 端用 BullMQ(Redis 持久化、重试 + 并发控制),Desktop 端用内联 fire-and-forget:
- agent 运行本身——
execute()把task.run入队到agent队列;worker 跑它并 produce 进响应读回的那个 Redis stream buffer(见 任务编排)。 - 轮末 job——一轮的流 drain 之后(
finalizeExecution,在task-runner.ts),三个 fire-and-forget job:扣费(credit.consume)、记忆提取(memory.extraction)、资源索引(resource.index)。
BullMQ 工作原理
BullMQ 是基于 Redis 的 Node.js 任务队列。先看懂它的 Job 生命周期,后面的重试行为和故障模式才推理得清楚。
三个角色:
- Producer——API 服务器调用
Queue.add(name, payload)入队 Job。数据写入 Redis Stream 后立即返回。Producer 不执行 Job。 - Redis——存储所有 Job 状态。Job 存放在 Redis Streams(等待列表)和 Sorted Sets(延迟/优先级)中。Redis 持久化(AOF/RDB)确保 Job 在 Redis 重启后存活。
- Worker
——独立的 Node.js 进程,通过阻塞读取(
BRPOPLPUSH)拉取 Job。每个 Worker 有并发限制——每队列最多并行处理 N 个 Job。
Job 状态流转: waiting → active → completed 或 failed。失败后如仍有重试次数,Job 进入
delayed(指数或固定退避),然后回到 waiting。所有重试耗尽后,Job 永久留在 failed——在 bull-board 仪表板可见,并触发 Webhook 告警。
为什么不直接用 setTimeout 或 Promise?
setTimeout或void promise期间进程崩溃 = Job 永久丢失。BullMQ Job 持久化在 Redis 中。setTimeout无重试。BullMQ 支持可配置退避的重试。- 裸 Promise 无并发控制。BullMQ 按队列限制并行执行数。
- Fire-and-forget 无可观测性。BullMQ 提供状态、耗时、尝试次数和失败原因。
JobQueue 接口
遵循项目的 infra 模式(TaskLock、StreamBuffer、KeyEncryption),队列定义为
@zapvol/backend/src/infra/job-queue.ts 中的接口,有两个实现:
export interface JobQueue {
enqueue(
jobName: string,
payload: Record<string, unknown>,
execute: () => Promise<void>,
options?: JobEnqueueOptions,
): void;
}
核心设计:enqueue 同时接收可序列化 payload 和 execute 闭包:
| 实现 | payload | execute |
|---|---|---|
createBullMQJobQueue() | 序列化到 Redis | 忽略(Worker 从 payload 重建) |
createInlineJobQueue() | 忽略 | 直接调用(fire-and-forget) |
调用方提供全部信息,实现方按需选取。两条路径共用同一调用点(task-runner.ts 中的 finalizeExecution)。
注入
// Server — 有 Redis 时用 BullMQ,否则内联回退
const queues = getQueues();
const jobQueue = queues ? createBullMQJobQueue(queues) : createInlineJobQueue(log);
// Desktop — 始终内联
const jobQueue = createInlineJobQueue(log);
队列拓扑
六个命名队列(queues.ts 里的 QUEUE_NAMES),拆开是为了一种负载不会饿死另一种:
| 队列 | 任务 | 为什么单独隔离 |
|---|---|---|
agent | task.run / chat.run | agent 运行本身——长时、IO-bound(等 LLM/工具);高并发,与其它队列隔开,免得一波运行饿死扣费 / 索引 |
critical | credit.consume | 扣费必须快速消费——绝不被慢速 LLM 任务饿死 |
llm | memory.extraction | 低并发——控制 provider 限流与 token 成本 |
indexing | resource.index | CPU / DB 密集型——独立伸缩 |
scheduling | 定时任务触发 | wait_and_resume / cron 续跑 |
nudge | nudge 触发 | 时间触发的 nudge |
queues.ts 里的懒加载单例 Queue 实例由 Job 队列、bull-board 仪表板与指标端点共享;每个队列的 Worker 在 worker.ts 里各带自己的并发。
Worker 进程
单个 Worker 进程(apps/server/src/worker.ts)为每个队列各起一个 Worker,按 job name 经 JOB_PROCESSORS 分发。独立于 API 服务器运行:pnpm worker 或
pnpm worker:dev。
apps/server/src/
worker.ts → 入口,每个队列一个 Worker(经 JOB_PROCESSORS 分发)
jobs/
task-run.ts → task.run(agent 运行 → produce 进 stream buffer)
credit-consumer.ts → credit.consume
memory-worker.ts → memory.extraction
resource-indexer.ts → resource.index
schedule-fire-runner.ts → 定时任务触发
nudge-fire-runner.ts → nudge 触发
每个处理器在模块级别创建 Service 实例(与路由文件相同的内联组装模式),导出一个普通异步函数:
export async function processCreditConsume(job: Job<CreditConsumePayload>) {
const { userId, taskId, messageId, totalTokens } = job.data;
await creditService.consume(userId, taskId, messageId, totalTokens);
}
Payload 契约
Payload 只包含可序列化的 ID,不含运行时对象。Worker 从这些 ID 重建服务依赖(sandbox、model 工厂、repo)——闭包和 NodeSandbox 文件句柄无法序列化。
Worker 通过 taskService.loadTaskData(taskId) 加载消息,通过 assistantMessageId
定位目标 assistant 消息(而非数组位置——入队到执行之间可能有新 round 启动)。
优雅停机
worker.close() 等待 in-flight Job 完成后停止拉取。SIGTERM/SIGINT 触发三个 Worker + Redis 连接的协调关闭。
幂等性
两层去重:
-
入队级——BullMQ
jobId(如credit:{taskId}:{messageId})防止重复入队。同一 Job 无法被添加两次。 -
重试级——Job 失败后 BullMQ 重试时处理器会再次执行。对于
credit.consume,repository 使用 insert-first 模式 + ledger 表的 uniquereferenceId约束(ON CONFLICT DO NOTHING)。ledger 行已存在时跳过 balance 扣减——无 TOCTOU 竞态。其他处理器天然幂等(upsert / 覆盖)。
可观测性
三个组件:
-
Bull-board 仪表板——
/admin/queues(需 admin 认证)。完整 Web UI,可查看 Job 详情、重试失败 Job、查看 Payload。文件:apps/server/src/routes/admin-queues.ts。 -
队列指标 API——
GET /api/admin/queue-metrics(需 admin 认证)。返回每个队列的计数:active、waiting、delayed、failed、completedTotal、failedTotal。 -
失败 Webhook——
JOB_FAILURE_WEBHOOK_URL环境变量(可选)。所有重试耗尽后发送 POST,包含 Job 元数据。Worker 同时为active、completed、failed、stalled、error事件输出结构化日志。
文件清单
| 文件 | 角色 |
|---|---|
packages/backend/src/infra/job-queue.ts | JobQueue 接口 + createInlineJobQueue() |
apps/server/src/lib/queues.ts | 共享 Queue 实例(懒加载单例) |
apps/server/src/infra/bullmq-job-queue.ts | createBullMQJobQueue() |
apps/server/src/worker.ts | Worker 入口 |
apps/server/src/jobs/*.ts | 三个 Job 处理器 |
apps/server/src/routes/admin-queues.ts | Bull-board 仪表板路由 |
设计约束
- Payload 必须可序列化。 闭包和文件句柄无法跨序列化边界。
- 按 ID 查找消息。 Worker 通过
assistantMessageId定位消息,非数组位置。 - Desktop 使用内联队列。 同一
JobQueue接口,fire-and-forget 执行,无需 Redis。 - 开发环境 Redis 可选。
REDIS_URL未设置时 Server 回退到createInlineJobQueue()。 - MCP 断开连接留在进程内。 依赖
mcpClientManager状态,不入队。