news 2026/9/25 5:11:51

BullMQ 官方 .NET 移植版完整指南:基于 Redis 与 PostgreSQL 的跨语言分布式队列

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
BullMQ 官方 .NET 移植版完整指南:基于 Redis 与 PostgreSQL 的跨语言分布式队列
  • 后端
  • 消息队列
  • 任务调度

【免费下载链接】bullmq

BullMQ - 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
点击查看免费下载

<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 定义了三个选项:

选项说明默认值
ConnectionStringNpgsql 连接字符串,如"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

参数默认值说明
Concurrency1并发处理的最大任务数(构造时校验必须>= 1)
LockDuration30000任务处理中的锁时长(毫秒)
LockRenewTimeLockDuration / 2锁续期频率(毫秒)
DrainDelay5Worker 循环前阻塞等待任务的秒数
Namenull可读的 Worker 名称(用于可观测性)
Autoruntrue创建后是否立即开始处理
MaxStalledCount1任务在判定失败前可从 stalled 恢复的最大次数
StalledInterval30000执行 stalled 检查的间隔(毫秒)
SkipStalledCheckfalse禁用 stalled 检查器
SkipLockRenewalfalse禁用周期性锁续期

在 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

项目地址:https://gitcode.com/gh_mirrors/bu/bullmq
点击查看免费下载

相关推荐

上一篇:NemoClaw 发布上下文检查:定位并审计最新全量 E2E main 运行的完整实践指南
下一篇:PyTextRank 教程

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/25 5:11:36

ADAU1701音频DSP实战:从SigmaStudio到硬件设计

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/25 5:09:53

Atlas 300V 24G部署YOLO推理加速卡实战指南

先说结论&#xff1a;Atlas 300V 24G就是一张AI运算加速卡&#xff0c;而且是专为“推理”场景设计的加速卡。如果你手头有训练好的YOLO模型&#xff0c;想在边缘服务器或者机房里面把它跑起来&#xff0c;用这张卡做在线推理&#xff0c;那路子基本是对的。但这里有个很多人刚…

作者头像 李华
网站建设 2026/9/25 5:08:52

nanoGPT GPT 训练、复现与微调实操指南:10 分钟跑通最小闭环

nanoGPT GPT 训练、复现与微调实操指南&#xff1a;10 分钟跑通最小闭环 【免费下载链接】nanoGPT The simplest, fastest repository for training/finetuning medium-sized GPTs. 项目地址: https://gitcode.com/GitHub_Trending/na/nanoGPT nanoGPT 是目前最精简、最…

作者头像 李华