- 后端
- 消息队列
- 任务调度
【免费下载链接】bullmq
BullMQ - Message Queue and Batch processing for NodeJS, Python, .NET, Elixir, Rust and PHP based on Redis or PostgreSQL
本篇快速入门指南以当前仓库 BullMQ(v6.x,基于 Redis 的消息队列与批处理库)为背景,带你从零搭建第一条可运行的队列链路:安装依赖、向队列投递任务、用 Worker 进程消费任务,并通过本地事件与全局QueueEvents监听任务全生命周期。阅读完成后,你将掌握 BullMQ 生产-消费模型的最小闭环,并能立刻在本地项目中复制运行。
安装 BullMQ
BullMQ 同时支持 npm 与 yarn 两种包管理器,在项目根目录执行其一即可:
$ npm install bullmq$ yarn add bullmq安装后,BullMQ 会以库的形式提供Queue、Worker、QueueEvents、Job等核心类。从当前仓库的 package.json 可以看到,BullMQ 使用 TypeScript 编写("source": "./src/index.ts"),对外同时发布 CJS 与 ESM 产物(main指向./dist/cjs/index.js,module指向./dist/esm/index.js),并随包附带完整的类型声明(types指向./dist/esm/index.d.ts)。运行环境要求 Node.js>= 14.17.0,其核心运行时依赖仅有cron-parser、msgpackr等少量库,而ioredis以可选 peer dependency 的形式声明(>= 5.0.0),因此你需要自行安装所选用的 Redis 客户端。
关于语言选择,官方文档提示:BullMQ 本身由 TypeScript 编写,虽然可以直接在原生 JavaScript 中使用,但本篇指南的示例统一采用 TypeScript 编写,以便充分利用类型推导与 IDE 提示。
前置条件:本地 Redis 服务
运行下方所有示例前,你必须在本地启动一个 Redis 服务。BullMQ 默认通过 ioredis 连接localhost:6379,若你使用自定义地址,可在Queue/Worker的构造选项中显式传入connection配置。更完整的连接方式(包括复用连接、node-redis 适配、Bun 内置 Redis 客户端适配等)可参考 连接指南。
创建队列并投递任务
导入Queue类并实例化一个名为foo的队列,随后即可通过add方法向队列投递任务。任务由「任务名称」与「数据载荷」两部分组成:
import { Queue } from 'bullmq'; const myQueue = new Queue('foo'); async function addJobs() { await myQueue.add('myJobName', { foo: 'bar' }); await myQueue.add('myJobName', { qux: 'baz' }); } await addJobs();从源码结构看,Queue.add的定义位于 src/classes/queue.ts,其签名为add(name: NameType, data: DataType, opts?: JobsOptions):第一个参数是任务名称(用于区分同一队列中的不同任务类型),第二个参数是任务数据(需为 JSON 可序列化的普通对象),第三个可选参数是任务选项(如重试次数、延迟、优先级等)。add会返回一个Job实例,其中包含系统分配的任务 ID。
上述代码执行后,两个任务即被写入 Redis 中对应foo队列的等待集合,等待任意 Worker 进程拾取。仓库测试 tests/queue.test.ts 中大量使用await queue.add(queueName, { foo: 'bar', bar: 1 })的模式验证「投递后可通过queue.getJob(job.id)重新读取任务数据」,这与本文示例的用法完全一致,可将其作为可运行的最小验证场景。
用 Worker 消费任务
任务进入队列后可以随时被处理,只要至少有一个 Node.js 进程在运行 Worker。Worker 通过阻塞式轮询从队列中取出任务并交给处理器(processor)执行:
import { Worker } from 'bullmq'; import IORedis from 'ioredis'; const connection = new IORedis({ maxRetriesPerRequest: null }); const worker = new Worker( 'foo', async job => { // 第一个任务会打印 { foo: 'bar'} // 第二个任务会打印 { qux: 'baz' } console.log(job.data); }, { connection }, );这里有几个关键点需要说明:
maxRetriesPerRequest: null是必须的:Worker 内部需要使用阻塞式 Redis 命令来等待新任务(默认阻塞上限为 10 秒,见 src/classes/worker.ts 中关于BZPOPMIN的注释)。若不加此配置,ioredis 默认的重试上限会在阻塞期间抛出异常。- 任务数据直接可见:处理器收到的
job参数是完整的Job实例,job.data即投递时写入的数据载荷。 - 多 Worker 水平扩展:你可以同时运行任意数量的 Worker 进程(甚至分布在多台机器上),BullMQ 会以轮询(round robin)方式将任务均衡地分发给各个 Worker,天然实现并行消费与水平扩容。
Worker 类的监听器接口定义在 src/classes/worker.ts,可以看到它作为EventEmitter提供completed、failed、active、progress、drained、stalled等事件,其中completed回调携带(job, result, prev),failed回调携带(job | undefined, error, prev)——注意当任务因removeOnFail被删除时,job可能为undefined,代码中需做好判空。
监听任务完成与失败
Worker 自带本地事件监听,可以直观地感知每个任务的执行结果:
worker.on('completed', job => { console.log(`${job.id} has completed!`); }); worker.on('failed', (job, err) => { console.log(`${job.id} has failed with ${err.message}`); });需要明确的是:Worker上的completed/failed事件属于进程本地事件,只有真正处理了该任务的 Worker 进程内才能监听到。BullMQ 还提供了大量其他事件,完整的分类说明见 事件指南。
用 QueueEvents 实现全局事件监听
在许多场景下,你希望在一个统一的位置监听所有 Worker 发出的事件(例如构建实时看板、WebSocket 推送)。为此 BullMQ 提供了专门的QueueEvents类:
import { QueueEvents } from 'bullmq'; const queueEvents = new QueueEvents('my-queue-name'); queueEvents.on('waiting', ({ jobId }) => { console.log(`A job with ID ${jobId} is waiting`); }); queueEvents.on('active', ({ jobId, prev }) => { console.log(`Job ${jobId} is now active; previous status was ${prev}`); }); queueEvents.on('completed', ({ jobId, returnvalue }) => { console.log(`${jobId} has completed and returned ${returnvalue}`); }); queueEvents.on('failed', ({ jobId, failedReason }) => { console.log(`${jobId} has failed with reason ${failedReason}`); });QueueEvents的内部实现基于 Redis Streams(事件写入以QUEUE_EVENT_SUFFIX命名的流中,相关常量定义于 src/utils 的引用中)。相比传统的 pub/sub,流式实现具备两个重要特性:
- 事件不丢失:断线重连期间产生的事件仍可被恢复消费,不会像 pub/sub 那样直接丢失。
- 自动裁剪:事件流默认保留约 10,000 条事件以防无限膨胀,可通过
streams.events.maxLen选项调整。
此外,每个事件回调还会附带第二个参数——事件的时间戳标识,其形态类似"1580456039332-0",可用于事件排序与去重:
import { QueueEvents } from 'bullmq'; const queueEvents = new QueueEvents('my-queue-name'); queueEvents.on('progress', ({ jobId, data }, timestamp) => { console.log(`${jobId} reported progress ${data} at ${timestamp}`); });QueueEvents 事件不携带 Job 实例
出于性能考虑,QueueEvents发出的事件只包含jobId字符串,而不会携带完整的Job实例。如果你需要获取任务对象,应使用Job.fromId静态方法按 ID 重新加载:
import { Job } from 'bullmq'; const job = await Job.fromId(queue, jobId);该方法定义于 src/classes/job.ts,接收(queue, jobId)两个参数,从队列后端按 ID 读取任务数据并重建Job实例。这是「全局事件 + 按需拉取任务详情」的标准搭配:事件流保持轻量,需要完整数据时再精准查询。
小结:一条完整的最小链路
将上述代码串联,你就拥有了一个完整的最小 BullMQ 链路:
npm install bullmq安装依赖,并确保本地 Redis 可用;new Queue('foo')创建队列,queue.add('myJobName', data)投递任务;new Worker('foo', processor, { connection })启动消费者,任务以 round robin 方式分发给所有 Worker;- Worker 本地监听
completed/failed事件感知单个进程内的执行结果; - 需要全局视角时,用
QueueEvents基于 Redis Streams 监听waiting/active/progress/completed/failed等事件,并结合Job.fromId按需加载任务详情。
在此基础上,你可以进一步阅读 事件指南 了解全部可用事件与流裁剪配置,或参考 连接指南 学习连接复用与多种 Redis 客户端适配方式,从而将这条最小链路扩展为生产可用的任务处理系统。
- 后端
- 消息队列
- 任务调度
【免费下载链接】bullmq
BullMQ - Message Queue and Batch processing for NodeJS, Python, .NET, Elixir, Rust and PHP based on Redis or PostgreSQL
相关推荐
Plate 项目 shadcn 风格组件编写规范:从语义内核到 Open UI 的三层架构与提取测试
Plate 项目 shadcn 风格组件编写规范:从语义内核到 Open UI 的三层架构与提取测试 Plate 是一个基于 shadcn/ui 构建的富文本编
后端消息队列任务调度如何快速上手bAbI-tasks:从安装到生成AI问答任务的完整指南
如何快速上手bAbI tasks:从安装到生成AI问答任务的完整指南 bAbI tasks是一个由Facebook AI Research开发的经典问答任务生成
Pipenv 快速上手指南:从安装到生产环境的完整实践
Pipenv 快速上手指南:从安装到生产环境的完整实践 本篇指南以 Pipenv(Python Development Workflow for Humans)
开发工具CLI包管理器
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考