资讯动态

BullMQ 快速上手指南:从安装、生产任务到 Worker 消费与事件监听

发布时间:2026/9/25 11:36:29 来源:尧图企业网站定制
后端消息队列任务调度【免费下载链接】bullmqBullMQ - Message Queue and Batch processing for NodeJS, Python, .NET, Elixir, Rust and PHP based on Redis or PostgreSQL项目地址https://gitcode.com/gh_mirrors/bu/bullmq点击查看免费下载本篇快速入门指南以当前仓库 BullMQv6.x基于 Redis 的消息队列与批处理库为背景带你从零搭建第一条可运行的队列链路安装依赖、向队列投递任务、用 Worker 进程消费任务并通过本地事件与全局QueueEvents监听任务全生命周期。阅读完成后你将掌握 BullMQ 生产-消费模型的最小闭环并能立刻在本地项目中复制运行。安装 BullMQBullMQ 同时支持 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.jsmodule指向./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 方式分发给所有 WorkerWorker 本地监听completed/failed事件感知单个进程内的执行结果需要全局视角时用QueueEvents基于 Redis Streams 监听waiting/active/progress/completed/failed等事件并结合Job.fromId按需加载任务详情。在此基础上你可以进一步阅读 事件指南 了解全部可用事件与流裁剪配置或参考 连接指南 学习连接复用与多种 Redis 客户端适配方式从而将这条最小链路扩展为生产可用的任务处理系统。赞分享后端消息队列任务调度【免费下载链接】bullmqBullMQ - Message Queue and Batch processing for NodeJS, Python, .NET, Elixir, Rust and PHP based on Redis or PostgreSQL项目地址https://gitcode.com/gh_mirrors/bu/bullmq点击查看免费下载相关推荐Plate 项目 shadcn 风格组件编写规范从语义内核到 Open UI 的三层架构与提取测试Plate 项目 shadcn 风格组件编写规范从语义内核到 Open UI 的三层架构与提取测试 Plate 是一个基于 shadcn/ui 构建的富文本编后端消息队列任务调度如何快速上手bAbI-tasks从安装到生成AI问答任务的完整指南如何快速上手bAbI tasks从安装到生成AI问答任务的完整指南 bAbI tasks是一个由Facebook AI Research开发的经典问答任务生成Pipenv 快速上手指南从安装到生产环境的完整实践Pipenv 快速上手指南从安装到生产环境的完整实践 本篇指南以 PipenvPython Development Workflow for Humans开发工具CLI包管理器上一篇窗口管理终极革命如何用PinWin打破你的多任务效率瓶颈下一篇librsvg XInclude 任意文件读取漏洞CVE-2023-38633复现与原理分析创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑