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 密集代码真正并行运行。
适合提高并发的任务
- 调用第三方 HTTP API;
- 查询数据库;
- 读写对象存储;
- 发送邮件;
- 大量时间处于异步等待的任务。
不适合盲目提高并发的任务
- 图片压缩;
- 视频转码;
- 大型加密或哈希计算;
- 大量同步 JSON 计算;
- 其他长时间占用 JavaScript 主线程的任务。
CPU 密集任务提高普通 concurrency 可能降低吞吐量,并阻塞事件循环,导致 BullMQ 无法及时更新任务锁。此类工作应考虑 sandboxed processor、Worker Threads、多个进程或专门的计算服务。
4. 怎样选择并发值
不要直接复制一个“万能值”。建议从小值开始,例如 5 或 10,再根据以下指标调整:
- 每秒完成任务数;
- 单任务耗时分位数;
- CPU 和内存使用率;
- 数据库连接池等待;
- 第三方服务限流;
- Redis 延迟;
- stalled、失败和重试数量。
并发受到最慢下游系统约束。假设数据库连接池只有 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 仍可继续处理任务。
扩展前需要确认:
- processor 不依赖单机内存状态;
- 外部副作用具有幂等性;
- 所有 Worker 使用相同代码版本或具备向后兼容的数据契约;
- 数据库和第三方服务能承受总并发;
- 日志包含 job ID、队列名和业务关联 ID。
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() 会停止领取新任务,并等待正在执行的任务结束。官方文档提醒,它自身不会超时。因此:
- processor 应为网络请求设置超时;
- 长任务应支持取消或分段执行;
- 容器的 termination grace period 应大于典型任务完成时间;
- 不要立即调用
process.exit(),否则会中断异步清理。
即使异常退出,BullMQ 也能通过 stalled 机制让其他 Worker 重新处理失去锁的任务,但这会带来重复执行的可能,所以优雅关闭和幂等处理都不可少。
9. 暂停单个 Worker
await worker.pause();
// 完成维护后:
worker.resume();
默认 pause() 会等待当前任务完成。worker.pause(true) 可以不等待当前任务便返回,但已经运行的 processor 不会被自动取消。暂停单个 Worker 适合滚动维护;暂停整个队列则使用 queue.pause()。
实用技巧与最佳实践
- Worker 作为独立进程部署,不与 HTTP 服务争抢 CPU 和内存。
- 从低并发开始,以生产指标调优,不凭感觉设成很大。
- I/O 密集任务可以提高并发;CPU 密集任务优先扩进程或使用 sandboxed processor。
- 为所有外部调用设置超时,避免
worker.close()永久等待。 - 每个 processor 只负责一个清晰业务动作,复杂流程拆成多个任务或 Flow。
- 捕获错误时保留原始错误并抛出,让 BullMQ 正确标记失败和执行重试。
- 让 processor 幂等,防范重试、stalled 恢复和进程崩溃造成的重复处理。
- 监控 waiting、active、failed、stalled 数量与处理时长,而不只看进程是否存活。
本章小结
Worker 的本地并发主要提升 I/O 密集任务吞吐量,多个 Worker 则同时提供并行处理和高可用。CPU 密集任务不能只靠增大 concurrency;所有 Worker 都应监听错误、具备幂等性,并在退出时调用 worker.close()。