BullMQ 是基于 Redis 构建的开源分布式任务队列系统,支持 Node.js、Python、Rust 等多种运行环境,保证任务操作的原子性,提供延迟调度、指数退避重试、并发控制以及父子任务编排能力,适用于 Web 服务削峰填谷、异步解耦及后台长时任务处理等场景。
官方技术规范与源码参考:BullMQ 官方文档。
在现代后端架构中,当面临耗时较长的异步任务或高并发削峰诉求时,通常需要引入异步队列机制。在技术选型中,BullMQ 与专用分布式消息队列(如 Kafka、RabbitMQ)的定位存在明确差异:
BullMQ 采用生产者-消费者模型,由四个核心组件协作驱动:
data)、运行参数(opts)、状态标记与执行结果的数据实体。任务在生命周期内按确定性的状态机流转:
各状态对应的触发条件与底层 Redis 存储结构如下:
| 状态名称 | 触发条件 | Redis 底层结构 | 机制说明 |
|---|---|---|---|
| waiting | 任务入队或延迟倒计时结束 | List(列表) | 待处理任务队列,Worker 依赖阻塞弹出获取任务 |
| active | Worker 成功获取并锁定任务 | Hash + Lock Token | 处于执行状态,Worker 维持心跳锁防止任务被重复分配 |
| delayed | 配置延迟执行或处于失败重试退避期 | Sorted Set(有序集合) | 以预定执行时间戳作为 Score 排序,定时检测就绪任务并迁移至 waiting |
| completed | 处理函数正常返回 | Hash | 存储执行结果,依据保留策略按时间或数量保留 |
| failed | 处理函数抛出异常且重试次数耗尽 | Hash | 存储错误堆栈与失败详情,供排查回溯或人工重放 |
在服务端工程实现中,将业务逻辑与队列组件紧密耦合会导致代码难以维护与测试;若在 Web 主进程中直接运行队列消费逻辑,长时任务还容易抢占正常请求的处理资源。
推荐采用分层解耦的架构模式,将职责划分为三个独立层次(下文代码实现以 TypeScript 为例,其设计思想适用于各语言栈):
job.data 解析与校验参数、上报进度,并将具体逻辑委托给领域服务。编写业务处理函数,保持入参为纯数据对象,不依赖 BullMQ 运行时:
export interface SendNotificationParams {
recipientId: string;
channel: "email" | "sms" | "in_app";
templateId: string;
variables: Record<string, string | number>;
}
export async function sendNotification(
params: SendNotificationParams
): Promise<{ deliveredAt: string }> {
// 业务实现:读取数据库、调用下游通道接口等
console.log(`[Service] 处理通知分发: ${params.channel} -> ${params.recipientId}`);
return {
deliveredAt: new Date().toISOString(),
};
}编写 Worker 处理器,仅处理任务参数校验并委托业务层执行:
import { Worker, type Job, type WorkerOptions } from "bullmq";
import { sendNotification } from "./notificationService";
import type { NotificationJobData } from "./notificationQueue";
// 薄层任务处理器
export async function processNotificationJob(
job: Job<NotificationJobData>
): Promise<{ deliveredAt: string }> {
const { recipientId, channel, templateId, variables } = job.data;
if (!recipientId || !channel) {
throw new Error("任务参数校验失败:缺失 recipientId 或 channel");
}
await job.updateProgress(20);
// 委托给领域服务层
const result = await sendNotification({
recipientId,
channel,
templateId,
variables,
});
await job.updateProgress(100);
return result;
}
// 消费者工厂
export function createNotificationWorker(options?: Partial<WorkerOptions>) {
const worker = new Worker<NotificationJobData, { deliveredAt: string }>(
"notification-dispatch",
processNotificationJob,
{
connection: {
host: process.env.REDIS_HOST ?? "127.0.0.1",
port: Number(process.env.REDIS_PORT ?? 6379),
maxRetriesPerRequest: null, // BullMQ 官方要求:避免阻塞命令在断连时抛出错误
},
concurrency: 5, // 控制单个 Worker 实例的并发处理上限
...options,
}
);
worker.on("completed", (job, result) => {
console.log(`[Worker] 任务 ${job.id} 处理完成,时间: ${result.deliveredAt}`);
});
worker.on("failed", (job, err) => {
console.error(`[Worker] 任务 ${job?.id} 第 ${job?.attemptsMade} 次执行失败: ${err.message}`);
});
return worker;
}通过泛型约束任务载荷结构,并配置通用重试策略:
import { Queue, type ConnectionOptions, type JobsOptions } from "bullmq";
export interface NotificationJobData {
recipientId: string;
channel: "email" | "sms" | "in_app";
templateId: string;
variables: Record<string, string | number>;
timestamp: string;
}
const connection: ConnectionOptions = {
host: process.env.REDIS_HOST ?? "127.0.0.1",
port: Number(process.env.REDIS_PORT ?? 6379),
maxRetriesPerRequest: null,
};
const defaultJobOptions: JobsOptions = {
attempts: 3,
backoff: {
type: "exponential",
delay: 2000, // 失败后按 2s, 4s, 8s 指数退避重试
},
removeOnComplete: {
count: 1000, // 仅保留最新 1000 条完成记录
},
removeOnFail: {
count: 5000, // 保留最近 5000 条失败记录供排查
},
};
export const notificationQueue = new Queue<NotificationJobData>(
"notification-dispatch",
{
connection,
defaultJobOptions,
}
);在生产端投递任务,通过指定业务唯一键作为 jobId 实现入队去重:
export async function enqueueNotification(
data: NotificationJobData,
eventId?: string
): Promise<string> {
// 若传入业务事件 ID 则作为幂等键,避免相同事件重复入队
const jobId = eventId ? `notify_${data.channel}_${eventId}_${data.recipientId}` : undefined;
const job = await notificationQueue.add("send-notification", data, {
jobId,
});
return job.id!;
}jobId 的任务在队列中处于 waiting 或 delayed 状态,重复调用 add 将不会生成新任务。通过配置 delay 与 repeat 参数支持时间维度的调度控制:
// 延迟执行:15 分钟后触发
await notificationQueue.add(
"timeout-reminder",
reminderData,
{ delay: 15 * 60 * 1000 }
);
// 周期性调度:每小时整点执行(基于 Cron 表达式)
await notificationQueue.add(
"hourly-digest",
digestData,
{
repeat: {
pattern: "0 * * * *",
},
}
);Redis 作为内存数据库,未清理的任务数据会持续占用内存资源。必须在队列或任务维度显式设置保留上限:
const productionJobOptions: JobsOptions = {
removeOnComplete: {
age: 24 * 3600, // 自动清理超过 24 小时的完成记录
count: 5000, // 最多保留最新 5000 条
},
removeOnFail: {
age: 7 * 24 * 3600, // 失败任务保留 7 天以供排查追溯
},
};在生产环境中,建议将 Worker 部署为与 Web 服务物理隔离的独立守护进程(如独立的容器或通过 PM2 托管的后台进程),避免长耗时任务阻塞 Web 服务的事件循环。
通过编写独立的启动脚本(如 scripts/start-worker.ts),汇集并实例化各 Worker,同时监听操作系统的退出信号以实现平滑退出(Graceful Shutdown):
// scripts/start-worker.ts
import { createNotificationWorker } from "../workers/notificationWorker";
console.log("[Worker] 启动后台消费者守护进程...");
// 批量实例化需要监听的 Worker
const notificationWorker = createNotificationWorker();
// 注册退出信号处理逻辑
async function gracefulShutdown(signal: string) {
console.log(`[Worker] 收到退出信号 [${signal}],开始平滑关闭...`);
try {
// 停止提取新任务,并等待正在执行中的任务处理完毕
await Promise.allSettled([
notificationWorker.close(),
]);
console.log("[Worker] 所有消费者已安全关闭");
process.exit(0);
} catch (err) {
console.error("[Worker] 关闭期间发生异常:", err);
process.exit(1);
}
}
process.on("SIGTERM", () => gracefulShutdown("SIGTERM"));
process.on("SIGINT", () => gracefulShutdown("SIGINT"));通过集成官方支持的看板组件 @bull-board,可可视化监控各队列的任务积压、失败率及重试状态:
import express from "express";
import { createBullBoard } from "@bull-board/api";
import { BullMQAdapter } from "@bull-board/api/bullMQAdapter";
import { ExpressAdapter } from "@bull-board/express";
import { notificationQueue } from "./notificationQueue";
const app = express();
const serverAdapter = new ExpressAdapter();
serverAdapter.setBasePath("/admin/queues");
createBullBoard({
queues: [new BullMQAdapter(notificationQueue)],
serverAdapter,
});
app.use("/admin/queues", serverAdapter.getRouter());
app.listen(3001, () => {
console.log("[Dashboard] 监控面板已启动: http://127.0.0.1:3001/admin/queues");
});BullMQ 为现代分布式任务调度与削峰填谷提供了轻量且完备的工程方案: