后端消息队列任务调度【免费下载链接】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点击查看免费下载本文围绕本项目文档库中保留的经典模式 Message Queue讲解如何把 Bull/BullMQ 当作持久化消息队列使用add即发送process即接收。文章先完整还原原文档的双服务器示例再给出 BullMQ 的现代等价实现并结合仓库源码queue.ts、worker.ts剖析消息从投递到处理的完整生命周期最后补充可靠投递、事件监听、去重、定时、集群等进阶能力帮助你构建真正可落地的跨服务通信方案。一、模式概述为什么用队列做服务间通信原文档开篇就点明了这个模式的核心价值Bull 也可以用作持久化消息队列persistent message queue这在某些场景下非常实用。典型的例子是两台服务器需要互相通信通过队列两台服务器不需要同时在线从而构成一条非常健壮的通信信道robust communication channel可以把add视为send发送把process视为receive接收。之所以持久化是因为任务消息会被完整写入 Redis本项目还支持 PostgreSQL 后端而不是像内存队列那样随进程消亡。发送方只要把消息写入队列即可立即返回接收方无论何时上线都能从队列中取到尚未被消费的消息。这为系统解耦、削峰填谷、故障恢复提供了基础。二、原版 Bull 双服务器示例原文档完整代码原文档给出了一个非常直观的双向通信示例每台服务器各持有一个发送队列和一个接收队列接收队列注册处理器发送队列投递消息。Server A服务 Aconst Queue require(bull); const sendQueue new Queue(Server B); const receiveQueue new Queue(Server A); receiveQueue.process(function (job, done) { console.log(Received message, job.data.msg); done(); }); sendQueue.add({ msg: Hello });Server B服务 Bconst Queue require(bull); const sendQueue new Queue(Server A); const receiveQueue new Queue(Server B); receiveQueue.process(function (job, done) { console.log(Received message, job.data.msg); done(); }); sendQueue.add({ msg: World });需要说明的是这段示例来自本项目docs/gitbook/bull/目录下保留的Bull上一代库文档使用的是require(bull)、回调式done的旧 API。本项目当前主库是BullMQ下面给出功能完全对等的现代写法。三、在 BullMQ 中的等价实现Queue WorkerBullMQ 把队列与处理者拆分为两个类Queue负责add发送/投递以及队列管理暂停、清空、删除等Worker负责process接收/消费实例化并连上 Redis 后即开始取任务处理。等价的双服务器示例TypeScript// server-a.ts —— 服务 A import { Queue, Worker } from bullmq; // 发送消息给服务 B const sendQueue new Queue(Server B); // 接收来自服务 B 的消息Worker 即消息接收者 const receiveWorker new Worker(Server A, async job { console.log(Received message, job.data.msg); }); // 发送一条消息返回 PromiseJob await sendQueue.add(message, { msg: Hello }); // 进程退出前优雅关闭等待正在处理的任务结束 await receiveWorker.close();// server-b.ts —— 服务 B import { Queue, Worker } from bullmq; const sendQueue new Queue(Server A); const receiveWorker new Worker(Server B, async job { console.log(Received message, job.data.msg); }); await sendQueue.add(message, { msg: World });使用 BullMQ 时需要注意的 API 差异均以本仓库源码为准add的签名不同BullMQ 的Queue.add(name, data, opts?)第一个参数是任务名见 queue.ts而旧版 Bull 是add(data, opts)。上例中的message就是任务名。处理器在Worker上而非Queue上BullMQ 中process的职责由Worker承担。从源码看Worker一旦实例化并连上 Redis 就会开始处理任务worker.ts且必须显式提供connection配置否则构造时会直接抛出Worker requires a connectionworker.ts。回调改为 Promise/asyncBullMQ 的处理器支持async函数直接return结果不再需要done()回调处理器返回值会通过completed事件对外发布worker.ts。数据必须是 JSON 可序列化的add传入的data会经过JSON.stringify后存入后端类实例的 prototype 方法与 getter 不会随消息传给消费者queue.ts。跨进程传消息时请使用纯对象或实现toJSON()。四、消息从发送到接收的生命周期源码视角要理解消息队列为何可靠可以顺着源码看一次完整投递投递调用queue.add(name, data, opts)后Queue会把默认任务选项defaultJobOptions与本次opts合并queue.ts通过Job.create创建任务并写入后端随后触发waiting事件queue.ts。任务此刻已经持久化发送方无需等待接收方在线。拉取Worker的主循环mainLoop持续从等待集合中取任务取到后进入active状态worker.ts任务被锁保护防止被多个 worker 重复处理。处理processJob调用处理器函数成功后handleCompleted调用job.moveToCompleted(result, token, ...)将任务移入completed状态并触发completed事件worker.ts处理器抛异常则进入handleFailed任务自动移入failed状态并触发failed事件worker.ts。结束只有当任务被成功处理并从活动状态移出后它才不再是待消费消息。对应到官方指南Workers 文档 明确写道Worker 等价于传统消息队列中的消息接收者任务成功则进入completed处理抛异常则自动进入failed。这条状态机正是持久化消息队列的骨架消息要么成功完成、要么失败留待处理而不会在投递与消费之间凭空丢失。五、持久化与可靠投递语义基于上面的生命周期可以总结出该模式提供的可靠性保障持久化存储消息写入 Redis或本项目支持的 PostgreSQL 后端见 postgresql.md 与后端抽象接口 queue-backend.ts发送方与接收方无需同时在线这是原文档强调的健壮通信信道的根本原因。失败可重试处理失败的 job 进入failed状态可通过 Retrying failing jobs 配置的自动重试策略attempts、backoff重新投递。防卡死机制stalled如果 worker 因 CPU 繁忙或进程崩溃无法及时续锁任务会被判定为 stalled 并移回等待队列由其他 worker 重新处理workers.md 的 Stalled jobs 一节。Worker的默认配置为lockDuration: 30000、stalledInterval: 30000、maxStalledCount: 1worker.ts。至少一次投递从源码结构看消息在完成前始终保留在队列的数据结构中处理完成后才从活动集合移除可以推断该模式具备至少一次投递语义如果业务对消息有幂等要求可配合去重能力见下文使用。六、事件驱动监听消息到达与处理结果在消息队列场景中除了消费消息本身往往还需要感知谁发送了消息处理结果如何。BullMQ 提供了两个层面的事件1. Worker 本地事件workers.mdreceiveWorker.on(completed, (job, returnvalue) { // 消息处理完成returnvalue 为处理器返回值 }); receiveWorker.on(failed, (job, error) { // 消息处理失败 }); receiveWorker.on(progress, (job, progress) { // 处理器内调用 job.updateProgress() 时触发 });此外还有active任务开始处理、drained等待列表被清空等事件完整列表见 worker.ts。2. 全局事件QueueEvents当接收方是独立服务、发送方或第三方服务想监听任意 worker 的处理结果时使用QueueEventsimport { QueueEvents } from bullmq; const queueEvents new QueueEvents(Server A); queueEvents.on(completed, ({ jobId, returnvalue }) { // 任何 worker 完成该队列的 job 时触发 }); queueEvents.on(failed, ({ jobId, failedReason }) { // 任何 worker 处理失败时触发 });同时发送方queue.add()成功后也会触发队列的waiting事件queue.ts可用于实时感知新消息已到达。七、组合模式任务队列 结果消息队列本仓库另一篇经典文档 Returning Job Completions 展示了消息队列模式的典型组合用法一个处理集群只负责高速消费任务另一个服务负责保存结果。最健壮、可扩展的实现方式就是把标准任务队列与消息队列模式结合起来服务端通过打开一个任务队列并add任务把工作分发给集群集群中的每个 job 完成时向一个结果消息队列发送一条包含结果数据的消息负责落库的服务监听这个结果队列把结果写入数据库。这样任务消费与结果处理完全解耦两边都可以独立伸缩这正是消息队列即通信信道模式的进阶体现。八、进阶消息能力让发消息更贴合业务在原文档发送/接收这一最小模型之上仓库提供了丰富的消息级能力可按需组合定时 / 周期消息add时传入opts.delay实现延迟投递需要周期发送时使用queue.upsertJobScheduler(jobSchedulerId, repeatOpts, jobTemplate)queue.ts支持 cron 表达式Bull 时代的写法见 quick-guide.md 的 Repeated jobs 一节。消息去重通过去重键deduplication key避免同一逻辑消息被重复投递与处理见 deduplication.mdQueue也提供removeDeduplicationKey(id)手动清理queue.ts。优先级opts.priority可让高优消息优先被消费。暂停 / 恢复queue.pause()/queue.resume()全局暂停整个队列底层通过原子 RENAME 把等待队列改为 paused见 queue.tsworker.pause()则只暂停当前 worker。速率限制与全局并发queue.setGlobalRateLimit(max, duration)与queue.setGlobalConcurrency(n)可控制队列整体的消费节奏queue.ts。沙箱化处理器把处理函数放到独立进程/线程中运行崩溃不影响 worker 主进程、可运行阻塞代码并更好利用多核见 quick-guide.md 的 Sandboxed processes 一节以及 worker.ts。集群并行消费多个进程/线程可以安全地并发消费同一个队列队列不会被破坏Bull 时代的cluster示例见 quick-guide.md 的 Cluster support 一节。这让一个队列、多个接收者的高吞吐消费成为可能。九、连接与部署注意事项在真实服务间通信场景中还需要留意连接管理复用 Redis 连接每个 Queue/Worker 实例默认都会建立新的 Redis 连接多队列场景建议复用连接详见 Reusing Redis Connections。持久连接消息队列通常要求进程长驻相关最佳实践见 Persistent connections。Redis Cluster在生产环境使用 Redis 集群时需要按集群要求配置连接见 Redis Cluster。后端选型本项目基于统一的IQueueBackend抽象同时提供 Redis 与 PostgreSQL 两套实现redis-queue-backend.ts、postgres-queue-backend.ts可按基础设施偏好选择。队列数量队列本身相对廉价需要多个通信信道时可直接用不同名字创建但队列过多会难以管理quick-guide.md 建议一般控制在十几个以内。另外从仓库结构看python/、dotnet/、elixir/、rust/、php/目录本项目还提供多种语言的 BullMQ 客户端实现。因此用队列做服务间通信天然可以跨语言例如 Node.js 服务向队列发送消息由 Python 或 .NET 的服务消费处理只需保证使用同一个队列名与同一套 Redis/PostgreSQL 后端即可。十、总结把 BullMQ 当作持久化消息队列使用时核心心智模型就是原文档那句精炼的总结add即发送process即接收。借助 Redis/PostgreSQL 的持久化存储、waiting → active → completed/failed的严谨状态机、stalled 防卡死机制和失败重试能力两端服务无需同时在线即可获得一条健壮的异步通信信道。在此基础上通过去重、延迟、定时、优先级、速率限制、全局事件与多语言客户端你可以把它扩展成覆盖任务分发、结果回传、跨语言解耦的完整消息通信方案。本文所有代码与源码引用均可在当前仓库对应路径中找到可对照 queue.ts 与 worker.ts 进一步深入阅读。赞分享后端消息队列任务调度【免费下载链接】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点击查看免费下载相关推荐3步搞定Figma中文界面设计师必备的高效汉化插件终极指南3步搞定Figma中文界面设计师必备的高效汉化插件终极指南 还在为Figma的全英文界面感到困扰吗作为全球顶尖的设计工具Figma的专业界面却让许多中文设后端消息队列任务调度3分钟免费美化foobar2000终极DUI配置让你的音乐播放器焕然一新3分钟免费美化foobar2000终极DUI配置让你的音乐播放器焕然一新 还在忍受foobar2000那个单调乏味的默认界面吗想让你的音乐播放器既有专业功能桌面应用音视频asyncpg与消息队列消息持久化asyncpg与消息队列消息持久化 你是否曾因系统崩溃丢失重要消息而烦恼作为运营人员或开发者消息丢失可能导致业务中断、数据不一致甚至客户投诉。本文将展示数据库后端上一篇Sunshine游戏串流终极指南3步搭建你的跨平台游戏共享网络下一篇Gatsby 私有源下的 CA 证书配置指南npm、yarn 与 Node 三端实战创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
