没问题!我把 BullMQ 四大核心类整理成了一份完整笔记,直接加到你之前的 Redis 笔记后面,形成一份 《Redis + BullMQ 全栈笔记》。复制保存,随时查阅!🚀
📚 BullMQ 四大核心类笔记(精华版)
一句话总览:BullMQ 是基于 Redis 的 Node.js 任务队列库。Queue 放任务,Worker 干活,QueueEvents 广播状态,FlowProducer 编排复杂依赖。
一、工厂流水线比喻(秒懂版)
| 核心类 |
角色比喻 |
一句话职责 |
Queue |
工厂的任务清单 |
负责创建任务,放到队列里。 |
Worker |
流水线上的工人 |
负责取出任务并执行。 |
QueueEvents |
工厂的广播系统 |
负责监听任务状态变化并通知你。 |
FlowProducer |
负责复杂组装的工头 |
负责编排有父子依赖的复杂任务流。 |
二、Queue(队列)—— 任务的生产者
作用:创建任务、管理队列。是任务的入口。
核心方法
| 方法 |
语法 |
说明 |
add |
queue.add(name, data, options) |
最核心。添加一个任务。 |
addBulk |
queue.addBulk([job1, job2]) |
批量添加多个任务。 |
pause |
queue.pause() |
暂停队列,暂停后 Worker 不再拉取新任务。 |
resume |
queue.resume() |
恢复队列。 |
getJobs |
queue.getJobs(['waiting', 'active']) |
获取指定状态的任务列表。 |
clean |
queue.clean(3600000, 1000) |
清理已完成/失败的任务(避免 Redis 内存爆了)。 |
add() 常用选项(options)
| 选项 |
类型 |
说明 |
delay |
number |
延迟执行(毫秒)。 |
attempts |
number |
失败后重试次数。 |
backoff |
object |
重试间隔策略。{ type: 'exponential', delay: 1000 } |
priority |
number |
优先级(1 最高,默认 0 无优先级)。 |
jobId |
string |
自定义任务 ID(默认自动生成)。 |
removeOnComplete |
boolean / number |
完成后自动删除(或保留数量)。 |
removeOnFail |
boolean / number |
失败后自动删除(或保留数量)。 |
代码示例
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16
| import { Queue } from 'bullmq';
const myQueue = new Queue('email-queue');
await myQueue.add('send-welcome', { userId: 123, email: 'user@example.com' });
await myQueue.add('send-reminder', { userId: 123 }, { delay: 5000 });
await myQueue.add('send-critical', { data: 'important' }, { attempts: 3, backoff: { type: 'exponential', delay: 1000 }, removeOnComplete: 10, });
|
三、Worker(工作器)—— 任务的执行者
作用:从队列取出任务并执行你定义的业务逻辑。是任务的消费者。
核心配置
| 参数 |
类型 |
说明 |
queueName |
string |
要监听的队列名称(必填)。 |
processor |
function |
处理任务的函数(必填)。可以返回 Promise。 |
concurrency |
number |
并发数。一个 Worker 同时能处理多少个任务。默认 1。 |
connection |
object |
Redis 连接配置(可自定义 host/port/password)。 |
autorun |
boolean |
是否自动开始拉取任务。默认 true。 |
limiter |
object |
速率限制。{ max: 10, duration: 1000 } 每秒最多处理 10 个。 |
核心方法
| 方法 |
说明 |
worker.on('completed', callback) |
监听任务完成事件。 |
worker.on('failed', callback) |
监听任务失败事件。 |
worker.on('progress', callback) |
监听任务进度更新事件。 |
worker.on('error', callback) |
监听 Worker 自身错误。 |
worker.close() |
优雅关闭 Worker(等待当前任务完成)。 |
代码示例
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28
| import { Worker } from 'bullmq';
const worker = new Worker( 'email-queue', async (job) => { console.log(`处理任务 ${job.id}: 发送邮件给 ${job.data.email}`);
await sendEmail(job.data.email);
await job.updateProgress(100);
return { sent: true, to: job.data.email }; }, { concurrency: 5, limiter: { max: 10, duration: 1000 }, } );
worker.on('error', (err) => { console.error('Worker 出错了:', err); });
|
四、QueueEvents(队列事件)—— 全局广播系统
作用:监听队列中所有任务的全局事件。需要独立的 Redis 连接,不占用 Worker 的资源。
核心事件
| 事件名 |
触发时机 |
回调参数 |
'completed' |
任务成功完成 |
{ jobId, returnvalue, prev, ... } |
'failed' |
任务失败 |
{ jobId, failedReason, ... } |
'progress' |
任务上报进度(job.updateProgress()) |
{ jobId, data, ... } |
'stalled' |
任务卡住(Worker 长时间未响应) |
{ jobId, ... } |
'delayed' |
任务被延迟执行 |
{ jobId, ... } |
'drained' |
队列中无等待任务 |
{ ... } |
'removed' |
任务被移除 |
{ jobId, ... } |
代码示例
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20
| import { QueueEvents } from 'bullmq';
const queueEvents = new QueueEvents('email-queue');
queueEvents.on('completed', ({ jobId, returnvalue }) => { console.log(`✅ 任务 ${jobId} 完成!返回值:`, returnvalue); });
queueEvents.on('failed', ({ jobId, failedReason }) => { console.error(`❌ 任务 ${jobId} 失败了,原因:`, failedReason); });
queueEvents.on('progress', ({ jobId, data }) => { console.log(`📊 任务 ${jobId} 进度: ${data}%`); });
|
五、FlowProducer(流式生产者)—— 复杂任务工头
作用:创建有父子依赖关系的任务流。父任务会等待所有子任务完成后才开始执行。
核心方法
| 方法 |
说明 |
add(flow) |
添加一个任务流。 |
addBulk(flows) |
批量添加多个任务流。 |
Flow 对象结构
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21
| { name: 'parent-job-name', queueName: 'parent-queue', data: { ... }, options: { ... }, children: [ { name: 'child-job-1', queueName: 'child-queue-1', data: { ... }, options: { ... }, }, { name: 'child-job-2', queueName: 'child-queue-2', data: { ... }, options: { ... }, children: [ ... ] } ] }
|
代码示例
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24
| import { FlowProducer } from 'bullmq';
const flowProducer = new FlowProducer();
const flow = await flowProducer.add({ name: 'generate-report', queueName: 'report-queue', data: { reportId: 123 }, children: [ { name: 'fetch-data', queueName: 'data-queue', data: { source: 'database' }, }, { name: 'process-data', queueName: 'process-queue', data: { transform: 'summarize' }, }, ], });
console.log('父任务 ID:', flow.job.id); console.log('子任务 IDs:', flow.children.map(c => c.job.id));
|
⚠️ 注意事项
- 子任务和父任务可以在不同的队列中。
- 原子性:
add() 要么整个任务流全部添加成功,要么全部失败(不会出现部分添加的情况)。
- 父子任务之间通过
jobId 关联,BullMQ 自动管理依赖关系。
- 支持深度嵌套(多层父子关系)。
六、完整工作流程图
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23
| ┌─────────────┐ ┌─────────────┐ ┌─────────────┐ │ 你的应用 │ │ Redis │ │ Worker │ │ (Producer) │ │ (存储队列) │ │ (Consumer) │ └──────┬──────┘ └──────┬──────┘ └──────┬──────┘ │ │ │ │ 1. queue.add() │ │ │───────────────────>│ │ │ │ 2. 任务进入等待 │ │ │ │ │ │ 3. worker 拉取任务 │ │ │<───────────────────│ │ │ │ │ │ 4. worker 执行 │ │ │ (业务逻辑) │ │ │ │ │ │ 5. worker 标记完成 │ │ │<───────────────────│ │ │ │ │ 6. QueueEvents │ │ │ 监听事件触发 │ │ │<───────────────────│ │ │ │ │ └────────────────────┴────────────────────┘
|
如果涉及父子依赖:
1 2 3 4 5 6 7 8
| ┌─────────────────────────────────────────────────┐ │ FlowProducer.add() │ │ ├── 子任务 A (在 queue-A 中) │ │ ├── 子任务 B (在 queue-B 中) │ │ └── 父任务 C (在 queue-C 中) │ │ ↓ 等待 A 和 B 都完成 │ │ 才出现在 queue-C 中等待 Worker 处理 │ └─────────────────────────────────────────────────┘
|
七、关键术语/缩写词典(补充)
| 术语 |
全称/含义 |
说明 |
| Job |
任务 |
一个任务的实例,包含 id、name、data、status 等。 |
| Producer |
生产者 |
使用 Queue 或 FlowProducer 添加任务的代码。 |
| Consumer |
消费者 |
使用 Worker 执行任务的代码。 |
| Concurrency |
并发数 |
一个 Worker 同时处理的任务数量。 |
| Backoff |
退避重试 |
任务失败后等待一段时间再重试的策略(固定、指数等)。 |
| Stalled |
卡住/僵持 |
Worker 长时间未上报心跳,任务被判定为“卡住”,会被重新放回队列。 |
| Rate Limiter |
速率限制 |
限制单位时间内处理的任务数量,保护下游系统。 |
八、快速上手 4 步曲(最小化示例)
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21
|
import { Queue } from 'bullmq'; const queue = new Queue('my-queue');
import { Worker } from 'bullmq'; const worker = new Worker('my-queue', async (job) => { console.log('处理:', job.data); await doSomething(job.data); });
await queue.add('my-job', { hello: 'world' });
import { QueueEvents } from 'bullmq'; const events = new QueueEvents('my-queue'); events.on('completed', ({ jobId }) => console.log(`${jobId} 完成!`));
|
九、常见面试题(自测用)
| 问题 |
答案要点 |
| BullMQ 的四大核心类是什么? |
Queue、Worker、QueueEvents、FlowProducer。 |
Queue 和 Worker 的区别? |
Queue 是生产者(创建任务),Worker 是消费者(执行任务)。 |
QueueEvents 和 Worker 的事件监听有什么区别? |
Worker.on() 只能监听该 Worker 处理的任务;QueueEvents 监听队列中所有任务,全局视角。 |
什么时候用 FlowProducer? |
需要处理有父子依赖关系的任务流,比如“生成报告”需要等“拉取数据”和“处理数据”都完成。 |
concurrency: 5 是什么意思? |
这个 Worker 同时最多处理 5 个任务。 |
总结: 你现在的笔记已经覆盖了 Redis 基础语法 + BullMQ 四大核心类,足够应对实际开发中的任务队列场景。接下来的学习路线可以按这个顺序走:
- 本地跑通最小示例(上面的 4 步曲)
- 实战一个场景:发邮件/生成报表/图片处理
- 学习高级特性:重试策略、速率限制、任务去重
需要我陪你跑通第一个完整示例,还是直接上实战场景?随时喊我!😄