- 后端
- 消息队列
- 任务调度
【免费下载链接】bullmq
BullMQ - Message Queue and Batch processing for NodeJS, Python, .NET, Elixir, Rust and PHP based on Redis or PostgreSQL
本指南以 BullMQ 官方文档 connections.md 为主体,结合当前仓库源码(src/classes/、src/interfaces/)深入讲解 BullMQ 的 Redis 连接模型:如何为Queue/Worker等类配置与复用连接、如何通过IRedisClient适配器接入 node-redis、Bun、Valkey Glide 等不同客户端、如何全局替换客户端工厂,以及maxRetriesPerRequest、keyPrefix、maxmemory-policy等关键参数的底层行为。读完本文,你将能根据生产者/消费者场景正确选择连接策略,并在多客户端环境中无障碍地接入 BullMQ。
连接模型总览:谁需要连接、如何复用
在 BullMQ 中,任何队列操作都离不开一条到 Redis 实例的连接。默认情况下,BullMQ 使用 ioredis 中构造连接时填充了port: 6379、host: '127.0.0.1',并附带一个指数退避的retryStrategy(Math.max(Math.min(Math.exp(times), 20000), 1000))。
连接使用有几个要点:
- 每个类至少消费一条连接:
Queue、Worker、QueueEvents、FlowProducer等类都会创建或复用连接。 - 支持连接复用:
Queue和Worker都接受一个已经构造好的(适配后的)Redis 客户端实例作为connection选项。 - 阻塞命令需要额外连接:
Worker与QueueEvents依赖阻塞式 Redis 命令(如BZPOPMIN、XREAD BLOCK),因此它们会在内部创建一条重复连接(duplicate)。这意味着你传入的客户端或适配器必须支持duplicate()。这一点在 IRedisClient 接口 中即为强制要求:duplicate(...args: any[]): IRedisClient。
换句话说:你可以把同一条 ioredis 连接传给两个Queue(生产者),也可以传给两个Worker(消费者),但 Worker 内部会为阻塞命令再开辟一条专用连接,并不会阻塞你共享的那条主连接。
每个实例独立创建连接
最简单的做法是让每个实例各自创建连接,连接选项直接写在connection对象里:
import { Queue, Worker } from 'bullmq'; // Create a new connection in every instance const myQueue = new Queue('myqueue', { connection: { host: 'myredis.taskforce.run', port: 32856, }, }); const myWorker = new Worker('myqueue', async job => {}, { connection: { host: 'myredis.taskforce.run', port: 32856, }, });这里connection选项支持host、port、db、username、password、tls、retryStrategy、maxRetriesPerRequest等字段,完整的选项类型定义见 redis-options.ts。
复用 ioredis 连接
当你的服务对连接数有硬性限制,或希望减少握手开销时,可以预先创建一个 ioredis 实例并在多个实例间共享:
import { Queue } from 'bullmq'; import IORedis from 'ioredis'; const connection = new IORedis(); // Reuse the ioredis instance in 2 different producers const myFirstQueue = new Queue('myFirstQueue', { connection }); const mySecondQueue = new Queue('mySecondQueue', { connection });import { Worker } from 'bullmq'; import IORedis from 'ioredis'; const connection = new IORedis({ maxRetriesPerRequest: null }); // Reuse the ioredis instance in 2 different consumers const myFirstWorker = new Worker('myFirstWorker', async job => {}, { connection, }); const mySecondWorker = new Worker('mySecondWorker', async job => {}, { connection, });注意第三个示例:虽然 ioredis 实例被两个 Worker 复用,但每个 Worker 仍会通过duplicate()在内部创建一条它自己需要的阻塞连接。同时注意 Worker 场景要求maxRetriesPerRequest: null(原因见后文"maxRetriesPerRequest 深度解析"一节)。
兼容性说明:透明代理包装
为了向后兼容,BullMQ 仍然接受原生的IORedis实例作为connection,尽管内部已统一走IRedisClient适配器接口。传入了原生 ioredis 实例时,它会被包装进一个透明 Proxy(实现见 ioredis-client.ts),该 Proxy 只做几件事:
- 新增
runCommand,用于按名称分发 Lua 脚本(defineCommand注册的脚本); - 为
hset、set、zrange、zrevrange、xadd、xread、xtrim、scan提供结构化选项形式的调用(原生 ioredis 的变参形式依然可用,Proxy 根据参数形状自动分发); pipeline()/multi()返回被增强过的事务对象(IRedisTransaction,含runCommand);duplicate()返回的是再次被包装的 Proxy,而不是裸的重复客户端;- 其余一切属性——事件、options、ioredis 特有的方法——都直接转发给你传入的底层实例,该底层实例永远不会被修改。
此外,ioredis-client.ts 中用一个WeakMap缓存了"原始实例 → Proxy"的映射,重复调用createIORedisClient会返回同一个 Proxy,保证事件监听器身份一致。如果你直接把原生 node-redis 或 Bun 客户端传给connection,redis-connection.ts 中的wrapRedisInstance会做结构化探测(node-redis 有sendCommand+isOpen/isReady,Bun 有send+connected)并自动套上对应适配器,因此这些用户甚至可以不安装 ioredis。
使用 node-redis 客户端
BullMQ不会直接替你创建 node-redis 客户端:你需要在应用里创建原生客户端,再用createNodeRedisClient包装后传给 BullMQ。
版本要求:使用 BullMQ 的 node-redis 适配器时,请安装
redisv5 或更新版本——BullMQ 为这个适配器声明了redis >= 5.0.0的 peer dependency。
import { Queue, Worker, createNodeRedisClient } from 'bullmq'; import { createClient } from 'redis'; const rawClient = createClient({ url: 'redis://localhost:6379', }); const connection = createNodeRedisClient(rawClient); const myQueue = new Queue('myqueue', { connection }); const myWorker = new Worker('myqueue', async job => {}, { connection });从源码看,createNodeRedisClient 返回一个完整的NodeRedisAdapter包装类(而非就地打补丁),因为它与 ioredis 的 API 结构差异太大。该适配器负责把status归一化为'wait'/'ready'/'end'语义(node-redis-client.ts),用 SHA1 对 Lua 脚本做defineCommand注册并通过EVALSHA(NOSCRIPT 时回退EVAL)执行(node-redis-client.ts),同时把 node-redis 的返回结构(如zRangeWithScores、hscan的{cursor, entries})转换成 ioredis 风格的扁平结构。duplicate()也会把已注册脚本复制给新适配器实例。
使用 Bun 内置 Redis 客户端
Bun 自带 Redis 客户端。同样地,BullMQ 不会替你实例化它:先创建原生 Bun 客户端,再用createBunRedisClient包装:
import { RedisClient } from 'bun'; import { Queue, Worker, createBunRedisClient } from 'bullmq'; const rawClient = new RedisClient('redis://localhost:6379'); const connection = createBunRedisClient(rawClient); const myQueue = new Queue('myqueue', { connection }); const myWorker = new Worker('myqueue', async job => {}, { connection });RedisClient类由 Bun 运行时提供,这段代码请用bun run ...运行,而不是纯 Node.js 环境。
BunRedisAdapter 的实现与 node-redis 适配器有几点显著差异,值得了解:
- Bun 的客户端没有 EventEmitter,而是用
onconnect/onclose回调,适配器负责把它们桥接成标准事件(bun-redis-client.ts); close()与quit()/disconnect()语义不同,send(command, args)是调用任意 Redis 命令的通用通道;- 原生
duplicate()是异步的,因此适配器用"延迟 materialize"的方式同步返回一个带rawFactory的新适配器(bun-redis-client.ts); - 脚本执行走
EVALSHA,NOSCRIPT 时回退EVAL(bun-redis-client.ts)。
Bun 连接的优雅关闭
当你把同一个包装连接共享给多个 Queue 和 Worker 时,关闭顺序很重要:
// Graceful shutdown await myWorker.close(); await myQueue.close(); connection.disconnect(); // or: await connection.quit();请通过包装器返回的连接(即createBunRedisClient的返回值)执行disconnect()或quit(),不要直接调用原生 BunRedisClient的close()。原因是:包装器无法把这种关闭标记为"有意为之",于是 in-flight 命令会以ConnectionClosedError被拒绝,包装器还会尝试重连。经由包装器关闭则能干净地排空这些命令,行为与 ioredis 的quit()一致。源码中这一点由 sendCommand 的关闭判定(this.closing || this.closed时吞掉连接关闭错误)保证。
使用 Valkey Glide
Valkey Glide 的 API 与 ioredis/node-redis 差异较大,需要用createValkeyGlideClient包装:
import { GlideClusterClient } from '@valkey/valkey-glide'; import { Queue, Worker, createValkeyGlideClient } from 'bullmq'; const rawClient = await GlideClusterClient.createClient({ addresses: [{ host: 'localhost', port: 6379 }], }); const connection = createValkeyGlideClient(rawClient); const myQueue = new Queue('myqueue', { connection }); const myWorker = new Worker('myqueue', async job => {}, { connection });Valkey Glide 适配器(valkey-glide-client.ts)同样实现了IRedisClient的全部契约,包括通过customCommand执行 Lua 脚本、把 Glide 的 Map/键值数组回复归一化为 ioredis 风格结构,以及使用 GlideBatch(原子或非原子)承载multi()/pipeline()语义。注意此适配器当前通过isCluster = false标记为单节点模式(Glide 集群客户端的能力边界以源码注释为准)。
全局替换客户端工厂:RedisConnection.clientFactory
如果你希望 BullMQ在需要新建连接时一律使用非 ioredis 客户端,可以在应用启动阶段设置RedisConnection.clientFactory。工厂接收合并后的连接选项,并必须返回一个已适配的IRedisClient:
import { Queue, RedisConnection, createNodeRedisClient } from 'bullmq'; import { createClient } from 'redis'; RedisConnection.clientFactory = opts => { const rawClient = createClient({ socket: { host: opts.host, port: opts.port, }, username: opts.username, password: opts.password, database: opts.db, }); return createNodeRedisClient(rawClient); }; const myQueue = new Queue('myqueue', { connection: { host: 'myredis.taskforce.run', port: 32856, }, });Bun 客户端同理:
import { RedisClient } from 'bun'; import { Queue, RedisConnection, createBunRedisClient } from 'bullmq'; RedisConnection.clientFactory = opts => { const host = opts?.host ?? 'localhost'; const port = opts?.port ?? 6379; const rawClient = new RedisClient(`redis://${host}:${port}`); return createBunRedisClient(rawClient); }; const myQueue = new Queue('myqueue', { connection: { host: 'myredis.taskforce.run', port: 32856, }, });这一机制在 redis-connection.ts 的init()中生效:当调用方既没有传入已构造的客户端实例,又设置了clientFactory时,就用工厂创建客户端;否则才回退到懒加载 ioredis。懒加载意味着使用 node-redis/Bun/PostgreSQL 后端的用户永远不会加载 ioredis 包(redis-connection.ts)。另外,clientFactory要求返回"已增强"(即已被createNodeRedisClient等包装过)的IRedisClient,这一点由 redis-connection.ts 的类型注释明确说明。
自定义 Redis 客户端:实现 IRedisClient 接口
任何 Redis 客户端,只要适配到 BullMQ 的IRedisClient接口就能使用。适配器需要暴露 BullMQ 用到的 Redis 命令、连接生命周期方法、事件、duplicate()、通过defineCommand()注册 Lua 脚本,以及通过multi()/pipeline()提供的流水线/事务能力。完整接口定义见 redis-client.ts,其中命令仅覆盖 BullMQ 实际使用的那部分(Hash/String/ZSet/List/Set/Stream/阻塞命令/服务管理/扫描),方法签名统一采用结构化选项对象而非 ioredis 风格的变参,这样每个适配器都能映射到自己的原生 API 而无需解析位置参数。
对大多数应用,优先使用内置适配器:
createIORedisClient:适用于 ioredis 的Redis与Cluster实例;createNodeRedisClient:适用于 node-redis 客户端;createBunRedisClient:适用于 Bun 内置 Redis 客户端;createValkeyGlideClient:适用于 Valkey Glide 客户端。
编写自定义适配器时需注意:bzpopmin的返回形状必须与 ioredis 原生一致——成功时是[key, member, score]元组,超时返回null(redis-client.ts),node-redis/Bun 等适配器都必须把原生返回值转换到这个元组形式,这样共享 ioredis 实例的用户代码不会因为返回形状被改动而破坏。
maxRetriesPerRequest 深度解析
maxRetriesPerRequest告诉 ioredis 客户端:一条命令在抛错之前最多重试多少次。即使 Redis 当前不可达或离线,命令也会持续重试,直到连接恢复或达到最大尝试次数。
- 对 Worker 而言:这保证了只要存在一条可用连接,Worker 就会一直处理下去(命令无限重试)。
- 手动创建 ioredis 客户端时:如果你把该客户端传给 Worker,BullMQ 会在
maxRetriesPerRequest未设为null时抛出异常。对应源码见 redis-connection.ts 的checkBlockingOptions:当blocking为 true(Worker/QueueEvents 场景)且选项中带有非空maxRetriesPerRequest时,要么抛错(传入已构造实例时),要么打印覆盖警告;而在blocking分支中,BullMQ 会直接强制this.opts.maxRetriesPerRequest = null(redis-connection.ts)。 - 使用其他客户端适配器时:请按照该客户端自身的文档配置重试与重连行为,让 Worker 连接能够持续重试。
相关行为在 tests/connection.test.ts 中有直接测试:blocking为 true 时maxRetriesPerRequest被置为null,为 false 时保留原值(如 10)。
底层架构:IQueueBackend 后端抽象
理解了IRedisClient之后,还需要知道 BullMQ 的连接体系还有更高一层抽象。IRedisClient抽象的是底层驱动(ioredis、node-redis、Bun、Glide),而Queue、Worker、FlowProducer、QueueEvents等高层类再往上坐一层:它们是数据存储无关的,只与实现了IQueueBackend契约的后端对话。后端持有连接,并实现所有队列操作:"add job"、"move to active"、"extend lock"、阻塞式"wait for next job"等。该接口刻意不暴露任何连接或事务类型——具体适配器自己拥有连接(例如 Redis 后端由一个提供IRedisClient的上下文构建,并为waitForJob配备专用阻塞客户端,见 queue-backend.ts)。
默认后端是 Redis 后端(RedisQueueBackend),因此日常使用中你完全不需要直接接触这一层——像本文前面那样传一个connection即可,BullMQ 会自动接好 Redis 后端(工厂为createRedisBackend,见 create-backend.ts)。
访问当前后端与后端专用客户端
高层类不再暴露clientgetter。getBackend()返回实际使用的后端(Redis、PostgreSQL 或自定义后端)。使用默认 Redis 后端时,你仍能按需拿到底层 Redis 客户端:
import { Queue, RedisClient } from 'bullmq'; const queue = new Queue('myqueue', { connection: { host: 'localhost', port: 6379 }, }); // By default BullMQ uses the Redis backend, so getBackend() exposes // Redis-specific escape hatches. const client: RedisClient = await queue.getBackend().client; await client.set('some-key', 'some-value');Redis 后端还暴露其他 Redis 特有细节,如redisVersion、databaseType和底层connection。对Worker而言,getBackend().blockingClient返回专用阻塞连接的客户端——即阻塞式"等待任务"原语所用的那条连接。
注意:尽量优先使用高层的
Queue/Worker/FlowProducerAPI。任何经由getBackend()触达的后端特有内容都不在数据存储无关契约之内,不同后端之间可能不一致。
注入自定义后端
所有高层类只依赖IQueueBackend接口,并接收一个构建它的后端工厂。默认工厂是createRedisBackend,但你可以把自定义工厂作为最后一个构造参数传入,用不同的数据存储或测试 Mock 支撑 BullMQ:
import { Queue, BackendFactory } from 'bullmq'; const myBackendFactory: BackendFactory = (name, opts, options) => { // return an object implementing IQueueBackend }; const queue = new Queue('myqueue', { connection: {} }, myBackendFactory);这些类是泛型于后端类型的,因此getBackend()会返回你提供的工厂产出的具体类型(默认是RedisQueueBackend)。非 Redis 用户可以这样写:new Queue<MyData, MyResult, string, MyBackend>(name, opts, createMyBackend)。BackendFactory的类型签名(含blocking/withBlockingConnection选项)见 queue-backend.ts。
警告:构建一个生产级后端是相当大的工程——你必须以正确的原子性、锁、时序和事件语义实现完整的
IQueueBackend契约。在把后端认定为可上线之前,请用 adapter-conformance 测试和完整的 BullMQ 测试套件来验证行为(仓库根目录的tests/adapter-conformance.test.ts即是此类验证的入口之一)。
内置 PostgreSQL 后端
BullMQ 自带一个现成的PostgreSQL 后端(createPostgresBackend),让完整的Queue/Worker/QueueEvents/FlowProducerAPI 运行在 PostgreSQL 之上而非 Redis。需求、连接选项、schema 与迁移说明见专门页面:PostgreSQL 后端指南。其 SQL 命令与迁移文件可在仓库的 src/postgres/commands/ 与 src/postgres/migrations/ 目录中查阅。
生产建议:Queue 与 Worker 的不同连接策略
需要特别留意:仅用于管理队列的简单Queue实例(添加任务、暂停、getters 等)与 Worker 的连接需求通常不同。
生产者场景:假设你通过一个 HTTP 端点往队列里加任务。调用方不能因为 Redis 恰好宕机而无限等待,因此maxRetriesPerRequest应保留默认值(当前为 20),或设置成一个较小的值,比如 1,让用户快速拿到错误、稍后重试。
消费者场景:如果你在 Worker 的处理器里加任务(后台进程),则可以共享同一条连接。
关于连接持久化的更多细节可参考仓库中 docs/gitbook/bull/patterns/ 目录下的模式文档(如手动重试、节流等主题),它们对生产环境下的连接与任务生命周期管理有进一步说明。
两个必须避免的陷阱
1. 不要使用 ioredis 的keyPrefix
使用 ioredis 连接时,切勿启用keyPrefix选项——它与 BullMQ 不兼容。BullMQ 有自己的键前缀机制,通过prefix选项实现(默认值为bull,见 queue-options.ts)。在 redis-connection.ts 中,如果你传入的客户端实例带有keyPrefix,会直接抛出错误:"BullMQ: ioredis does not support ioredis prefixes, use the prefix option instead."
2. Redis 必须设置maxmemory-policy=noeviction
请确保你的 Redis 实例配置了:
maxmemory-policy=noeviction否则 Redis 可能自动淘汰(evict)BullMQ 的键,造成意想不到的错误。这一检查同样内置于源码:init()阶段读取INFO时,redis-connection.ts 会解析maxmemory_policy字段,若不为noeviction则打印警告IMPORTANT! Eviction policy is ... It should be "noeviction"。测试 tests/connection.test.ts 对此有专门用例:当把策略改为volatile-lru时触发警告,改回noeviction后消除。
此外,redis-connection.ts 定义了版本基线:BullMQ 要求 Redis 版本>= 5.0.0(minimumVersion),并强烈建议>= 6.2.0(recommendedMinimumVersion);低于下限会抛错,低于建议值会打印警告。连接就绪后还会根据版本计算能力位(如canDoubleTimeout需要 6.0.0+、canBlockFor1Ms需要 7.0.8+,见 redis-connection.ts)。
最后,如果连接数对你不是问题,那就大胆地用——Redis 连接的开销很低,除非服务商施加了硬性限制,否则通常不需要刻意复用连接。
- 后端
- 消息队列
- 任务调度
【免费下载链接】bullmq
BullMQ - Message Queue and Batch processing for NodeJS, Python, .NET, Elixir, Rust and PHP based on Redis or PostgreSQL
相关推荐
Node-Redis连接池监控终极指南:连接数、空闲连接与超时配置完全解析
Node Redis连接池监控终极指南:连接数、空闲连接与超时配置完全解析 Node Redis作为Redis官方推荐的Node.js客户端,在高性能应用开发中
后端数据库客户端缓存Node-Redis连接池优化:管理高并发Redis连接的终极指南
Node Redis连接池优化:管理高并发Redis连接的终极指南 在Node.js应用中处理高并发Redis访问时, node redis连接池 是提升性能和
后端数据库客户端缓存Node-Steam-Guide数据库集成:使用MongoDB存储Steam物品数据完整指南
Node Steam Guide数据库集成:使用MongoDB存储Steam物品数据完整指南 Node Steam Guide是一个使用Node.js创建Ste
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考