06 - Worker、并发与水平扩展

Worker 是真正执行任务的消费者。它从队列领取 Job,调用 processor;processor 正常返回时任务进入 completed,抛出异常时进入 failed 或按配置重试。

1. 编写类型明确的 Worker

import { Job, Worker } from "bullmq";

interface ExportJobData {
  tenantId: string;
  reportId: string;
}

interface ExportJobResult {
  objectKey: string;
}

type ExportJobName = "export-report";

const worker = new Worker<
  ExportJobData,
  ExportJobResult,
  ExportJobName
>(
  "report-export",
  async (
    job: Job<ExportJobData, ExportJobResult, ExportJobName>,
  ): Promise<ExportJobResult> => {
    console.log(`开始导出 ${job.data.reportId}`);

    const objectKey = await exportReport(job.data);
    return { objectKey };
  },
  {
    connection: {
      host: "127.0.0.1",
      port: 6379,
    },
    concurrency: 5,
  },
);

processor 必须返回 Promise 或使用 async。创建 Worker 后默认会立即开始消费;需要延后时可以设置 autorun: false,初始化完成后调用 worker.run()

2. 必须监听的重要事件

worker.on("completed", (job, result) => {
  console.log(`任务 ${job.id} 完成`, result);
});

worker.on("failed", (job, error: Error) => {
  console.error(`任务 ${job?.id ?? "unknown"} 失败`, error);
});

worker.on("error", (error: Error) => {
  console.error("Worker 内部或连接错误", error);
});

failed 表示某个 Job 处理失败,error 则可能是 Worker 自身或 Redis 连接错误。官方文档特别要求为 Worker 添加 error 监听器;缺少监听器可能导致 Worker 停止处理任务。

3. concurrency 是什么

const worker = new Worker("email", processEmail, {
  connection,
  concurrency: 20,
});

concurrency: 20 表示这个 Worker 实例最多同时推进 20 个 Job。它主要利用 Node.js 事件循环,在任务等待数据库、HTTP、磁盘等 I/O 时处理其他任务。

它不代表单个 Node.js 线程能让 20 段 CPU 密集代码真正并行运行。

适合提高并发的任务

不适合盲目提高并发的任务

CPU 密集任务提高普通 concurrency 可能降低吞吐量,并阻塞事件循环,导致 BullMQ 无法及时更新任务锁。此类工作应考虑 sandboxed processor、Worker Threads、多个进程或专门的计算服务。

4. 怎样选择并发值

不要直接复制一个“万能值”。建议从小值开始,例如 5 或 10,再根据以下指标调整:

并发受到最慢下游系统约束。假设数据库连接池只有 10 个连接,将 Worker 并发设置为 200 往往只是制造排队和超时。

运行时也可以修改本地并发值:

worker.concurrency = 10;

但生产环境更常通过配置和重新部署管理,便于审计和容量规划。

5. 严格依次执行

如果这个 Worker 每次只应处理一个任务:

const worker = new Worker("legacy-system-sync", processSync, {
  connection,
  concurrency: 1,
});

注意:concurrency: 1 只约束当前 Worker 实例。如果启动了三个同队列 Worker,最多仍可能同时执行三个任务。需要整个队列全局只执行一个任务时,应使用 Queue 的 global concurrency 能力,或确保部署层只有一个 Worker 实例。

另外,失败重试、优先级和延迟任务可能影响观察到的完成顺序。不要把“FIFO 入队”误解为所有场景下都严格按顺序完成。

6. 水平扩展多个 Worker

                 ┌─ Worker A(concurrency: 10)
Producer → Redis ├─ Worker B(concurrency: 10)
                 └─ Worker C(concurrency: 10)

多个 Worker 可以运行在不同 Node.js 进程、容器或机器上。BullMQ 会把任务分发给可用 Worker。

官方推荐使用多个 Worker,因为这不仅增加吞吐量,也提高可用性:一个 Worker 下线后,其他 Worker 仍可继续处理任务。

扩展前需要确认:

7. 一个 Worker 处理多种任务

type NotificationJobName = "send-email" | "send-sms";

interface NotificationJobData {
  recipient: string;
  content: string;
}

const worker = new Worker<
  NotificationJobData,
  void,
  NotificationJobName
>("notification", async (job): Promise<void> => {
  switch (job.name) {
    case "send-email":
      await sendEmail(job.data);
      return;
    case "send-sms":
      await sendSms(job.data);
      return;
    default: {
      const exhaustiveCheck: never = job.name;
      throw new Error(`未知任务类型:${exhaustiveCheck}`);
    }
  }
});

任务共享相似的资源、并发和重试策略时可以放在同一队列。若邮件和短信的限流、故障域或扩容方式明显不同,拆成两个队列会更清晰。

8. 优雅关闭

进程收到终止信号时,不应直接退出:

let isShuttingDown = false;

async function shutdown(signal: NodeJS.Signals): Promise<void> {
  if (isShuttingDown) {
    return;
  }

  isShuttingDown = true;
  console.log(`收到 ${signal},停止领取新任务`);

  try {
    await worker.close();
    console.log("Worker 已安全关闭");
  } catch (error: unknown) {
    console.error("Worker 关闭失败", error);
    process.exitCode = 1;
  }
}

process.once("SIGINT", () => {
  void shutdown("SIGINT");
});

process.once("SIGTERM", () => {
  void shutdown("SIGTERM");
});

worker.close() 会停止领取新任务,并等待正在执行的任务结束。官方文档提醒,它自身不会超时。因此:

即使异常退出,BullMQ 也能通过 stalled 机制让其他 Worker 重新处理失去锁的任务,但这会带来重复执行的可能,所以优雅关闭和幂等处理都不可少。

9. 暂停单个 Worker

await worker.pause();

// 完成维护后:
worker.resume();

默认 pause() 会等待当前任务完成。worker.pause(true) 可以不等待当前任务便返回,但已经运行的 processor 不会被自动取消。暂停单个 Worker 适合滚动维护;暂停整个队列则使用 queue.pause()

实用技巧与最佳实践

本章小结

Worker 的本地并发主要提升 I/O 密集任务吞吐量,多个 Worker 则同时提供并行处理和高可用。CPU 密集任务不能只靠增大 concurrency;所有 Worker 都应监听错误、具备幂等性,并在退出时调用 worker.close()

官方资料