后端消息队列任务调度【免费下载链接】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 官方基准测试文档elixir/guides/benchmarks.md为核心系统讲解其在 Redis 之上的单/多 Worker 吞吐量表现、LockManager 单定时器架构、批量入队add_bulk的 MULTI/EXEC 与连接池加速原理并给出完整可复现的基准脚本与生产级优化建议。读完本文你将掌握如何在自身环境中复现基准、读懂吞吐量拐点背后的瓶颈成因并依据实测数据为队列配置合理的 Worker 数量与并发度。测试环境与适用前提文档中记录的基准测试环境如下项目版本CPUApple M2 ProRedis7.4.0Elixir1.15Erlang/OTP26文档另一处对运行环境的补充描述为MacBook ProApple Silicon、Redis 7.xDocker 运行、Elixir 1.18.x、Erlang/OTP 27.x。两处版本表述略有差异说明基准数值会随硬件、Redis 部署方式与语言版本变化复现时应以自身环境为准。从仓库源码看redis_connection.ex 中定义minimum_redis_version {6, 2, 0}连接建立时会通过INFO server校验 Redis 版本低于 6.2 会直接抛出ArgumentError因此复现基准前请确保 Redis ≥ 6.2。此外Worker 在无任务时会通过 Redis 的BZPOPMIN进行阻塞等待见 worker.ex 中get_next_job的文档说明这也是依赖 Redis 6.2 浮动超时能力的原因之一。单 Worker 吞吐量并发度与饱和拐点文档使用瞬时no-op任务对单个 Worker 在不同并发度下进行测试结果如下ConcurrencyJobsTimeThroughput100500206ms2,427 j/s2001,000357ms2,801 j/s5002,500510ms4,901 j/s关键结论单个 Worker 在约 500 并发时趋于饱和其瓶颈并非 CPU 或调度器而是从 Redis 顺序拉取任务sequential job fetching的串行路径。Elixir 虽然能用轻量进程并行处理任务但任务供给fetch仍由单个 Worker 进程完成因此并发度从 100 提升到 500 时吞吐只从约 2,400 j/s 增长到约 4,900 j/s增幅明显放缓。从源码结构可以印证这一点worker.ex 中handle_info(:fetch_jobs, ...)会计算available_slots concurrency - (active_jobs in_flight_workers)再由 Worker 进程统一发起任务抓取任务处理则由独立进程执行供给与消费天然分离单 Worker 的抓取频率成为上限约束。多 Worker 吞吐量近乎线性的横向扩展将多个 Worker每个 500 并发指向同一队列结果如下WorkersConc/WorkerTotal ConcJobsTimeThroughput15005002,500608ms4,111 j/s55002,50012,5001,011ms12,363 j/s105005,00025,0001,515ms16,501 j/s关键结论多个 Worker 几乎线性扩展10 个 Worker 时达到 16,500 j/s。每个 Worker 拥有独立的 Redis 连接、独立的 LockManager1 个定时器而非 N 个、并独立处理任务从而绕开了单 Worker 的串行抓取瓶颈。扩展效率相对单 Worker 基准如下WorkersThroughputvs Single WorkerEfficiency14,111 j/s1.0x100%512,363 j/s3.0x60%1016,501 j/s4.0x40%高 Worker 数量下效率下降的原因在文档中有明确结论所有 Worker 竞争同一个队列Redis 成为瓶颈。也就是说横向扩展受限于 Redis 单实例的处理能力而非 Elixir 侧。架构原理LockManager 为何能支撑高并发BullMQ Elixir 使用LockManager模块与 Node.js 版设计类似其核心优化是每个 Worker 只维护一个定时器用于续期所有活跃任务的锁而不是为每个任务创建独立定时器Without LockManager: 500 concurrent jobs 500 timers With LockManager: 500 concurrent jobs 1 timer (per worker)查看 lock_manager.ex 源码可以确认具体机制单定时器驱动init/1中调用schedule_renewal/1定时器按lock_renew_time / 2间隔触发lock_duration默认 30,000ms则lock_renew_time默认 15,000ms即每 7.5 秒检查一次批量续期handle_info(:extend_locks, ...)遍历tracked_jobs将“快过期”ts threshold now的任务一次性交给Backend.extend_locks/4批量续期而不是逐任务下发命令失败回调续期失败的任务 ID 会触发on_lock_renewal_failed回调在 worker.ex 中该回调会把对应任务取消reason 为{:lock_lost, job_id}防止锁丢失后出现重复处理。在 Worker 启动流程handle_info(:start, ...)中每个 Worker 都会LockManager.start_link(...)并显式Process.link到自身LockManager 崩溃会连带 Worker 重启由监督树兜底。这正是“多个 Worker 各自持有 1 个定时器、互不干扰”的实现基础。为什么多 Worker 优于单 Worker 超高并发每个 Worker 具备三个独立维度共同打破单点串行独立的 Redis 连接并行抓取任务避免单个连接上的请求排队独立的 LockManager1 个定时器而非 N 个锁续期开销与并发任务数解耦独立的处理进程任务在各自进程内并行执行互不阻塞。因此文档给出的实践结论是想提升吞吐优先加 Worker而不是把单个 Worker 的并发调得过高。复现基准从零跑通三个基准脚本前置准备# 启动 Redis文档使用 redis:7 镜像 docker run -d --name redis -p 6379:6379 redis:7 # 安装依赖 cd elixir mix deps.get单 Worker 基准mix run -e alias BullMQ.{Queue, Worker, RedisConnection} configs [ {100, 500}, # {concurrency, job_count} {200, 1000}, {500, 2500}, ] IO.puts(| Concurrency | Jobs | Time | Throughput |) IO.puts(|-------------|--------|---------|------------|) for {concurrency, job_count} - configs do conn_name :bench_#{:erlang.unique_integer([:positive])} {:ok, _} RedisConnection.start_link(host: localhost, port: 6379, name: conn_name) queue_name bench_#{:erlang.unique_integer([:positive])} completed :counters.new(1, []) processor fn _job - :counters.add(completed, 1, 1) :ok end jobs for i - 1..job_count, do: {job-#{i}, %{}, []} {:ok, _} Queue.add_bulk(queue_name, jobs, connection: conn_name) start_time System.monotonic_time(:millisecond) {:ok, worker} Worker.start_link( queue: queue_name, connection: conn_name, concurrency: concurrency, processor: processor ) # Wait for completion wait fn wait_fn - Process.sleep(50) if :counters.get(completed, 1) job_count, do: wait_fn.(wait_fn) end wait.(wait) elapsed System.monotonic_time(:millisecond) - start_time throughput trunc(job_count / elapsed * 1000) IO.puts(| #{concurrency} | #{job_count} | #{elapsed}ms | #{throughput} j/s |) GenServer.stop(worker) {:ok, keys} RedisConnection.command(conn_name, [KEYS, bull:#{queue_name}:*]) if length(keys) 0, do: RedisConnection.command(conn_name, [DEL | keys]) end 脚本要点用:counters原子计数记录完成数并轮询等待吞吐量按job_count / elapsed * 1000计算每个队列跑完后清理bull:queue:*键避免污染下一次测试。多 Worker 基准mix run -e alias BullMQ.{Queue, Worker, RedisConnection} configs [ {1, 500, 2500}, # {workers, concurrency, jobs} {5, 500, 12500}, {10, 500, 25000}, ] IO.puts(| Workers | Conc/W | Total Conc | Jobs | Time | Throughput |) IO.puts(|---------|--------|------------|--------|----------|-------------|) for {num_workers, concurrency, job_count} - configs do conn_name :bench_#{:erlang.unique_integer([:positive])} {:ok, _} RedisConnection.start_link(host: localhost, port: 6379, name: conn_name) queue_name bench_#{:erlang.unique_integer([:positive])} completed :counters.new(1, []) processor fn _job - :counters.add(completed, 1, 1) :ok end jobs for i - 1..job_count, do: {job-#{i}, %{}, []} {:ok, _} Queue.add_bulk(queue_name, jobs, connection: conn_name) start_time System.monotonic_time(:millisecond) workers for _ - 1..num_workers do {:ok, w} Worker.start_link( queue: queue_name, connection: conn_name, concurrency: concurrency, processor: processor ) w end # Wait for completion wait fn wait_fn - Process.sleep(100) if :counters.get(completed, 1) job_count, do: wait_fn.(wait_fn) end wait.(wait) elapsed System.monotonic_time(:millisecond) - start_time throughput trunc(job_count / elapsed * 1000) IO.puts(| #{num_workers} | #{concurrency} | #{num_workers * concurrency} | #{job_count} | #{elapsed}ms | #{throughput} j/s |) Enum.each(workers, GenServer.stop/1) {:ok, keys} RedisConnection.command(conn_name, [KEYS, bull:#{queue_name}:*]) if length(keys) 0, do: RedisConnection.command(conn_name, [DEL | keys]) end 注意该脚本中多个 Worker 共享同一个conn_name连接连接池会提供多条底层连接见下文 redis_connection.ex 中基于NimblePool的连接池实现仓库内的完整基准套件 suite.exs 则为每个 Worker 单独创建连接实测效果更接近“并行抓取”的理想形态。真实负载基准10ms 任务mix run -e alias BullMQ.{Queue, Worker, RedisConnection} job_duration_ms 10 job_count 5000 concurrency 500 conn_name :bench_conn {:ok, _} RedisConnection.start_link(host: localhost, port: 6379, name: conn_name) queue_name realistic_bench completed :counters.new(1, []) processor fn _job - Process.sleep(job_duration_ms) :counters.add(completed, 1, 1) :ok end jobs for i - 1..job_count, do: {job-#{i}, %{}, []} {:ok, _} Queue.add_bulk(queue_name, jobs, connection: conn_name) start_time System.monotonic_time(:millisecond) {:ok, worker} Worker.start_link( queue: queue_name, connection: conn_name, concurrency: concurrency, processor: processor ) # Wait wait fn wait_fn - Process.sleep(100) if :counters.get(completed, 1) job_count, do: wait_fn.(wait_fn) end wait.(wait) elapsed System.monotonic_time(:millisecond) - start_time throughput trunc(job_count / elapsed * 1000) # Theoretical max: job_count / (job_duration_ms / concurrency) theoretical_max trunc(concurrency / job_duration_ms * 1000) IO.puts(Jobs: #{job_count}, Concurrency: #{concurrency}, Job duration: #{job_duration_ms}ms) IO.puts(Time: #{elapsed}ms, Throughput: #{throughput} j/s) IO.puts(Theoretical max: #{theoretical_max} j/s, Efficiency: #{trunc(throughput / theoretical_max * 100)}%) GenServer.stop(worker) {:ok, keys} RedisConnection.command(conn_name, [KEYS, bull:#{queue_name}:*]) if length(keys) 0, do: RedisConnection.command(conn_name, [DEL | keys]) 该脚本同时输出理论吞吐上限concurrency / job_duration_ms * 1000与实际效率用于判断任务是否真的“吃满”了并发度。仓库自带的进阶基准脚本除文档内嵌脚本外仓库还提供了可直接运行的基准程序均位于 elixir/benchmark 目录throughput_benchmark.exs可配置的吞吐量基准。支持命令行参数--jobs 5000 --concurrencies 10,50,100,200 --job-duration 500 --workers 1内置预热warmup、进度输出并会同时输出 CSV 与 Markdown 表格方便直接引用到文章或对比报告其默认并发序列为[1, 5, 10, 25, 50, 100, 150, 200, 250, 300, 400, 500]add_job_benchmark.exs批量入队饱和度测试详见下一节测试100,000个任务、连接池规模 1/2/4/8/16/32/64并自动定位“最佳”与“饱和”连接数redis_baseline.exsRedis 基线测试分别测量 PING、单条 SET、Pipeline SET、简单 LuaEVAL的吞吐8 条连接并行、各 10,000 次脚本注释给出的判断方法是若 BullMQ 吞吐显著低于 LuaEVAL速率说明瓶颈在moveToActive脚本复杂度若两者相近说明 Redis 本身是瓶颈suite.exs一键完整套件覆盖单 Worker、多 Worker、真实负载与汇总支持环境变量REDIS_HOST、REDIS_PORT、QUICKtrue快速模式切换配置。优化建议从基准数据反推生产配置1. 优先使用多个 Worker实测对比no-op 任务# Good: 10 workers × 500 concurrency 16,500 j/s for _ - 1..10 do Worker.start_link(queue: myqueue, connection: conn, concurrency: 500, processor: process/1) end # Less optimal: 1 worker × 5000 concurrency ~5,000 j/s Worker.start_link(queue: myqueue, connection: conn, concurrency: 5000, processor: process/1)2. 单 Worker 并发甜点区200~500超过 500 后因串行抓取而收益递减。这一区间与 worker.ex 中concurrency选项默认 1、:pos_integer的语义一致并发代表 Worker 可同时运行的处理进程数。3. 任务耗时决定瓶颈位置瞬时任务no-opRedis 成为瓶颈单 Worker 约 5,000 j/s 封顶有真实计算/IO 的任务通常在触达 Redis 上限之前CPU 或 IO 就已先成为瓶颈此时多 Worker 高并发的价值更明显。4. 生产环境使用监督树管理多个 Workerchildren [ {BullMQ.RedisConnection, name: :redis, host: localhost}, # Multiple workers for the same queue Supervisor.child_spec( {BullMQ.Worker, queue: jobs, connection: :redis, concurrency: 500, processor: MyApp.process/1}, id: :worker_1 ), Supervisor.child_spec( {BullMQ.Worker, queue: jobs, connection: :redis, concurrency: 500, processor: MyApp.process/1}, id: :worker_2 ), # ... more workers ] Supervisor.start_link(children, strategy: :one_for_one)从 worker.ex 的init/1可知Worker 支持name、lock_duration默认 30,000ms、stalled_interval默认 30,000ms、max_stalled_count默认 1、limiter限流%{max:, duration:}等选项配合Supervisor.child_spec/2可做到崩溃自动重启。批量添加任务add_bulk性能MULTI/EXEC 与连接池入队侧同样有独立基准add_bulk/3借助 Redis MULTI/EXEC 事务实现原子批量写入全有或全无并通过并行处理达到极高的入队速率。10 万任务实测MethodConnectionsThroughputSpeedupSequential15,700 j/s1.0xAtomic124,000 j/s4.2xAtomic239,000 j/s6.8xAtomic454,000 j/s9.5xAtomic858,000 j/s10.2xAtomic1656,000 j/s9.8x关键发现事务带来约 4 倍加速把 Redis 命令批量包裹进原子 MULTI/EXEC 事务原子性保证每个批次的任务要么全部写入、要么全部不写4~8 条连接达到饱和超过 8 条连接后吞吐不再上升16 条连接反而略降默认配置即最优atomic: true与max_pipeline_size: 10_000即可命中峰值性能。工作原理对比Sequential (5,700 j/s): Job1 → Redis → Response → Job2 → Redis → Response → ... Atomic (24,000 j/s): MULTI → [Job1, Job2, ..., JobN] → EXEC → [Response1, ..., ResponseN] (all jobs in batch added atomically) Parallel Atomic (58,000 j/s): Conn1: MULTI [Jobs batch 1] EXEC → Redis → Conn2: MULTI [Jobs batch 2] EXEC → Redis → All in parallel Conn3: MULTI [Jobs batch 3] EXEC → Redis → ... (each connections batch is atomic)从 redis.ex 后端实现 可以看到这一流程的落地细节add_jobs/3先确保加载add_standard_job脚本再调用Scripts.build_bulk_add_commands/2生成批量命令无连接池时按max_pipeline_size分批atomic时用execute_transaction即 MULTI/EXEC否则用execute_pipeline有连接池时按div(length(commands), pool_size)切分再通过Task.async_stream(..., max_concurrency: pool_size)并行分发到各连接每连接内部仍按max_pipeline_size二次分批。使用连接池写入大批次# Create a connection pool pool for i - 1..8 do name :redis_pool_#{i} {:ok, _} BullMQ.RedisConnection.start_link(name: name, host: localhost) name end # Add jobs with parallel processing # Each batch is added atomically (default: atomic: true) jobs for i - 1..100_000, do: {job, %{index: i}, []} {:ok, added} BullMQ.Queue.add_bulk(my-queue, jobs, connection: :redis, connection_pool: pool )注意使用connection_pool时每条连接上的批次各自原子跨连接的整体批量操作并不具备全局原子性若业务要求整体原子应保持单连接 atomic: true。选项参考OptionDefaultDescriptionpipelinetrueUse pipelining for efficiencyatomictrueWrap batches in MULTI/EXEC transactions. Withconnection_pool, each connections batch is atomic independently.connection_poolnilList of connections for parallel processingmax_pipeline_size10_000Maximum jobs per pipeline batch在 queue.ex 的add_bulk/3文档中同样注明了这些选项并额外说明标准任务无 delay、无 priority走优化的批量命令路径带 delay 或 priority 的任务会自动回退到顺序添加若部分任务失败返回{:error, {:partial_failure, results}}。连接池内的连接在生产环境中应放入监督树管理以保证生命周期与断线重连。运行批量入队基准cd elixir mix run benchmark/add_job_benchmark.exs该脚本默认测试 100,000 个任务、chunk_size 100、连接池规模[1, 2, 4, 8, 16, 32, 64]先以pipeline: false跑出顺序基线再逐级放大连接池最终打印包含“BEST / SATURATES”标记的结果表与汇总行。结果解读与注意事项这些基准以 no-op 任务测量原始吞吐上限真实业务吞吐取决于任务处理时间每个任务实际执行的工作量Redis 延迟到 Redis 服务器的网络距离任务数据体积更大的 payload 会拉长序列化/反序列化耗时外部依赖第三方 API 调用、数据库查询等。对 IO 密集型负载Elixir 的轻量进程优势最为明显——成千上万个并发任务可同时等待外部资源而不互相阻塞。另外需要留意本文引用的数据来自文档记录的特定环境Apple M2 Pro / Redis 7.4.0换用不同 CPU、Redis 实例或网络拓扑后绝对数值会变化但“单 Worker 约 500 并发饱和、多 Worker 近线性扩展、8 连接左右入队饱和”这些相对规律通常仍然成立。若需定位自己环境下的瓶颈可先跑 redis_baseline.exs 测量 Redis 原始能力再对比 BullMQ 实测吞吐即可判断瓶颈在 Redis 还是脚本复杂度。小结BullMQ Elixir 在 Erlang/OTP 进程模型之上实现了高吞吐任务处理单 Worker 因串行抓取在约 500 并发饱和约 5,000 j/s而通过多个 Worker 各自持有独立连接与单定时器 LockManager可近乎线性扩展到 16,500 j/s10 Worker 实测入队侧凭借add_bulk的 MULTI/EXEC 事务与连接池并行可从顺序写入的约 5,700 j/s 提升至约 58,000 j/s。生产调优的核心口诀是并发甜点区取 200~500、吞吐不够加 Worker、批量写入保持默认atomic: true并按需扩展连接池。赞分享后端消息队列任务调度【免费下载链接】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点击查看免费下载相关推荐LMCache Transfer Channel 吞吐基准测试工具全解析架构、用法与性能调优LMCache Transfer Channel 吞吐基准测试工具全解析架构、用法与性能调优 导读 本文深入剖析 LMCache 项目自带的 Transfer人工智能大模型缓存抽象模型推理服务LMCache Transfer Channel 吞吐量基准测试工具详解原理、用法与 NUMA 性能调优LMCache Transfer Channel 吞吐量基准测试工具详解原理、用法与 NUMA 性能调优 导读 本文围绕 LMCache 官方提供的 Tra人工智能大模型缓存抽象模型推理服务实测Hyperswitch支付吞吐量基准测试全指南从环境搭建到性能调优实测Hyperswitch支付吞吐量基准测试全指南从环境搭建到性能调优 Hyperswitch作为一款高性能的支付编排平台其吞吐量表现直接影响支付系统的稳后端金融科技上一篇Termshark用户界面定制从字体大小到颜色方案的全攻略下一篇如何高效使用AWS Amplify构建数据模型GraphQL与REST API的实战指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考