发邮件、生成报表、处理图片、调用 LLM 批量摘要……这些慢任务不该阻塞 HTTP 响应。BullMQ 是 Node 生态最流行的 Redis 任务队列:支持优先级、延迟任务、重复调度、失败重试、并发沙箱。本文从 Redis 队列模型讲起,覆盖 Producer/Worker、调度、Backoff、沙箱与可观测性。
1. 为什么用任务队列
1.1 同步阻塞的代价
app.post('/api/export', async (req, res) => {
await buildBigReport(req.body.userId); // 可能要跑 2 分钟
res.json({ ok: true }); // 客户端等到超时
});
慢任务放进 HTTP 同步链路,服务吞吐被拖垮、网关超时、用户体验差。正确姿势是立即受理,后台执行。
1.2 队列的收益
| 收益 | 说明 |
|---|---|
| 快速响应 | 请求立刻返回,任务后台执行 |
| 削峰填谷 | 高峰排队,低谷消化 |
| 失败重试 | 任务失败自动重跑 |
| 水平扩展 | 加 Worker 实例即扩并发 |
| 解耦 | 生产与消费分离 |
一句话:任务队列 = 「立即受理 + 后台执行」——慢任务出 HTTP 链路,进队列,由 Worker 慢慢消化。
2. Redis 队列模型
2.1 BullMQ 的数据结构
BullMQ 每个队列在 Redis 里用多组 key 组织:
bull:mail:wait 待执行任务列表
bull:mail:active 正在执行的任务
bull:mail:delayed 延迟/定时任务
bull:mail:failed 失败任务
bull:mail:completed 已完成任务
bull:mail:repeat 重复任务定义
任务的流转:wait → active → completed/failed,失败进入重试逻辑。
2.2 安装与连接
npm i bullmq ioredis
import { Queue, Worker, QueueEvents } from 'bullmq';
import { Redis } from 'ioredis';
const connection = new Redis({ host: '127.0.0.1', port: 6379, maxRetriesPerRequest: null });
2.3 任务三要素
Job 名称 → 决定由哪个处理器消费
Job 数据 → payload,序列化成 JSON
Job 选项 → 优先级、延迟、重试、超时
一句话:BullMQ = Redis 里一组 key 构成的队列状态机,job 沿
wait → active → completed/failed流转;共享连接时maxRetriesPerRequest: null必须设。
3. Producer 与 Worker
3.1 Producer 生产任务
const queue = new Queue('mail', { connection });
// 最简单的入队
await queue.add('send-welcome', { userId: 123, template: 'welcome' });
3.2 任务选项
await queue.add('send-batch', { fileId: 'abc' }, {
priority: 5, // 数字越小越先执行
delay: 60 * 1000, // 延迟 1 分钟执行
attempts: 3, // 失败重试 3 次
removeOnComplete: true, // 完成后自动清理
removeOnFail: 100, // 失败最多保留 100 条
});
3.3 Worker 消费任务
const worker = new Worker('mail', async (job) => {
const { userId, template } = job.data;
await sendMail(userId, template);
return { sent: userId }; // 返回值写入 job.returnvalue
}, { connection, concurrency: 5 });
worker.on('completed', (job) => console.log(`完成 ${job.id}`));
worker.on('failed', (job, err) => console.error(`失败 ${job.id}: ${err.message}`));
Worker 处理器必须处理完再 resolve,异常会触发失败与重试逻辑。
一句话:Producer
queue.add(name, data, options)入队,Workernew Worker(name, handler)消费;handler 抛异常即失败,可重试。
4. 调度与重复任务
4.1 一次性延迟任务
await queue.add('remind', { orderId: 7 }, { delay: 30 * 60 * 1000 });
任务入队后进 delayed,到点转 wait。适合订单超时提醒、支付倒计时。
4.2 cron 重复任务
const repeatable = await queue.add('daily-report', { type: 'sales' }, {
repeat: {
pattern: '0 8 * * *', // 每天 8 点,cron 表达式
tz: 'Asia/Shanghai',
},
jobId: 'daily-report', // 固定 jobId 保证幂等
});
// 停止重复
await queue.removeRepeatable('daily-report', {
pattern: '0 8 * * *', tz: 'Asia/Shanghai',
});
4.3 重复任务 vs 调度器
需要「精确到秒、动态跳过」 → 自定义循环任务(setTimeout + 重新入队)
固定周期执行 → repeat cron 足够
一句话:
delay做一次性延迟,repeat.pattern做 cron 周期任务;重复任务固定 jobId + 幂等消费,否则会重复入队。
5. 失败重试与 Backoff
5.1 重试配置
await queue.add('convert-video', { src: 'a.mp4' }, {
attempts: 5, // 含首次,共尝试 5 次
backoff: {
type: 'exponential', // 指数退避
delay: 1000, // 首次等待 1 秒,之后 2 倍递增
},
jobId: 'video-a', // 幂等:同 ID 不重复入队
});
5.2 区分可重试与不可重试
处理器里对「永久性错误」主动抛特殊标记,避免无效重试:
class PermanentError extends Error {
constructor(message) { super(message); this.permanent = true; }
}
const worker = new Worker('convert-video', async (job) => {
const src = job.data.src;
if (!await fileExists(src)) throw new PermanentError('源文件不存在');
await convert(src);
});
worker.on('failed', async (job, err) => {
if (err.permanent) await job.remove(); // 永久错误,移除不再重试
});
5.3 停滞任务保护
Worker 崩溃导致任务卡在 active,BullMQ 会检测 stalled 并重新入队:
const worker = new Worker('mail', handler, {
stalledInterval: 30 * 1000, // 每 30 秒检查停滞
lockDuration: 60 * 1000, // 任务锁时长,超时判定停滞
});
一句话:失败重试 =
attempts+ 指数退避 + 固定 jobId 幂等;永久错误要主动标记并移除,避免无效重试刷爆日志。
6. 并发与沙箱
6.1 并发控制
// 单 Worker 内并发
const worker = new Worker('image', handler, { concurrency: 10 });
// 横向扩展:多 Worker 进程共享同一队列
// 每个进程 new Worker(...),Redis 负责分发
concurrency 控制单进程并行数;多进程部署时每个进程都是独立消费者,Redis 保证同一任务只被一个 Worker 拿走。
6.2 沙箱处理器隔离崩溃
处理器里 process.exit、未捕获异常会拖垮 Worker。用沙箱隔离:
// 处理器文件 processor.ts
export default async function (job) {
if (job.data.malicious) process.exit(1); // 只会杀沙箱,不影响主进程
return heavyWork(job.data);
}
import { Job } from 'bullmq';
const worker = new Worker('sandboxed', new Job.SandboxedJobProcessor('processor.ts'), {
useWorkerThreads: true, // 用 worker_threads 隔离
});
6.3 沙箱的限制
沙箱处理器与主进程隔离内存与状态,不能访问主进程变量;依赖与配置要在处理器文件内自给自足。
一句话:并发 =
concurrency控单进程 + 多 Worker 进程横向扩展;处理器有崩溃风险时用 SandboxedJobProcessor 隔离,防止拖垮主进程。
7. 可观测性与运维
7.1 事件监听
const events = new QueueEvents('mail', { connection });
events.on('completed', ({ jobId, returnvalue }) => {
console.log(`job ${jobId} 完成,结果:${JSON.stringify(returnvalue)}`);
});
events.on('failed', ({ jobId, failedReason }) => {
console.error(`job ${jobId} 失败:${failedReason}`);
});
events.on('progress', ({ jobId, data }) => {
console.log(`job ${jobId} 进度 ${data}%`);
});
处理器里上报进度:
await job.updateProgress(50); // 触发 progress 事件
7.2 队列状态统计
const counts = await queue.getJobCounts('wait', 'active', 'delayed', 'failed');
console.log(counts); // { wait: 12, active: 3, delayed: 1, failed: 2 }
配合定时任务上报指标,或接入 Prometheus 做告警:failed 持续增长说明系统异常。
7.3 运维命令
await queue.obliterate({ force: true }); // 清空队列(含正在执行的)
await queue.clean(3600 * 1000, 'completed'); // 清理 1 小时前的完成记录
常用 Grafana 方案是 bull-monitor,可视化队列长度、失败率、重试次数。
| 指标 | 正常值 | 报警值 |
|---|---|---|
| wait 堆积 | 秒级消化 | 持续增长超阈值 |
| failed 率 | < 1% | 持续 > 5% |
| stalled | 极少 | 频繁出现 |
一句话:可观测 = QueueEvents 监听 completed/failed/progress + 定时统计各状态计数;
wait堆积和failed攀升是队列告警的核心指标。
8. 踩坑清单
| 坑 | 现象 | 对策 |
|---|---|---|
| maxRetriesPerRequest 未设 | Redis 重连时任务报错 | 连接设 null |
| 重复任务不设 jobId | 每次 add 都新增 | 固定 jobId + 幂等 |
| 处理器抛永久错误重试 | 无效重试刷爆 | PermanentError 标记移除 |
| 慢任务占满 concurrency | 其他任务饿死 | 拆分队列 + 调优先级 |
| 大 payload 入队 | Redis 内存暴涨 | payload 只放引用 ID |
| 不监听 failed | 任务静默消失 | QueueEvents failed 告警 |
| 忘记 clean | completed 无限堆积 | removeOnComplete + clean |
| 处理器污染主进程 | Worker 无端退出 | SandboxedJobProcessor |
| 时区错配 cron | 执行时间不对 | repeat.tz 显式指定 |
9. 总结
| 环节 | 要点 |
|---|---|
| 场景 | 慢任务移出 HTTP 链路 |
| 模型 | Redis key:wait/active/delayed/failed |
| Producer | queue.add(name, data, options) |
| Worker | new Worker + handler,并发控制 |
| 调度 | delay 延迟 / repeat cron |
| 重试 | attempts + 指数退避 + 幂等 jobId |
| 并发 | concurrency + 多进程扩展 |
| 沙箱 | SandboxedJobProcessor 隔离 |
| 观测 | QueueEvents + 状态计数 + 告警 |
一句话记住:BullMQ 是把「慢任务」从 HTTP 响应里解放出来的 Redis 任务队列——Producer 快速入队、Worker 后台消化、指数退避扛住失败、沙箱隔离崩溃,再配一套事件与计数观测,后台任务从此可控可查。
延伸阅读
- Node.js 消息队列实践 — 消息队列与任务队列的定位差异
- Node.js 异步与并发 — Worker 消费背后的异步模型
- Node.js 可观测性 — 任务处理的指标与追踪
- Node.js 错误处理与日志 — 失败任务的异常归类
- Node.js 缓存策略 — 结果缓存与队列组合
- Node.js Docker 与 Kubernetes — Worker 独立部署与扩缩容
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。