本文汇总 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()
Worker
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 时,它们会竞争消费同一个队列中的任务。正常情况下,一条 Job 由其中一个 Worker 领取,而不是广播给所有 Worker。
BullMQ 有类似 RabbitMQ ACK 的语义,但普通 Worker 会自动完成确认。
| RabbitMQ | BullMQ |
|---|---|
| 收到消息 | Worker 领取 Job |
ack() | 处理函数正常结束 |
nack() | 处理函数抛出异常 |
| 未 ACK 后连接中断 | Job 锁过期并成为 stalled |
正常执行完成,即使没有显式 returncompleted
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()JobwaitUntilFinished()
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()QueueEventsqueueEvents.waitUntilReady()waitUntilFinished()QueueEventsqueueEvents.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条
}
}
如果只记录错误但不继续 throwcompleted
BullMQ 内置两种 backoff(任务执行失败后,失败任务的重试方式。达到重试标准后,改任务仍然会入队,按照优先级priority来执行):
| 类型 | 含义 |
|---|---|
fixed | 每次失败后等待固定时间 |
exponential | 每次失败后等待时间指数增加 |
固定退避:
jsbackoff: {
type: 'fixed',
delay: 5000,
}
指数退避:
jsbackoff: {
type: 'exponential',
delay: 1000,
jitter: 0.3,
}
指数退避公式:
markdown 2 ^ (attemptsMade - 1) × delay
jitter01
还可以在 Worker 的 settings.backoffStrategy
jssettings: {
backoffStrategy: (attemptsMade, type, error, job) => {
if (type === 'business-api') {
return attemptsMade * 5000;
}
throw new Error(`Unknown backoff type: ${type}`);
},
}
自定义策略返回:
0-1faileddelay
jsawait queue.add(
'close-unpaid-order',
{ orderId: 'ORDER-1001' },
{ delay: 30 * 60 * 1000 },
);
markdown 入队 → delayed → 到期 → waiting/prioritized → active
delay
区别:
markdown delay = 第一次执行前等待多久 backoff = 执行失败后,重试前等待多久
未设置 priority
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 的任务优先。
未设置 prioritypriority: 0job.opts.priorityundefined
普通无优先级任务会先于正数优先级任务。因此,一旦业务采用优先级,同类任务最好统一设置正数优先级。
delaypriority
markdown delay → 什么时候获得排队资格(表示当前任务入队后,多长时间后才能被消费者取到。在没到指定时间时,任务状态为等待,到了时间后,会按照当前队列任务中的优先级priority来执行BullMQ是按照优先级来执行任务的。) priority → 获得资格后排在什么位置
优先级只影响 Worker 获取任务的顺序,不会抢占已经处于 active
如果 Worker 的 concurrency > 1
默认 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
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 可以通过 SharedArrayBufferAtomicsuseWorkerThreadsworkerThreadsOptions.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()推荐给 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 配置:
attemptsbackoffjobIddeduplicationdelaypriorityremoveOnCompleteremoveOnFail自动删除 Job 后,相同 jobIdjobId
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,
},
);
生产建议:
concurrencycompletedfailedstallederrorSIGTERMSIGINTworker.close()prefixRedis 生产环境还应:
maxmemory-policy noevictionFlowProducer
普通 Queue.add()
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
生产环境应明确子任务最终失败时如何影响父任务:
failParentOnFailureignoreDependencyOnFailureremoveDependencyOnFailurecontinueParentOnFailuremarkdown Queue.add = 添加一个独立 Job Worker = 常驻消费者,领取并执行 Job QueueEvents = 跨进程监听任务状态 waitUntilFinished = 通过 QueueEvents 等待 Job 的最终结果 delay = 什么时候获得第一次执行资格 priority = 获得资格后排在什么位置 backoff = 失败后多久再次获得执行资格 FlowProducer = 创建具有父子依赖关系的任务树
BullMQ 的核心价值不是简单地把数组放进 Redis,而是围绕 Redis 提供任务状态、原子领取、锁、自动续锁、失败重试、延迟、优先级、结果保存、并发控制和依赖编排。