BullMQ for Elixir 1.0 首发版全景解析:核心模块、架构设计与 Redis/PostgreSQL 双后端实现
后端消息队列任务调度【免费下载链接】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 仓库中 elixir/CHANGELOG.md 记录的 1.0.0 首发版内容为主体骨架逐项剖析 BullMQ 在 Elixir 生态中的首次落地从 Queue/Worker/Job 三大核心到 Backoff、限流、JobScheduler、FlowProducer、StalledChecker、QueueEvents、Telemetry 等高级能力再到 Config、Keys、Scripts、RedisConnection 等基础设施层。全文结合仓库内 Elixir 源码elixir/lib/bullmq与测试文件进行印证读者可在阅读后掌握该 Elixir 端口各模块的 API 形态、底层实现机制与 Node.js 生态互操作的边界条件。版本背景BullMQ 生态中的 Elixir 原生端口BullMQ 是一个基于 Redis 或 PostgreSQL 的分布式消息队列与批处理库原生实现为 Node.js当前为 v5.x。Elixir 端口的目标是用 Erlang/OTP 的进程模型提供与 Node.js 版完全兼容的队列语义让同一批 BullMQ 队列可以被两种运行时的 Worker 混合消费。根据 elixir/CHANGELOG.md 的记录1.0.0是 BullMQ for Elixir 的初始发布版发布于 2025-12-04一次性补齐了生产级任务队列所需的全部能力面核心队列入队add/3、add_bulk/3、暂停/恢复、按 ID 取任务、drain 与 obliterateWorker可配置并发、自动锁续期、优雅关闭、限流支持Job优先级、延迟执行、带退避的自动重试、进度上报、自定义任务 ID调度与编排cron/间隔调度JobScheduler、父子任务依赖FlowProducer可靠性stalled 任务检测与自动恢复可观测性QueueEvents 实时事件流、Telemetry 追踪集成。当前仓库中 elixir/mix.exs 的版本号已是2.2.3说明自 1.0.0 首发之后该项目仍在持续迭代但 1.0.0 所奠定的模块划分与兼容性原则始终是理解后续版本的基础。CHANGELOG 同时明确了首发版的兼容性基线与 Node.js BullMQ v5.x 兼容、要求 Elixir 1.15、Erlang/OTP 26、Redis 6.0PostgreSQL 后端则要求 PostgreSQL 13见 elixir/README.md。核心队列BullMQ.Queue入队单任务与批量Queue是添加任务的主入口。源码 提供了两种调用形态无状态函数式调用推荐大多数场景传入队列名字符串并通过connection:选项指定后端连接GenServer 形态将 Queue 启动为受监督的进程之后用原子名称调用。# 添加普通任务 {:ok, job} BullMQ.Queue.add(emails, send-welcome, %{ to: userexample.com, template: welcome }, connection: :my_redis) # 延迟任务delay 单位为毫秒 {:ok, job} BullMQ.Queue.add(emails, reminder, %{message: Dont forget!}, connection: :my_redis, delay: 60_000 # 1 分钟 ) # 优先级任务数值越小优先级越高 {:ok, job} BullMQ.Queue.add(emails, urgent, %{}, connection: :my_redis, priority: 1 )批量入队add_bulk/3是吞吐优化的关键路径。从 queue.ex 的源码实现看它会在内部按任务选项将任务分成两类标准任务无delay、无priority走后端优化过的批量命令Redis 后端使用流水线 pipeline延迟/优先级任务则回退到顺序单条添加。源码注释给出的参考量级为单连接约 5 万任务/秒、4 连接池约 7 万 任务/秒——这些是仓库源码中标注的基准表述具体数值取决于硬件与网络环境。jobs [ {email, %{to: user1example.com}, []}, {email, %{to: user2example.com}, []}, {email, %{to: user3example.com}, [priority: 1]} ] {:ok, added_jobs} BullMQ.Queue.add_bulk(emails, jobs, connection: :my_redis)批量场景还支持:connection_pool连接池并行分发、:max_pipeline_size默认 10_000 条/批以及:atomic事务开关若部分任务失败返回{:error, {:partial_failure, results}}便于调用方做部分成功的补偿处理。队列运维操作CHANGELOG 首发版包含的队列管理操作在 queue.ex 中均有对应实现# 暂停 / 恢复队列 :ok BullMQ.Queue.pause(emails, connection: :my_redis) :ok BullMQ.Queue.resume(emails, connection: :my_redis) {:ok, is_paused} BullMQ.Queue.paused?(emails, connection: :my_redis) # 清空等待队列 :ok BullMQ.Queue.drain(emails, connection: :my_redis) # 删除 / 重试单个任务 :ok BullMQ.Queue.remove_job(emails, job-id-123, connection: :my_redis) :ok BullMQ.Queue.retry_job(emails, job-id-123, connection: :my_redis)队列状态查询方面get_counts/2会返回waiting/active/delayed/prioritized/completed/failed/paused/waiting_children的完整计数映射其中paused会并入waitingget_jobs/3支持按单个状态或状态列表分页拉取start/end/asc选项另有get_waiting、get_active、get_delayed、get_completed、get_failed等便捷函数以及get_meta/2读取队列元数据暂停标记、Lua 脚本能力版本号、并发与限流配置。Worker并发、锁续期与优雅关闭BullMQ.Worker是消费端核心。worker.ex 的文档与实现展示了其与 Node.js 版的本质差异Node.js 靠单线程异步并发而 Elixir Worker 将每个并发任务放进独立进程在 Worker 的监督之下运行实现真正的并行。defmodule MyApp.EmailWorker do def process(%BullMQ.Job{name: send-welcome, data: data}) do MyApp.Mailer.send_welcome(data[to], data[template]) {:ok, %{sent: true}} end def process(%BullMQ.Job{name: name}) do {:error, Unknown job type: #{name}} end end {:ok, worker} BullMQ.Worker.start_link( queue: emails, connection: :my_redis, processor: MyApp.EmailWorker.process/1, concurrency: 5, on_completed: fn job, result - IO.puts(Job #{job.id} completed with #{inspect(result)}) end, on_failed: fn job, reason - IO.puts(Job #{job.id} failed: #{reason}) end )处理器函数的返回值契约处理器processor的返回约定在 worker.ex 模块文档中有明确定义这是理解重试/流转语义的关键返回值语义{:ok, result}或:ok任务成功完成result会被保存为返回数据{:error, reason}任务失败计入尝试次数触发重试逻辑{:delay, ms}延迟后重试不递增尝试次数{:rate_limit, ms}因限流将任务移入延迟队列:waiting将任务放回等待队列:waiting_children移入 waiting-children 状态等待子任务完成关键配置项与默认值Worker 的opts_schema经 NimbleOptions 校验核心参数及默认值如下:queue必填— 要消费的队列名:connection必填— 后端连接原子名或 pid:processor必填autorun: false时可传nil配合手动处理— 处理函数支持 1/2/3 元函数三参数形态可拿到BullMQ.CancellationToken实现取消:concurrency默认 1— 最大并发任务数:lock_duration默认 30_000 ms— 任务锁 TTLWorker 会自动续期:stalled_interval默认 30_000 ms— stalled 检查周期:max_stalled_count默认 1— 任务被判定 stalled 达到该次数后转为 failed:limiter— 限流配置%{max: n, duration: ms}:drain_delay默认 5 秒— 队列空时的阻塞等待超时。事件回调方面除了示例中的on_completed/on_failed还支持on_active、on_progress、on_stalled、on_error、on_lock_renewal_failed。其中锁续期失败回调接收一批 job ID受影响任务会自动以{:lock_lost, job_id}理由取消防止重复处理。优雅关闭Worker 的close/1,2默认会等待所有 active 任务完成后再退出force: true可跳过等待立即关闭:ok BullMQ.Worker.close(worker) :ok BullMQ.Worker.close(worker, force: true)Job 数据模型与任务生命周期BullMQ.Job结构体承载任务的全部元数据。job.ex 定义了 7 种状态waiting等待处理、active处理中、delayed延迟到未来时间、prioritized优先级队列、completed成功、failed重试耗尽后失败、waiting-children父任务等待子任务。Job 选项参考入队时可用的选项queue.ex 模块文档:job_id— 自定义任务 ID缺省自动生成支持字符串/整数:delay— 延迟毫秒数到期前任务停留在 delayed 状态:priority— 优先级0为最高默认:attempts— 最大重试次数默认 1即不重试:backoff— 重试退避配置:lifo— 使用 LIFO 而非默认 FIFO 排序:timeout— 任务超时毫秒数:remove_on_complete/:remove_on_fail— 完成/失败后清理支持布尔值、数量整数或%{age: ..., count: ..., limit: ...}映射:repeat— 重复调度选项:parent— 用于 Flow 的父任务引用。一个组合示例BullMQ.Queue.add(tasks, process-data, %{data: ...}, connection: :my_redis, priority: 1, # 数值越小优先级越高 delay: 60_000, # 延迟 60 秒 attempts: 5, # 最多重试 5 次 backoff: %{type: exponential, delay: 1000}, remove_on_complete: true, # 完成后立即清理 remove_on_fail: 100 # 失败任务最多保留最近 100 个 )与 Node.js 存储格式的互操作值得注意的实现细节Elixir 版在把 Job 选项写入存储时会做短键编码job.ex 中的opts_encode_map/opts_decode_map例如deduplication → de、fail_parent_on_failure → fpof、keep_logs → kl、telemetry_metadata → tm并在读取时解码回 snake_case。这正是为了保证与 Node.js BullMQ 的 Redis 数据结构完全一致让跨运行时读写成为可能。进度上报处理函数内可通过BullMQ.Worker.update_progress(job, progress)上报 0~100 的进度值def process(%BullMQ.Job{} job) do Enum.each(1..100, fn i - do_work(i) BullMQ.Worker.update_progress(job, i) end) {:ok, done} end退避策略BullMQ.Backoff首发版内置两种退避策略并支持注册自定义策略backoff.exfixed每次重试固定间隔与尝试次数无关exponential第n次重试的延迟为2^(n-1) × base_delay如 base1000ms 时序列为 1s、2s、4s、8s…自定义通过BullMQ.Backoff.register(:linear, fun)注册函数签名为(attempt, base_delay, error, job) - delay_ms可使用Agent注册表动态增删Jitterjitter: 0.0~1.0在计算出的延迟上叠加 ±jitter 比例的随机扰动避免重试风暴同步化源码中apply_jitter/2实现为在delay×(1-jitter)与delay×(1jitter)区间内取随机值。# 指数退避 抖动 BullMQ.Queue.add(my_queue, job, %{}, connection: :my_redis, attempts: 5, backoff: %{type: :exponential, delay: 1_000, jitter: 0.2} ) # 自定义线性退避 BullMQ.Backoff.register(:linear, fn attempt, delay, _error, _job - attempt * delay end) BullMQ.Queue.add(my_queue, job, %{}, connection: :my_redis, attempts: 5, backoff: %{type: :linear, delay: 1_000} )未知的自定义策略会安全回退为 fixed 行为calculate_from_config/2支持从配置 map%{type:, delay:, jitter:}直接推导延迟供重试调度内部使用。速率限制Worker 通过:limiter选项启用队列级限流worker.ex{:ok, worker} BullMQ.Worker.start_link( queue: api-calls, connection: :my_redis, processor: process/1, limiter: %{max: 100, duration: 60_000} # 每分钟最多 100 个 )CHANGELOG 首发版描述的限流能力面包括队列级限流如上、基于分组的限流以及手动触发限流即处理器返回{:rate_limit, ms}将任务移入延迟队列不视为失败、不消耗 attempts。限流元数据max/duration会写入队列元数据可经Queue.get_meta/2读取。定时任务调度BullMQ.JobSchedulerJobScheduler 取代了旧式 repeatable jobs 概念支持 cron 表达式与固定间隔两种调度模式job_scheduler.ex# cron 模式每小时整点 {:ok, job} BullMQ.JobScheduler.upsert(:my_redis, maintenance, cleanup, %{pattern: 0 * * * *}, cleanup-job, %{type: hourly}, prefix: bull ) # 间隔模式每 60 秒 {:ok, job} BullMQ.JobScheduler.upsert(:my_redis, heartbeats, ping, %{every: 60_000}, heartbeat, %{}, prefix: bull ) # 管理操作 {:ok, schedulers} BullMQ.JobScheduler.list(:my_redis, maintenance, prefix: bull) {:ok, removed} BullMQ.JobScheduler.remove(:my_redis, maintenance, cleanup, prefix: bull)Repeat 选项包括:patterncron、:every毫秒间隔与 pattern 互斥、:limit最大执行次数、:start_date/:end_date起止时间支持毫秒时间戳或DateTime、:tzcron 时区默认 UTC、:immediately立即先执行一次仅 pattern 模式、:offset间隔模式的相位偏移。跨运行时 cron 兼容性边界重要job_scheduler.ex 模块文档用醒目的警告块列出了 Elixir 版与 Node.js 版在 cron 语义上的差异这是混合部署时最容易踩坑的点特性ElixirNode.js是否兼容5 字段 cron无秒✅✅✅6 字段 cron含秒❌✅❌星期日 7✅✅✅星期日 0❌✅❌原因是 Elixir 使用的crontab库默认只解析标准 5 字段格式而 Node.js 支持可选的 6 字段秒。跨平台兼容建议一律使用 5 字段表达式星期日用7而非0或者干脆使用:every间隔调度。cron 字段速查5 字段分钟(0-59)、小时(0-23)、日(1-31)、月(1-12)、星期(1-7周一至周日)。流程编排BullMQ.FlowProducerFlowProducer 用于创建父子任务依赖关系且整个流程的创建是原子的——要么全部入队要么全部失败Redis 后端基于 MULTI/EXEC 事务实现见 flow_producer.ex 模块文档{:ok, flow} BullMQ.FlowProducer.add(%{ name: process_order, queue_name: orders, data: %{order_id: 123}, children: [ %{name: validate, queue_name: validation, data: %{order_id: 123}}, %{name: check_inventory, queue_name: inventory, data: %{order_id: 123}}, %{name: process_payment, queue_name: payments, data: %{order_id: 123}} ] }, connection: :my_redis)嵌套流程flow 节点支持递归的children字段可构建任意深度的任务树批量创建add_bulk/2把多个 flow 放在单个事务中执行全成功或全失败失败策略默认子任务失败会导致父任务失败可通过:fail_parent_on_failure、:ignore_dependency_on_failure等选项调整这些选项会以短键形式编码存储见上文 Job 互操作一节。父任务处理器内可用BullMQ.Job的 Flow 方法获取子任务结果get_children_values/1成功的子任务返回值、get_ignored_children_failures/1被忽略的失败、get_dependencies/1未完成依赖。可靠性Stalled 任务检测与恢复Worker 崩溃、机器断电、网络抖动都可能导致任务被取出后失联。stalled_checker.ex 采用两阶段检测算法防止误判标记阶段Mark把没有有效锁的任务移入 stalled 集合恢复阶段Recover下次检查时仍在 stalled 集合中的任务依据max_stalled_count被重新入队或标记失败。{BullMQ.Worker, queue: emails, connection: :my_redis, processor: MyApp.send_email/1, lock_duration: 30_000, # 锁 TTL 默认 30s stalled_interval: 30_000, # 检查周期默认 30s max_stalled_count: 1 # 默认 1 }模块文档对参数调优给出明确建议max_stalled_count默认 1 是因为 stalled 属于罕见事件若同一任务反复 stalled通常意味着更严重的问题处理器崩溃、资源耗尽、外部服务故障、处理逻辑 bug应排查根因而非盲目调大阈值lock_duration也仅在任务合法地需要超过 30 秒才续期时才需要调大。此外BullMQ.StalledChecker.check(conn, queue)可手动触发检查返回%{recovered: n, failed: n}。实时事件BullMQ.QueueEventsQueueEvents 基于 Redis Streams 提供队列级实时事件订阅elixir/README.md{:ok, events} BullMQ.QueueEvents.start_link( queue: tasks, connection: :my_redis ) BullMQ.QueueEvents.subscribe(events) receive do {:bullmq_event, :completed, %{jobId id}} - IO.puts(Job #{id} completed!) {:bullmq_event, :failed, %{jobId id, failedReason reason}} - IO.puts(Job #{id} failed: #{reason}) end事件流长度可通过 Queue 的streams选项中的max_len默认 10_000控制。除订阅外部事件外Worker 自身也提供on_*回调作为进程内的事件钩子与 QueueEvents 形成进程内 跨进程两种互补的观察方式。可观测性Telemetry 集成CHANGELOG 首发版明确列出了 Telemetry 的三种事件类别任务生命周期事件completed/failed/active/stalled 等Worker 事件限流事件。实现上BullMQ.Telemetry通过BullMQ.Telemetry.Behaviour行为定义telemetry/behaviour.ex并提供基于 OpenTelemetry 的默认实现BullMQ.Telemetry.OpenTelemetrytelemetry/opentelemetry.ex支持span 级别的分布式追踪。Queue 与 Worker 均接受telemetry:选项如BullMQ.Telemetry.OpenTelemetry来开启埋点opentelemetry_api是可选依赖见 mix.exs不强制引入。基础设施层Config / Keys / Scripts / RedisConnectionConfigNimbleOptions 驱动的配置校验CHANGELOG 强调所有配置均基于NimbleOptions schema做编译期与运行期校验。这在实际代码中随处可见Queue 与 Worker 都通过NimbleOptions.new!定义opts_schema非法参数会在启动时立即报错而非静默失效。例如 Worker 的:processor只接受 1/2/3 元函数或nillimiter必须是含:max/:duration的 map。Keys一致的键命名BullMQ.Keys负责队列所有 Redis 键的构造保证命名与 Node.js 版一致keys.ex。默认前缀为bull可通过prefix:选项全局或按调用覆盖——这一点对共享 Redis 实例、多租户隔离至关重要。Scripts原子操作、SHA 缓存与 EVAL 回退Elixir 版直接复用 Node.js 版同一套 Lua 脚本发布时由mix scripts.copy别名从rawScripts复制到priv/scripts见 mix.exs。执行机制上SHA 缓存脚本先SCRIPT LOAD后以 SHA 调用减少网络传输EVAL 回退缓存失效时回退到EVAL重新加载。整套 Lua 脚本还同时被 Python/.NET/Rust/PHP 端口共享这正是行为一致、无 SQL/Lua 分叉架构的基石。RedisConnectionNimblePool 连接池连接层基于 NimblePool 实现redis_connection.ex 与 redis_connection/pool.ex支持可配置池大小并区分普通命令连接与专用阻塞连接Worker 的空队列阻塞等待使用后者。双后端Redis 与 PostgreSQL 的可插拔架构虽然 CHANGELOG 首发版的核心围绕 Redis 后端展开但当前仓库中 Elixir 端口已具备可插拔后端能力elixir/README.md所有高层模块Queue、Worker、Job、QueueEvents、FlowProducer、JobScheduler均面向BullMQ.Backend行为编程与具体存储解耦。选择方式有两种# 全局配置 config :bullmq, :backend, BullMQ.Backends.Redis # 或按调用/进程覆盖 BullMQ.Queue.add(emails, send, %{to: userexample.com}, connection: :my_pg, backend: BullMQ.Backends.Postgres )PostgreSQL 后端复用与 Node.js 完全相同的 SQL 产物src/postgres/migrations/*.sql与src/postgres/commands/*.sql。在 Hex 包与 Mix release 中从elixir/priv/postgres读取仓库开发模式下回退到src/postgres迁移由BullMQ.Backends.Postgres.Connection.start_link/1自动执行{:ok, _conn} BullMQ.Backends.Postgres.Connection.start_link( name: :my_pg, url: postgresql://localhost:5432/bullmq, schema: bullmq ) {:ok, worker} BullMQ.Worker.start_link( queue: emails, connection: :my_pg, backend: BullMQ.Backends.Postgres, processor: MyApp.EmailWorker.process/1 )依赖与工程实践elixir/mix.exs 展示了实现这些能力所用的依赖栈可作为技术选型参考RedixRedis 客户端、Postgrex可选 PostgreSQL 客户端、NimblePool连接池、NimbleOptions配置校验、JasonJSON 编解码、crontabcron 解析、msgpaxLua 脚本的 MessagePack 编码、elixir_uuid任务 ID 生成、telemetry 与可选的 opentelemetry_api追踪。CHANGELOG 首发版还特别强调了两项工程化交付物完整文档体系Getting Started、Job Options 参考、Worker 配置、Rate Limiting 指南、Flow 模式、Telemetry 配置对应仓库 elixir/guides 下的同名指南与测试套件各模块单元测试 依赖 Redis 的集成测试见 elixir/test/bullmq 下的worker_integration_test.exs、queue_integration_test.exs、flow_producer_test.exs、job_scheduler_integration_test.exs等。结语BullMQ for Elixir 1.0 首发版不是 Node.js 版的简单翻译而是在保持 Redis 数据结构、Lua 脚本与队列语义完全兼容的前提下用 Erlang/OTP 的进程模型重构了 Worker 的并发模型每任务独立进程用 NimbleOptions 强化了配置校验用 NimblePool 重写了连接管理。对于已有 Node.js BullMQ 基础设施的团队Elixir 端口提供了一条渐进式迁移路径混合 Worker 并存、共享同一队列再逐步将处理器重写为 Elixir 实现最终享受 BEAM 虚拟机带来的并发与容错特性。赞分享后端消息队列任务调度【免费下载链接】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点击查看免费下载相关推荐Paper2GUI 后端架构C核心模块的设计与实现Paper2GUI 后端架构C核心模块的设计与实现 Paper2GUI作为一款致力于让普通人简单便捷使用前沿人工智能技术的桌面应用工具箱其后端架构的设计人工智能AI 应用桌面应用RedisDesktopManager源码解析架构设计与核心模块实现RedisDesktopManager源码解析架构设计与核心模块实现 RedisDesktopManager是一个功能强大的Redis数据库管理工具提供直观数据库客户端桌面应用Steam挂刀行情站24小时自动追踪四大平台饰品价格的专业指南 Steam挂刀行情站24小时自动追踪四大平台饰品价格的专业指南 SteamTradingSiteTrackerSteam挂刀行情站是一个强大的开源工后端网页爬虫上一篇JavaScript Standard Style 代码规范详解下一篇JavaScript Standard Style 代码规范详解创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考