后台任务队列

基于 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 + JobQueue 接口 — 双路径分发 onEnd — 后台任务分发 credit.consume | memory.extraction | resource.index JobQueue.enqueue( jobName, payload, execute, options ) Server (Redis) Desktop / Dev createBullMQJobQueue() payload 序列化到 Redis (Queue.add) createInlineJobQueue() 直接调用 execute() Redis (BullMQ Queues) Worker 进程 critical (C=10) credit.consume llm (C=2) memory.extraction indexing (C=5) resource.index 从 payload ID 重建依赖 → DB / LLM / Sandbox jobId 入队去重 + unique ledger 约束 (credit 重试幂等) 可观测性 (Observability) bull-board UI | queue-metrics API | failure webhook 1 2 3 4 5

BullMQ 工作原理

BullMQ 是基于 Redis 的 Node.js 任务队列。先看懂它的 Job 生命周期,后面的重试行为和故障模式才推理得清楚。

BullMQ Job 生命周期 Producer → Redis Streams → Worker — 含重试与退避 Producer API Server / Orchestrator Queue.add(name, payload) Redis Streams + Sorted Sets 持久化 Job 存储 Worker 独立进程 阻塞拉取 (BRPOPLPUSH) Job 状态机 (State Machine) waiting Worker 拉取 active 成功 completed 异常 failed attempts < max? delayed 退避等待 exhausted Webhook 告警 持久化 (Persistence) Job 在进程重启后存活 Redis 通过 AOF/RDB 将 waiting/delayed 持久化到磁盘 vs 内联: 进程崩溃 = Job 丢失 vs setTimeout: 无持久化 vs DB 轮询: 延迟更高 并发控制 (Concurrency) 每个队列有并发限制 (C) Worker 每队列最多并行 处理 C 个 Job critical C=10: 快速消费扣费 llm C=2: 控制 provider 限流 indexing C=5: DB 吞吐量约束

三个角色:

  • 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 状态流转: waitingactivecompletedfailed。失败后如仍有重试次数,Job 进入 delayed(指数或固定退避),然后回到 waiting。所有重试耗尽后,Job 永久留在 failed——在 bull-board 仪表板可见,并触发 Webhook 告警。

为什么不直接用 setTimeoutPromise

  • setTimeoutvoid promise 期间进程崩溃 = Job 永久丢失。BullMQ Job 持久化在 Redis 中。
  • setTimeout 无重试。BullMQ 支持可配置退避的重试。
  • 裸 Promise 无并发控制。BullMQ 按队列限制并行执行数。
  • Fire-and-forget 无可观测性。BullMQ 提供状态、耗时、尝试次数和失败原因。

JobQueue 接口

遵循项目的 infra 模式(TaskLockStreamBufferKeyEncryption),队列定义为 @zapvol/backend/src/infra/job-queue.ts 中的接口,有两个实现:

export interface JobQueue {
  enqueue(
    jobName: string,
    payload: Record<string, unknown>,
    execute: () => Promise<void>,
    options?: JobEnqueueOptions,
  ): void;
}

核心设计:enqueue 同时接收可序列化 payloadexecute 闭包

实现payloadexecute
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),拆开是为了一种负载不会饿死另一种:

队列任务为什么单独隔离
agenttask.run / chat.runagent 运行本身——长时、IO-bound(等 LLM/工具);高并发,与其它队列隔开,免得一波运行饿死扣费 / 索引
criticalcredit.consume扣费必须快速消费——绝不被慢速 LLM 任务饿死
llmmemory.extraction低并发——控制 provider 限流与 token 成本
indexingresource.indexCPU / DB 密集型——独立伸缩
scheduling定时任务触发wait_and_resume / cron 续跑
nudgenudge 触发时间触发的 nudge

queues.ts 里的懒加载单例 Queue 实例由 Job 队列、bull-board 仪表板与指标端点共享;每个队列的 Workerworker.ts 里各带自己的并发。

Worker 进程

单个 Worker 进程(apps/server/src/worker.ts)为每个队列各起一个 Worker,按 job name 经 JOB_PROCESSORS 分发。独立于 API 服务器运行:pnpm workerpnpm 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 连接的协调关闭。

幂等性

两层去重:

  1. 入队级——BullMQ jobId(如 credit:{taskId}:{messageId})防止重复入队。同一 Job 无法被添加两次。

  2. 重试级——Job 失败后 BullMQ 重试时处理器会再次执行。对于 credit.consume,repository 使用 insert-first 模式 + ledger 表的 unique referenceId 约束(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 认证)。返回每个队列的计数: activewaitingdelayedfailedcompletedTotalfailedTotal

  • 失败 Webhook——JOB_FAILURE_WEBHOOK_URL 环境变量(可选)。所有重试耗尽后发送 POST,包含 Job 元数据。Worker 同时为 activecompletedfailedstallederror 事件输出结构化日志。

文件清单

文件角色
packages/backend/src/infra/job-queue.tsJobQueue 接口 + createInlineJobQueue()
apps/server/src/lib/queues.ts共享 Queue 实例(懒加载单例)
apps/server/src/infra/bullmq-job-queue.tscreateBullMQJobQueue()
apps/server/src/worker.tsWorker 入口
apps/server/src/jobs/*.ts三个 Job 处理器
apps/server/src/routes/admin-queues.tsBull-board 仪表板路由

设计约束

  1. Payload 必须可序列化。 闭包和文件句柄无法跨序列化边界。
  2. 按 ID 查找消息。 Worker 通过 assistantMessageId 定位消息,非数组位置。
  3. Desktop 使用内联队列。 同一 JobQueue 接口,fire-and-forget 执行,无需 Redis。
  4. 开发环境 Redis 可选。 REDIS_URL 未设置时 Server 回退到 createInlineJobQueue()
  5. MCP 断开连接留在进程内。 依赖 mcpClientManager 状态,不入队。
这页有帮助吗?