本文汇总 BullMQ 学习过程中涉及的主要知识点,包括核心架构、任务执行、失败重试、排序规则、Worker Thread、生产配置以及 FlowProducer。
BullMQ 是一个基于 Redis、主要面向 Node.js/TypeScript 的后台任务队列库。
它更关注:
典型场景包括发送邮件、生成报表、处理图片、调用 Webhook、定时任务和异步数据处理。
| 产品 | 核心定位 | 典型用途 |
|---|---|---|
| BullMQ | Redis 后台任务队列 | Node.js 异步任务、重试、延迟和任务编排 |
| RabbitMQ | 通用消息代理 | 跨语言通信、Exchange 和复杂路由 |
| Kafka | 分布式事件流平台 | 高吞吐事件、长期保存和消息回放 |
| SQS | 云托管队列 | AWS 中无需自行维护消息基础设施 |
| Redis Pub/Sub | 实时广播 | 不需要持久化和 ACK 的即时通知 |
| Redis Streams | Redis 原生消息流 | 自行构建消费者组和消息处理逻辑 |
BullMQ 更像“应用任务队列”,不是 RabbitMQ 或 Kafka 的直接替代品。
markdown 生产者 │ Queue.add() ▼ Redis │ ▼ Worker │ ├── completed ├── failed └── retry
Queue 用于添加和管理任务:
jsconst queue = new Queue('email-queue', { connection });
await queue.add('send-email', {
email: 'user@example.com',
});
Queue.add() 不执行任务,只把 Job 数据和状态写入 Redis。
Worker 是常驻的任务消费者(这里的Worker不是Threads Worker,而是一个消费者)。
也可以让消费任务采用threads worker 来执行,BullMQ对齐进行了封装,实现起来很简单:
jsconst worker = new Worker(
'email-queue',
async (job) => {
await sendEmail(job.data.email);
return {
delivered: true,
};
},
{ connection },
);
Worker 会持续连接 Redis、领取任务、执行处理函数并更新任务状态。它默认不是临时创建的工作线程,而是一个逻辑上的消费者。
QueueEvents 用于跨进程监听队列中的全局状态事件:
jsconst queueEvents = new QueueEvents('email-queue', {
connection,
});
queueEvents.on('completed', ({ jobId }) => {
console.log(`Job ${jobId} completed`);
});
生产者可以连续添加多个 Job:
jsawait queue.add('send-email', { email: 'a@example.com' });
await queue.add('send-email', { email: 'b@example.com' });
await queue.add('send-email', { email: 'c@example.com' });
这些 Job 会保存到 Redis,Worker 持续领取并执行。
如果 Worker 尚未启动,任务会停留在 waiting。Worker 以后启动时仍可继续处理。
启动多个 Worker 时,它们会竞争消费同一个队列中的任务。正常情况下,一条 Job 由其中一个 Worker 领取,而不是广播给所有 Worker。
BullMQ 有类似 RabbitMQ ACK 的语义,但普通 Worker 会自动完成确认。
| RabbitMQ | BullMQ |
|---|---|
| 收到消息 | Worker 领取 Job |
ack() | 处理函数正常结束 |
nack() | 处理函数抛出异常 |
| 未 ACK 后连接中断 | Job 锁过期并成为 stalled |
正常执行完成,即使没有显式 return,也会进入 completed:
jsnew Worker('email-queue', async (job) => {
await sendEmail(job.data.email);
// 隐式返回 undefined,仍然视为成功
});
return 的作用是保存处理结果,不是决定是否成功。
必须正确等待异步操作:
js// 错误:BullMQ 可能在邮件发送完成前把 Job 标记为 completed
sendEmail(job.data.email);
// 正确
await sendEmail(job.data.email);
queue.add() 返回一个 Job 对象。生产者可以通过 waitUntilFinished() 等待 Worker 的返回值:
jsconst queueEvents = new QueueEvents('email-queue', {
connection,
});
await queueEvents.waitUntilReady();
const job = await queue.add('send-email', {
email: 'user@example.com',
});
const result = await job.waitUntilFinished(
queueEvents,
10_000,
);
Worker 返回:
jsreturn {
delivered: true,
email: job.data.email,
};
执行链路:
markdownWorker return
↓
结果写入 Redis,Job 变为 completed
↓
QueueEvents 收到 completed
↓
waitUntilFinished() 返回结果
注意事项:
waitUntilFinished() 必须配合相同队列的 QueueEvents。queueEvents.waitUntilReady()。waitUntilFinished() 会抛出异常。QueueEvents 可以复用来等待同一队列中的多个 Job。queueEvents.close()。Worker 抛出异常时,BullMQ 根据 attempts 判断是否重试:
jsawait queue.add(
'send-email',
{ email: 'user@example.com' },
{
attempts: 3,
backoff: {
type: 'exponential',
delay: 1000,
},
},
);
attempts: 3 表示最多执行三次,包含第一次执行。
markdown 执行失败 ├── 还有次数 → delayed/waiting → 再次执行 └── 次数耗尽 → failed
最终失败任务默认不会被删除:
jsconst failedJobs = await queue.getFailed();
for (const job of failedJobs) {
console.log({
id: job.id,
data: job.data,
failedReason: job.failedReason,
attemptsMade: job.attemptsMade,
stacktrace: job.stacktrace,
});
}
失败任务可以人工重试:
jsawait job.retry();
不要吞掉异常:
jstry {
await sendEmail(job.data.email);
} catch (error) {
console.error(error);
throw error;
}
其它配置策略
js // 配置失败任务不重试,直接从队列中删除:
await queue.add(
'send-email',
{ email: 'user@example.com' },
{
removeOnFail: true,
},
);
// 或者保留最近的1000条失败任务
{
removeOnFail: 1000
}
// 按照数量和时间一起来清理
{
removeOnFail: {
age: 24 * 60 * 60, // 保留一天,单位为秒
count: 1000, // 最多保留1000条
}
}
如果只记录错误但不继续 throw,处理函数会正常结束,BullMQ 会把 Job 标记为 completed。
BullMQ 内置两种 backoff(任务执行失败后,失败任务的重试方式。达到重试标准后,改任务仍然会入队,按照优先级priority来执行):
| 类型 | 含义 |
|---|---|
fixed | 每次失败后等待固定时间 |
exponential | 每次失败后等待时间指数增加 |
固定退避:
jsbackoff: {
type: 'fixed',
delay: 5000,
}
指数退避:
jsbackoff: {
type: 'exponential',
delay: 1000,
jitter: 0.3,
}
指数退避公式:
markdown 2 ^ (attemptsMade - 1) × delay
jitter 为 0 到 1,用于给等待时间增加随机性,避免大量失败任务同时重试。
还可以在 Worker 的 settings.backoffStrategy 中定义自定义策略:
jssettings: {
backoffStrategy: (attemptsMade, type, error, job) => {
if (type === 'business-api') {
return attemptsMade * 5000;
}
throw new Error(`Unknown backoff type: ${type}`);
},
}
自定义策略返回:
0:立即重新进入等待队列。-1:停止重试并进入 failed。delay 表示 Job 至少等待多久才具备第一次执行资格:
jsawait queue.add(
'close-unpaid-order',
{ orderId: 'ORDER-1001' },
{ delay: 30 * 60 * 1000 },
);
markdown 入队 → delayed → 到期 → waiting/prioritized → active
delay 不保证准时执行,只保证不会早于指定时间执行。如果 Worker 繁忙或不可用,实际执行时间会更晚。
区别:
markdown delay = 第一次执行前等待多久 backoff = 执行失败后,重试前等待多久
未设置 priority 时,任务默认按照 FIFO 领取:
markdown Job 1 → Job 2 → Job 3
设置 lifo: true 后使用后进先出:
markdown Job 3 → Job 2 → Job 1
正数 priority 越小,优先级越高:
markdown priority 1 → 高 priority 5 → 中 priority 10 → 低
没有设置 priority 的普通任务,比设置了正数 priority 的任务优先。
未设置 priority 可以在调度概念上视为 priority: 0,但 job.opts.priority 可能仍是 undefined。
普通无优先级任务会先于正数优先级任务。因此,一旦业务采用优先级,同类任务最好统一设置正数优先级。
delay 与 priority 的关系:
markdown delay → 什么时候获得排队资格(表示当前任务入队后,多长时间后才能被消费者取到。在没到指定时间时,任务状态为等待,到了时间后,会按照当前队列任务中的优先级priority来执行BullMQ是按照优先级来执行任务的。) priority → 获得资格后排在什么位置
优先级只影响 Worker 获取任务的顺序,不会抢占已经处于 active 的任务。
如果 Worker 的 concurrency > 1 或有多个 Worker,任务完成顺序不一定等于获取顺序。
默认 FIFO、单并发且没有显式优先级时,立即重试或 backoff 到期后的 Job 会重新进入等待队列,通常排在当前等待任务之后。
markdown Job 1 失败 Job 2 → Job 3 → ... → Job 10 Job 1 重试
重试不会自动提高 Job 的优先级;如果 Job 原本设置了 priority,重试时会保留该优先级。
普通写法默认在 Worker 所在 Node.js 进程的主线程运行:
jsnew Worker('email-queue', async (job) => {
await sendEmail(job.data.email);
}, {
concurrency: 20,
});
concurrency: 20 表示最多同时处理 20 个任务,不表示创建了 20 个线程。它主要利用 Node.js 异步 I/O,适合邮件、HTTP 和数据库操作。
CPU 密集任务可以使用外部 processor 和 Worker Thread:
jsnew Worker(
'image-queue',
processorFile,
{
connection,
useWorkerThreads: true, // 开启此选项,让消费者函数在真正的worker中执行
},
);
BullMQ 会在内部封装:
postMessage() 的线程通信。生产者并不直接与处理线程通信:
markdown 生产者 → Redis → BullMQ Worker 主线程 → Worker Thread Worker Thread → Worker 主线程 → Redis → QueueEvents/生产者
处理器的返回值最终写入 Redis,因此生产者和 Worker 可以位于不同进程或服务器。
原生 Node.js Worker Thread 通常需要显式通信:
js// 主线程
const worker = new Worker('./processor.js');
worker.postMessage(jobData);
worker.on('message', (result) => {
console.log(result);
});
js// processor.js
const { parentPort } = require('node:worker_threads');
parentPort.on('message', async (jobData) => {
const result = await processJob(jobData);
parentPort.postMessage(result);
});
注册了 parentPort.on('message') 的线程会持续等待新消息,不会在单次任务完成后自动退出。可以由主线程调用:
jsawait worker.terminate();
也可以让线程调用 parentPort.close(),在没有其他活跃资源时自然退出。
Node.js Worker Thread 可以通过 SharedArrayBuffer 和 Atomics 共享内存。BullMQ 使用 useWorkerThreads 时,也可以通过 workerThreadsOptions.workerData 把共享内存交给其创建的线程:
jsconst sharedBuffer = new SharedArrayBuffer(
Int32Array.BYTES_PER_ELEMENT,
);
new Worker('thread-queue', processorFile, {
connection,
useWorkerThreads: true,
concurrency: 4,
workerThreadsOptions: {
workerData: {
sharedBuffer,
},
},
});
在线程处理器中:
jsconst { workerData } = require('node:worker_threads');
const counter = new Int32Array(workerData.sharedBuffer);
const value = Atomics.add(counter, 0, 1) + 1;
共享范围仅限同一个 Node.js 进程创建的 Worker Thread:
queue.add() 把共享内存引用传给远程 Worker。推荐给 Queue 设置统一默认策略:
jsconst queue = new Queue('email-queue', {
connection: {
host: process.env.REDIS_HOST,
port: Number(process.env.REDIS_PORT),
password: process.env.REDIS_PASSWORD,
maxRetriesPerRequest: 1,
connectTimeout: 5000,
},
prefix: 'myapp',
defaultJobOptions: {
attempts: 4,
backoff: {
type: 'exponential',
delay: 1000,
jitter: 0.3,
},
removeOnComplete: {
age: 24 * 3600,
count: 10_000,
},
removeOnFail: {
age: 7 * 24 * 3600,
count: 50_000,
},
},
});
常用 Job 配置:
attempts:最大执行次数。backoff:重试等待策略。jobId:业务唯一 ID,防止现存 Job 重复入队。deduplication:一定时间内去重或节流。delay:第一次执行前延迟。priority:执行资格相同时的优先级。removeOnComplete:成功任务保留策略。removeOnFail:失败任务保留策略。自动删除 Job 后,相同 jobId 可以再次入队,因此 jobId 不能代替业务幂等设计。
jsconst worker = new Worker(
'email-queue',
async (job) => {
return await sendEmail(job.data);
},
{
connection: {
host: process.env.REDIS_HOST,
port: Number(process.env.REDIS_PORT),
password: process.env.REDIS_PASSWORD,
maxRetriesPerRequest: null,
},
prefix: 'myapp',
concurrency: 20,
limiter: {
max: 100,
duration: 1000,
},
lockDuration: 30_000,
stalledInterval: 30_000,
maxStalledCount: 1,
},
);
生产建议:
concurrency。completed、failed、stalled 和 error。SIGTERM 或 SIGINT 时调用 worker.close()。prefix 必须与 Queue、QueueEvents 和 FlowProducer 保持一致。Redis 生产环境还应:
maxmemory-policy noeviction。FlowProducer 用于原子地创建具有父子依赖关系的任务树。
普通 Queue.add() 创建独立任务;FlowProducer 创建依赖关系:
markdown 父任务
/ \
子任务 A 子任务 B
执行顺序是:
markdown子任务 A、B 执行
↓
全部成功完成
↓
父任务 waiting-children → waiting
↓
父任务 Worker 执行
子任务和父任务本质上都是普通 Job:
父任务读取直属子任务结果:
jsconst childrenValues = await job.getChildrenValues();
const results = Object.values(childrenValues);
父任务默认只能直接获取直属子任务结果。多层 Flow 应由中间节点汇总后逐层向上返回。
markdown最终任务 A
├── 子任务 B
│ ├── 子任务 D
│ └── 子任务 E
└── 子任务 C
└── 子任务 F
执行顺序从叶子向根节点:
markdown D、E → B F → C B、C → A
并列子任务可以设置 priority,但优先级只影响领取顺序,不能保证一个任务完成后另一个才开始。
如果要求:
markdown 子任务 2 完成 → 子任务 1 开始 → 最终父任务开始
应该用嵌套依赖:
jsawait flowProducer.add({
name: 'final-parent',
queueName: 'parent-queue',
children: [
{
name: 'child-1',
queueName: 'child-queue',
children: [
{
name: 'child-2',
queueName: 'child-queue',
},
],
},
],
});
Flow 会按照:
markdown child-2 → child-1 → final-parent
生产环境应明确子任务最终失败时如何影响父任务:
failParentOnFailure:让父任务失败。ignoreDependencyOnFailure:记录失败,但允许父任务继续。removeDependencyOnFailure:移除失败子任务依赖。continueParentOnFailure:子任务失败后立即让父任务处理失败情况。markdown Queue.add = 添加一个独立 Job Worker = 常驻消费者,领取并执行 Job QueueEvents = 跨进程监听任务状态 waitUntilFinished = 通过 QueueEvents 等待 Job 的最终结果 delay = 什么时候获得第一次执行资格 priority = 获得资格后排在什么位置 backoff = 失败后多久再次获得执行资格 FlowProducer = 创建具有父子依赖关系的任务树
BullMQ 的核心价值不是简单地把数组放进 Redis,而是围绕 Redis 提供任务状态、原子领取、锁、自动续锁、失败重试、延迟、优先级、结果保存、并发控制和依赖编排。