- 后端
- 消息队列
- 任务调度
【免费下载链接】bullmq
BullMQ - Message Queue and Batch processing for NodeJS, Python, .NET, Elixir, Rust and PHP based on Redis or PostgreSQL
<output_article>
BullMQ for .NET 开发实战指南:跨 Redis 与 PostgreSQL 的统一分布式队列
BullMQ for .NET 是 BullMQ 官方维护的 .NET 移植版——一套基于 Redis 与 PostgreSQL 的快速、健壮的分布式任务队列。与 Node.js、Python、PHP、Elixir、Rust 等运行时共享同一套 Lua(Redis)与 SQL(PostgreSQL)脚本,保证跨语言的原子性与语义完全一致:Node.js 产出的队列可以被 .NET Worker 消费,反之亦然。读完本文,你将掌握 .NET 端队列/Worker 的完整接入方式、连接共享技巧、PostgreSQL 后端切换方法,以及底层的脚本共享机制与开发测试流程。
项目定位:glue code 之上的跨语言一致性
从 dotnet/README.md 的定义来看,这个 .NET 绑定本质上是胶水代码(glue code):它不重新实现队列语义,而是直接复用仓库根目录下rawScripts/中的 Lua 脚本,以及src/postgres/下的 SQL 脚本。共享脚本正是“跨语言完全一致”的保障——不同语言的运行时执行的是同一份经过实战检验的原子脚本:
- Redis 后端使用与 Node.js 主库 相同的 Lua 脚本(如
moveToActive-11.lua、addStandardJob-9.lua); - PostgreSQL 后端使用 src/postgres/commands 下的参数化 SQL 命令(如
add_job.sql、move_job_to_active.sql等)。
当前状态(README 声明,结合 dotnet/src/BullMQ 源码可验证):数据库无关的queue-backend 契约已同时为 Redis 与 PostgreSQL 完整实现,并有一套共享一致性测试套件(同一批测试同时跑在两个适配器上);高层的Queue、Worker、Job类在任一后端上均可运行,FlowProducer、QueueEvents、JobScheduler亦已实现。选择后端只需在选项中设置Postgres。
环境要求
- .NET 8.0 或更高版本;
- 使用 Redis 后端时:Redis 6.2+(或兼容服务器,如 Valkey、Dragonfly);
- 使用 PostgreSQL 后端时:PostgreSQL 13+(最低版本校验实现在 PostgresConnection.cs 中,
MinimumPostgresMajor = 13,启动时会读取server_version_num校验,可用SkipVersionCheck跳过)。
安装
通过 NuGet 一键引入:
dotnet add package BullMQ快速开始:队列 + Worker 最小可运行示例
以下示例完整来自 dotnet/README.md,演示了创建队列、入队、启动 Worker 处理并监听完成事件:
using BullMQ; // 创建队列并添加一个任务。 await using var queue = new Queue("emails", new QueueOptions { Connection = ConnectionOptions.FromString("localhost:6379"), }); await queue.AddAsync("welcome", new { to = "user@example.com" }); // 用 Worker 处理任务。 await using var worker = new Worker("emails", async (job, ct) => { Console.WriteLine($"Processing {job.Name} #{job.Id}"); // ... 执行业务逻辑 ... return "sent"; }, new WorkerOptions { Connection = ConnectionOptions.FromString("localhost:6379"), Concurrency = 5, }); worker.Completed += (job, result) => Console.WriteLine($"Job {job.Id} completed with '{result}'");几个关键点:
await using确保队列/Worker 随作用域退出被正确释放;Queue.AddAsync(name, data)的返回值为Job对象,处理器返回的对象会作为任务的完成结果持久化;- Worker 处理器签名是 Worker.cs 中定义的委托:
public delegate Task<object?> Processor(Job job, CancellationToken cancellationToken),ct用于协作式取消(如强制关闭); - Worker 事件(
Active、Completed、Failed、Drained、Stalled、Error、LockRenewalFailed、LocksRenewed等)定义在 Worker.cs,用于观察与可观测性集成。
Worker 的阻塞式取任务与锁机制
从 Worker.cs 的注释可以看到底层实现原理:任务获取使用后端的阻塞等待原语——Redis 端是对 marker 集合的BZPOPMIN,PostgreSQL 端是LISTEN;因此空闲 Worker不会忙轮询。锁由LockManager周期续期,后台还有一个 stalled-job 检查器负责恢复锁已过期的任务。
共享连接:复用 IConnectionMultiplexer
每个队列/Worker 都维护自己的连接,但可以通过传入已有的IConnectionMultiplexer在多个队列与 Worker 间复用连接。BullMQ 不会释放它不拥有的连接(README 明确说明),因此复用后由你自己管理其生命周期:
var mux = await StackExchange.Redis.ConnectionMultiplexer.ConnectAsync("localhost:6379"); var options = new QueueOptions { Connection = new ConnectionOptions { Multiplexer = mux }, };对应 Options.cs 中的ConnectionOptions:提供ConnectionString、Configuration(已解析的ConfigurationOptions)或Multiplexer三选一,Multiplexer优先级最高;另有Database(逻辑库索引,默认-1)与便捷工厂方法FromString(connectionString)。
使用 PostgreSQL 后端
在选项上设置Postgres,即可用完全相同的Queue/Worker/JobAPI 切换到 PostgreSQL。Schema(默认bullmq)作为所有队列的命名空间,所需表与函数会在首次使用时自动迁移创建:
using BullMQ; using BullMQ.Postgres; var pg = new PostgresOptions { ConnectionString = "Host=localhost;Database=bullmq;Username=postgres;Password=postgres", // Schema = "bullmq", // 可选,此为默认值 }; await using var queue = new Queue("emails", new QueueOptions { Postgres = pg }); await queue.AddAsync("welcome", new { to = "user@example.com" }); await using var worker = new Worker("emails", async (job, ct) => "sent", new WorkerOptions { Postgres = pg, Concurrency = 5 });PostgresOptions 与 Schema 语义
PostgresOptions.cs 定义了三个选项:
| 选项 | 说明 | 默认值 |
|---|---|---|
ConnectionString | Npgsql 连接字符串,如"Host=localhost;Database=bullmq_test" | 空字符串 |
Schema | 命名所有队列的 Schema,等价于 Redis 键前缀的 SQL 原生替代 | "bullmq" |
SkipVersionCheck | 跳过启动时的 PostgreSQL 最低版本检查 | false |
Schema 的妙处在于:它被钉在连接的search_path上(PostgresConnection.cs),因此共享的.sql文件全部使用不带 Schema 限定的可移植名称,既跨库可移植又天然防注入。Schema 名称必须匹配^[A-Za-z_][A-Za-z0-9_$]*$且不超过 63 字符(QuoteSchemaName)。
自动迁移的实现细节
PostgresConnection.cs 展示了迁移的健壮设计:
- 迁移使用事务级 advisory lock(键为
0x42554C4C,即字符串BULL,所有语言运行时使用完全相同的键)串行化跨进程并发启动,保证迁移恰好执行一次; - 迁移文件(见 src/postgres/migrations)按文件名中的数字前缀排序执行,记录在
bullmq_migration表中; - 迁移后调用
ReloadTypesAsync()重新加载 Npgsql 的类型目录(迁移会创建如 job-state 枚举等自定义类型); - 每个连接字符串在进程内只迁移一次(
s_migrated静态集合 + 静态s_migrateLock),避免多实例并发启动时触发 Npgsql 类型目录重载风暴。
阻塞等待:LISTEN/NOTIFY
PostgreSQL 后端使用专用LISTEN连接实现 Worker 的阻塞取任务(PostgresConnection.cs):订阅bullmq_jobs频道等待 NOTIFY,超时或关闭时通过链接的CancellationTokenSource优雅结束,不会阻塞进程退出。QueueEvents则订阅bullmq_events频道(JobsChannel/EventsChannel常量见 PostgresConnection.cs)。
核心选项速查:从源码注释出发的参数手册
Options.cs 是全部高层类共享的选项定义,以下参数均可在配置时直接使用:
QueueBaseOptions(Queue / Worker / QueueEvents 共用)
Connection:Redis 后端连接设置(默认后端);Prefix:Redis 键前缀,默认"bull"(与主库一致);Postgres:PostgreSQL 后端设置。一旦设置,队列/Worker 改用 PostgreSQL,且Connection被忽略。
QueueOptions
DefaultJobOptions:合并进该队列每个任务的默认JobsOptions;SkipMetasUpdate:为true时启动不写队列 meta hash。从 Queue.cs 可见默认写入opts.maxLenEvents = 10000与version = "bullmq:{Version.Value}"。
WorkerOptions
| 参数 | 默认值 | 说明 |
|---|---|---|
Concurrency | 1 | 并发处理的最大任务数(构造时校验必须>= 1) |
LockDuration | 30000 | 任务处理中的锁时长(毫秒) |
LockRenewTime | LockDuration / 2 | 锁续期频率(毫秒) |
DrainDelay | 5 | Worker 循环前阻塞等待任务的秒数 |
Name | null | 可读的 Worker 名称(用于可观测性) |
Autorun | true | 创建后是否立即开始处理 |
MaxStalledCount | 1 | 任务在判定失败前可从 stalled 恢复的最大次数 |
StalledInterval | 30000 | 执行 stalled 检查的间隔(毫秒) |
SkipStalledCheck | false | 禁用 stalled 检查器 |
SkipLockRenewal | false | 禁用周期性锁续期 |
在 Worker.cs 中可以确认:构造时若既无Postgres也无Connection会抛出ArgumentException;Concurrency < 1抛ArgumentOutOfRangeException;Autorun为true(默认)时立即启动后台消费循环。
QueueEventsOptions
Autorun(默认true):创建后立即开始消费事件;LastEventId:起始游标,默认"$"(只接收监听开始后的事件),可传入已知事件 id 恢复;BlockingTimeout(默认10000):每次对事件流的阻塞读超时(毫秒)。
JobsOptions(单任务配置)
JobId:显式任务 id;不能为"0"或以"0:"开头(Queue.cs 会在AddAsync时校验);Delay:任务变为可用前的延迟(毫秒);Priority:任务优先级(1最高;0或null表示无优先级);Attempts:任务完成前的总尝试次数;Backoff:自动重试的退避配置(BackoffOptions:Type如"fixed"/"exponential",Delay为毫秒基数);Lifo:为true时任务加入等待列表右侧(后进先出);Timestamp:创建时间戳(毫秒),默认当前时间;RemoveOnComplete/RemoveOnFail:完成后/失败后移除策略——true移除、数字表示保留条数、或KeepJobs策略对象;KeepLogs:任务保留的最大日志条数(0表示不限)。
KeepJobs支持Count(保留最大数量)与Age(保留最大秒数)两种裁剪维度。
高可用语义:Queue 的管理能力
Queue.cs 提供了与主库对齐的管理 API:
PauseAsync()/ResumeAsync()/IsPausedAsync():全局暂停/恢复队列处理;GetJobCountsAsync(params string[] types):按状态统计数量,未传参时默认统计waiting, active, completed, failed, delayed, paused六种状态;GetWaitingCountAsync():等待处理的任务数;GetJobAsync(jobId):按 id 获取任务;AddAsync(name, data, opts):入队并返回Job。
开发与测试:共享脚本的生成与验证流程
README 的 Development 一节说明了共享脚本的组织方式(dotnet/README.md):
- 仓库根目录
rawScripts/下的 Lua 脚本会被复制到dotnet/src/BullMQ/Commands/(嵌入程序集); src/postgres/**会被复制到dotnet/src/BullMQ/Postgres/;- 这些副本都是生成的、被 git 忽略的(SqlLoader.cs 通过程序集内嵌资源加载,命令与迁移均以内嵌
.sql资源形式提供,缺副本时会提示运行copy:sql:dotnet)。
在仓库根目录执行:
# 生成脚本副本(需先执行一次 `yarn install`)。 yarn copy:lua:dotnet yarn copy:sql:dotnet # 构建与测试。 cd dotnet dotnet build dotnet test # 需要 localhost:6379 上的 Redis 服务器或直接使用一键辅助脚本 dotnet/scripts/test.sh(它会自动补全缺失的脚本副本、设置默认测试连接,并把任意参数直接透传给dotnet test):
cd dotnet ./scripts/test.sh # 全量测试(Redis + PostgreSQL) ./scripts/test.sh --filter Name~Flow # 只跑子集测试连接可通过环境变量指向其他服务器:
BULLMQ_TEST_REDIS(默认localhost:6379)BULLMQ_TEST_POSTGRES(默认Host=localhost;Database=bullmq_test;Username=$USER)
测试脚本还处理了dotnet不在 PATH 的场景(回退到/opt/homebrew/opt/dotnet/libexec),并导出DOTNET_CLI_TELEMETRY_OPTOUT=1、DOTNET_NOLOGO=1保持输出干净。
共享一致性测试套件对应 dotnet/tests/BullMQ.Tests(如BackendConformanceTests.cs同时针对两个适配器运行),这正是 README 所述“同一套测试验证两个后端”的落地体现。
跨语言互操作的实战意义
由于所有运行时共享同一批 Lua/SQL 脚本,.NET 可以:
- 消费 Node.js/Python/Elixir/Rust/PHP 生产者入队的任务;
- 反向为其他语言的 Worker 提供任务;
- 在 Redis 与 PostgreSQL 之间切换而不改业务代码(仅改选项配置)。
这意味着团队可以按服务选择语言,而队列语义、原子性、幂等与重试行为保持一致。从仓库的 README.md(根目录)可见,BullMQ 生态覆盖 NodeJS、Python、.NET、Elixir、Rust、PHP 六种运行时,.NET 绑定正是其中一环。
许可证
BullMQ for .NET 采用 MIT 许可证,详见仓库根目录 LICENSE。 </output_article>
- 后端
- 消息队列
- 任务调度
【免费下载链接】bullmq
BullMQ - Message Queue and Batch processing for NodeJS, Python, .NET, Elixir, Rust and PHP based on Redis or PostgreSQL
相关推荐
BullMQ 实战指南:基于 Redis/PostgreSQL 的多语言分布式任务队列全面解析
BullMQ 实战指南:基于 Redis/PostgreSQL 的多语言分布式任务队列全面解析 本指南以 BullMQ 仓库根目录 README.md http
后端消息队列任务调度Bull 队列完全指南:基于 Redis 的 Node.js 分布式任务与消息队列实战
Bull 队列完全指南:基于 Redis 的 Node.js 分布式任务与消息队列实战 本文围绕开源仓库 bull(一个基于 Redis 的 Node.js 任
任务调度后端BullMQ 是什么:基于 Redis 的分布式任务队列核心特性与设计原理
BullMQ 是什么:基于 Redis 的分布式任务队列核心特性与设计原理 本文以官方文档《What is BullMQ》为核心骨架,结合仓库源码(src 目录
后端消息队列任务调度
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考