
后端消息队列任务调度【免费下载链接】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点击查看免费下载全局并发Global Concurrency是 BullMQ 提供的一项队列级Queue 级配置它从队列维度统一约束所有 Worker 实例的总并行处理数量是分布式环境下实现整队列限流的关键手段。本文基于 docs/gitbook/guide/queues/global-concurrency.md 展开并结合仓库源码queue.ts、queue-getters.ts、isQueueMaxed.lua剖析其底层实现帮助你在多 Worker、多实例部署中精确控制同一队列同时被处理的作业数理解它与 Worker 级concurrency选项的区别并掌握设置、读取与移除的完整用法。什么是全局并发因子全局并发因子是一个队列选项它决定在所有 Worker 实例之间同一时刻最多允许有多少个 Job 被并行处理。举个直观的例子假设你有 3 台服务器每台运行 1 个 Worker 实例每个实例内部又配置了concurrency: 4。在没有全局并发限制时整个系统最多可以同时处理 3 × 4 12 个 Job。而一旦给队列设置globalConcurrency 2那么无论有多少个 Worker、每个 Worker 内部并发多高全系统同一时刻最多只有 2 个 Job 处于处理中状态。从源码看该值的语义定义在 src/interfaces/queue-meta.ts/** * Maximum number of jobs that can be processed concurrently across all * workers attached to this queue. Set via Queue.setGlobalConcurrency. */ concurrency?: number;这段注释明确说明concurrency字段代表连接到该队列的所有 Worker 可并发处理的作业总数上限只能通过Queue.setGlobalConcurrency设置。它是队列元数据queue meta的一部分与全局速率限制max/duration、暂停状态paused等一同存放在队列的meta哈希键中。设置全局并发在 Node.js/TypeScript 中通过Queue实例调用setGlobalConcurrency即可启用并设置全局并发值import { Queue } from bullmq; const queue new Queue(my-queue); // 全局并发为 4所有 Worker 合计最多同时处理 4 个 Job await queue.setGlobalConcurrency(4);方法签名与参数说明见 src/classes/queue.ts方法async setGlobalConcurrency(concurrency: number): Promisenumber参数concurrency—— 所有 Worker 可同时处理的作业最大数量。例如设置为1即可保证同一时刻最多只有一个 Job 被处理实现整队列串行如果该值未定义即从未设置或已被移除则对并发数量没有任何限制。底层实现写入队列 meta 哈希setGlobalConcurrency本身并不直接操作 Redis 客户端而是委托给队列后端backend写入元数据async setGlobalConcurrency(concurrency: number) { return this.backend.setQueueMeta({ concurrency }); }在 Redis 后端中src/classes/redis-queue-backend.tssetQueueMeta最终执行的是对队列 meta 键的HSET操作async setQueueMeta(values: Recordstring, string | number): Promisenumber { const client await this.queue.client; return client.hset(this.queue.keys.meta, values); }也就是说await queue.setGlobalConcurrency(4)等价于向bull:my-queue:meta这个哈希键写入concurrency 4。由于写入的是队列元数据而非某个 Worker 的本地状态因此任何 Worker 实例都能感知到该限制——这正是全局二字的来源。与其他后端的关系BullMQ 当前支持基于 Redis 与 PostgreSQL 的多种后端。从源码结构看src/classes/queue.ts 统一调用this.backend.setQueueMetasetGlobalConcurrency是一个后端无关backend-agnostic的 API无论底层使用 redis-queue-backend.ts 还是 Postgres 后端src/postgres/postgres-queue-backend.tsQueue 层暴露给使用者的方法签名保持一致。本文后续以 Redis 实现为例展开。读取全局并发值通过getGlobalConcurrency可以读取当前队列已设置的全局并发值const globalConcurrency await queue.getGlobalConcurrency(); console.log(globalConcurrency); // 例如输出 4未设置时输出 null该方法的实现位于 src/classes/queue-getters.ts/** * Get global concurrency value. * Returns null in case no value is set. */ async getGlobalConcurrency(): Promisenumber | null { const concurrency await this.backend.getQueueMetaField(concurrency); if (concurrency) { return Number(concurrency); } return null; }需要注意两个关键细节返回类型是Promisenumber | null当队列从未设置过全局并发meta 哈希中没有concurrency字段时返回null而不是0或undefined字符串到数字的转换Redis 哈希中存储的字段值本质是字符串因此读取后需要通过Number(concurrency)显式转换。底层getQueueMetaField执行的是HGET bull:my-queue:meta concurrency见 src/classes/redis-queue-backend.ts。移除全局并发当不再需要全局限制时调用removeGlobalConcurrency即可移除await queue.removeGlobalConcurrency(); // 移除后再次读取将得到 null const value await queue.getGlobalConcurrency(); // null实现同样位于 src/classes/queue.ts/** * Remove global concurrency value. */ async removeGlobalConcurrency() { return this.backend.removeQueueMetaFields([concurrency]); }底层执行的是HDEL bull:my-queue:meta concurrency见 src/classes/redis-queue-backend.ts即从队列 meta 哈希中删除concurrency字段。移除之后队列恢复为不限制并发的默认行为——注意是恢复为无限制而不是恢复为某个默认数值。全局并发与 Worker 级 concurrency 的区别使用全局并发时最容易混淆的一点是它与 Worker 的concurrency选项之间的关系。原文档中的提示hint非常关键注意如果你在自己的 Worker 中设置了 concurrency 等级它不会覆盖全局并发值它只是单个 Worker 最多能并行处理的作业数但永远不会超过全局值。两者的定位可以这样理解配置项作用范围语义Worker 级concurrency如new Worker(queue, processor, { concurrency: 4 })单个 Worker 实例该 Worker 内部最多并行处理的作业数是上限全局setGlobalConcurrency(4)整个队列的所有 Worker所有 Worker 加在一起的并行处理总数上限是总闸门Worker 级 concurrency 决定单点能力即使没有全局限制单个 Worker 也不会超过自己配置的并发数全局并发决定集群总量它像一个全局闸门任何 Worker 在领取新任务前都要检查当前全队列的活跃数是否已到顶。因此最终同一时刻实际被并行处理的作业数为min(全局并发值, Σ 各 Worker 的本地 concurrency)设置全局并发并不会修改或覆盖任何 Worker 的本地配置两者是取最小值的叠加关系。底层原理Worker 如何看到全局并发全局并发之所以能在所有 Worker 之间生效是因为它的检查发生在领取作业claim阶段而非 Worker 本地调度阶段。核心逻辑在 Lua 脚本isQueueMaxedsrc/commands/includes/isQueueMaxed.lualocal function isQueueMaxed(queueMetaKey, activeKey) local maxConcurrency rcall(HGET, queueMetaKey, concurrency) if maxConcurrency then local activeCount rcall(LLEN, activeKey) if activeCount tonumber(maxConcurrency) then return true end end return false end这段脚本揭示了完整的判定流程从队列 meta 键中读取concurrency字段若该字段不存在说明未设置全局并发直接返回false不限制若已设置则统计当前活跃列表active list即正在被处理的作业列表的长度当活跃作业数全局并发上限时判定队列已满maxed返回true。从脚本可以看到该函数同时接收queueMetaKey和activeKey两个参数——活跃作业列表是队列级别的 Redis 数据结构所有 Worker 共享同一条列表。正因为如此无论有多少个 Worker 实例在消费该队列LLEN activeKey统计的都是全集群的真实处理中数量从而保证全局并发在多实例下严格生效。这一检查会被嵌入到作业领取类命令如moveToActive的执行链路中当队列被判为 maxed 时Worker 便无法再领取新的作业从而把整个集群的并行处理数量压制在全局上限之内。这同时也解释了为什么Worker 的本地 concurrency 永远不会超过全局值——因为超出的部分在领取阶段就被拦下了。完整实战示例下面给出一个端到端的示例演示从设置、验证到移除的完整流程import { Queue, Worker } from bullmq; const queue new Queue(tickets); // 1. 全局并发设为 2全系统同一时刻最多处理 2 个 Job await queue.setGlobalConcurrency(2); // 2. 读取并校验 const concurrency await queue.getGlobalConcurrency(); console.log(concurrency); // 2 // 3. 启动多个 Worker每个本地并发可以大于全局值 // 但实际并行数永远被全局值钳制在 2 const workerA new Worker(tickets, async job { /* 处理逻辑 */ }, { concurrency: 5 }); const workerB new Worker(tickets, async job { /* 处理逻辑 */ }, { concurrency: 5 }); // 4. 业务结束后移除全局并发限制 await queue.removeGlobalConcurrency(); console.log(await queue.getGlobalConcurrency()); // null典型应用场景串行化处理设置setGlobalConcurrency(1)强制整个队列同一时间只处理一个作业适用于对共享资源如单一数据库连接、有状态外部系统有严格串行要求的场景控制下游负载多个 Worker 横向扩展时用全局并发统一限制对下游 API/数据库的并发压力避免因扩容导致下游过载分批放量在业务高峰前调高全局并发、低谷时调低动态控制整条流水线的吞吐。小结全局并发是队列级选项通过Queue.setGlobalConcurrency(n)设置控制所有 Worker 实例合计的并行处理上限使用Queue.getGlobalConcurrency()读取当前值未设置时返回null使用Queue.removeGlobalConcurrency()移除限制恢复为无限制状态全局并发不覆盖 Worker 本地concurrency两者取最小值生效底层通过队列 meta 哈希RedisHSET/HGET/HDEL存储并在 Worker 领取作业的 Lua 脚本isQueueMaxed中按活跃列表长度实时判定是否已达上限从而保证多实例下的全局一致性。相关源码与文档队列 API 实现见 src/classes/queue.ts读取逻辑见 src/classes/queue-getters.tsRedis 元数据操作见 src/classes/redis-queue-backend.ts队列元数据字段定义见 src/interfaces/queue-meta.tsLua 判定脚本见 src/commands/includes/isQueueMaxed.lua。赞分享后端消息队列任务调度【免费下载链接】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点击查看免费下载相关推荐BullMQ Pro Groups 并发控制指南按组限制并行任务数Group ConcurrencyBullMQ Pro Groups 并发控制指南按组限制并行任务数Group Concurrency 导读 本文聚焦 BullMQ Pro taskf后端消息队列任务调度BullMQ Python 全局并发与全局速率限制实战指南BullMQ Python 全局并发与全局速率限制实战指南 全局并发Global Concurrency与全局速率限制Global Rate Limit后端消息队列任务调度BullMQ并行与并发机制深度解析BullMQ并行与并发机制深度解析 前言 在现代分布式系统中任务队列的性能优化是开发者必须掌握的技能。本文将深入探讨BullMQ任务队列系统中的并行 Para后端消息队列任务调度上一篇OpenCode与其他终端AI工具对比为什么开发者应该选择这款终极终端助手下一篇inshellisense与静态站点生成Next.js/Gatsby命令补全创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考