没问题!我把 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' });

// 添加一个延迟 5 秒执行的任务
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, // 只保留最近 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) => {
// job.name: 'send-welcome'(add 时传的 name)
// job.data: { userId: 123, email: 'user@example.com' }
console.log(`处理任务 ${job.id}: 发送邮件给 ${job.data.email}`);

// 模拟耗时的邮件发送
await sendEmail(job.data.email);

// 任务进度上报(可以被 QueueEvents 监听到)
await job.updateProgress(100);

// 返回结果(会被记录到任务中)
return { sent: true, to: job.data.email };
},
{
concurrency: 5, // 同时处理 5 个任务
limiter: { max: 10, duration: 1000 }, // 每秒最多 10 个
}
);

// 监听 worker 自己的错误
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: { ... }, // 父任务选项(delay, attempts 等)
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 任务 一个任务的实例,包含 idnamedatastatus 等。
Producer 生产者 使用 QueueFlowProducer 添加任务的代码。
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
// 1. 安装依赖
// npm install bullmq ioredis

// 2. 创建 Queue(生产者)
import { Queue } from 'bullmq';
const queue = new Queue('my-queue');

// 3. 创建 Worker(消费者)
import { Worker } from 'bullmq';
const worker = new Worker('my-queue', async (job) => {
console.log('处理:', job.data);
await doSomething(job.data);
});

// 4. 添加任务并监听
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 的四大核心类是什么? QueueWorkerQueueEventsFlowProducer
QueueWorker 的区别? Queue 是生产者(创建任务),Worker 是消费者(执行任务)。
QueueEventsWorker 的事件监听有什么区别? Worker.on() 只能监听该 Worker 处理的任务;QueueEvents 监听队列中所有任务,全局视角。
什么时候用 FlowProducer 需要处理有父子依赖关系的任务流,比如“生成报告”需要等“拉取数据”和“处理数据”都完成。
concurrency: 5 是什么意思? 这个 Worker 同时最多处理 5 个任务。

总结: 你现在的笔记已经覆盖了 Redis 基础语法 + BullMQ 四大核心类,足够应对实际开发中的任务队列场景。接下来的学习路线可以按这个顺序走:

  1. 本地跑通最小示例(上面的 4 步曲)
  2. 实战一个场景:发邮件/生成报表/图片处理
  3. 学习高级特性:重试策略、速率限制、任务去重

需要我陪你跑通第一个完整示例,还是直接上实战场景?随时喊我!😄