BullMQ 连接配置全指南:ioredis / node-redis / Bun / Valkey Glide 适配器与 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点击查看免费下载本指南以 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并附带一个指数退避的retryStrategyMath.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含runCommandduplicate()返回的是再次被包装的 Proxy而不是裸的重复客户端其余一切属性——事件、options、ioredis 特有的方法——都直接转发给你传入的底层实例该底层实例永远不会被修改。此外ioredis-client.ts 中用一个WeakMap缓存了原始实例 → Proxy的映射重复调用createIORedisClient会返回同一个 Proxy保证事件监听器身份一致。如果你直接把原生 node-redis 或 Bun 客户端传给connectionredis-connection.ts 中的wrapRedisInstance会做结构化探测node-redis 有sendCommandisOpen/isReadyBun 有sendconnected并自动套上对应适配器因此这些用户甚至可以不安装 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注册并通过EVALSHANOSCRIPT 时回退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.tsclose()与quit()/disconnect()语义不同send(command, args)是调用任意 Redis 命令的通用通道原生duplicate()是异步的因此适配器用延迟 materialize的方式同步返回一个带rawFactory的新适配器bun-redis-client.ts脚本执行走EVALSHANOSCRIPT 时回退EVALbun-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 GlideValkey 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。工厂接收合并后的连接选项并必须返回一个已适配的IRedisClientimport { 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]元组超时返回nullredis-client.tsnode-redis/Bun 等适配器都必须把原生返回值转换到这个元组形式这样共享 ioredis 实例的用户代码不会因为返回形状被改动而破坏。maxRetriesPerRequest 深度解析maxRetriesPerRequest告诉 ioredis 客户端一条命令在抛错之前最多重试多少次。即使 Redis 当前不可达或离线命令也会持续重试直到连接恢复或达到最大尝试次数。对 Worker 而言这保证了只要存在一条可用连接Worker 就会一直处理下去命令无限重试。手动创建 ioredis 客户端时如果你把该客户端传给 WorkerBullMQ 会在maxRetriesPerRequest未设为null时抛出异常。对应源码见 redis-connection.ts 的checkBlockingOptions当blocking为 trueWorker/QueueEvents 场景且选项中带有非空maxRetriesPerRequest时要么抛错传入已构造实例时要么打印覆盖警告而在blocking分支中BullMQ 会直接强制this.opts.maxRetriesPerRequest nullredis-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 支撑 BullMQimport { 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 QueueMyData, 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-policynoeviction请确保你的 Redis 实例配置了maxmemory-policynoeviction否则 Redis 可能自动淘汰evictBullMQ 的键造成意想不到的错误。这一检查同样内置于源码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.0minimumVersion并强烈建议 6.2.0recommendedMinimumVersion低于下限会抛错低于建议值会打印警告。连接就绪后还会根据版本计算能力位如canDoubleTimeout需要 6.0.0、canBlockFor1Ms需要 7.0.8见 redis-connection.ts。最后如果连接数对你不是问题那就大胆地用——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点击查看免费下载相关推荐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上一篇Deep Learning for Coders 2020零基础入门AI的终极指南下一篇SumatraPDF 终极指南从快速入门到高效使用的 10 个技巧创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考