
后端消息队列任务调度【免费下载链接】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点击查看免费下载导读本指南以 Elixir 版 BullMQ 的BullMQ.Queue.add/4与add_bulk/3为切入点系统讲解任务Job的优先级、延迟调度、重试退避、自定义 ID、去重Deduplication、LIFO 处理、自动清理等全部配置项并结合仓库源码elixir/lib/bullmq/queue.ex、job.ex、backoff.ex、keys.ex、scripts.ex与基准测试数据帮助你为每个任务精确配置行为构建高吞吐、可观测、不重复的队列流水线。读完本文你将能独立为 Elixir 项目中的任何任务设计该何时执行、失败如何重试、如何防重、何时清理的完整策略。添加任务一切配置的入口任务是 BullMQ 队列处理的基本单元所有任务级配置都通过BullMQ.Queue.add/4的第四个参数options 关键字列表传入{:ok, job} BullMQ.Queue.add(queue_name, job_name, data, opts)四个参数的含义参数类型说明queue_namestring队列名称决定任务进入哪个队列job_namestring任务类型/名称Worker 处理器据此做模式匹配分发datamap任务负载payload任意可 JSON 序列化的数据optskeyword list任务选项其中:connection为必填项其余为任务行为配置在 queue.ex 中可以看到当以字符串形式传入队列名时add/4会取出:connection与:prefix默认bull构造队列上下文后通过Job.new/4创建任务结构体并写入后端存储def add(queue, name, data, opts) when is_binary(queue) do conn Keyword.fetch!(opts, :connection) prefix Keyword.get(opts, :prefix, bull) ctx Keys.new(queue, prefix: prefix) job Job.new(queue, name, data, opts) add_job(conn, ctx, job) end这里有几个值得注意的底层事实opts中的选项会被转换为 map 存入任务结构体。在 job.ex 中Job.new/4会读取:job_id、:timestamp默认取当前系统毫秒时间戳、:delay默认0、:priority默认0、:parent用于 Flow 父子任务等字段。选项在落库时会编码为短键以兼容 Node.js BullMQ。例如deduplication→de、keep_logs→kl、telemetry_metadata→tm、omit_context→omc见 job.ex 的opts_encode_map。prefix决定所有 Redis 键的前缀默认bull所有访问同一队列的组件必须使用相同前缀否则会看不见彼此。如果你希望队列级统一下发默认选项还可以在启动 Queue GenServer 时配置default_job_opts见 queue.ex所有通过该队列添加的任务都会自动合并这些默认选项。优先级让重要任务先执行BullMQ 使用数值越小优先级越高的规则。默认优先级为0最高因此只有当你显式传入大于0的:priority时任务才会进入优先队列# High priority job (processed first) BullMQ.Queue.add(tasks, urgent-task, %{}, connection: :my_redis, priority: 1 ) # Normal priority (default) BullMQ.Queue.add(tasks, normal-task, %{}, connection: :my_redis ) # Low priority job (processed last) BullMQ.Queue.add(tasks, batch-task, %{}, connection: :my_redis, priority: 100 )底层原理优先任务存放在 Redis 的有序集合sorted set中而不是普通的等待列表。在 keys.ex 中可以确认每个队列都有一个独立的#{base}:prioritized键Worker 通过优先级分数priority score按升序取出任务因此只要队列中还有优先任务就会始终按优先级顺序被处理。这也意味着如果你从未设置过:priority所有任务都走普通 FIFO 等待列表不会产生额外开销。延迟调度让任务在指定时间后执行:delay以毫秒为单位控制任务在多久之后变为可处理状态。两类典型用法# Run in 5 minutes BullMQ.Queue.add(reminders, send-reminder, %{message: Dont forget!}, connection: :my_redis, delay: 5 * 60 * 1000 # 5 minutes in milliseconds ) # Run at a specific time future_time DateTime.utc_now() | DateTime.add(3600, :second) delay DateTime.diff(future_time, DateTime.utc_now(), :millisecond) BullMQ.Queue.add(reports, scheduled-report, %{}, connection: :my_redis, delay: delay )从源码看延迟任务会进入 Redis 的#{base}:delayed有序集合分数为timestamp delay见 scripts.ex 中add_delayed_job的计算逻辑delayed_timestamp timestamp delay。Worker 端会有定时扫描将该集合中到期的任务提升promote回等待列表。因此delay与任务的timestamp创建时间相加决定了任务的绝对到期时刻延迟任务与优先任务一样不会占用普通等待列表的 FIFO 位置你可以通过BullMQ.Queue.get_delayed/2、get_delayed_count/2观察延迟队列见 queue.ex。重试与退避失败后的自动恢复策略基础用法默认情况下任务只尝试 1 次attempts默认1即首次尝试本身。设置:attempts可以指定包含首次尝试在内的总尝试次数配合:backoff控制每次重试之间的等待时间# 3 retries with exponential backoff BullMQ.Queue.add(api-calls, call-api, %{url: ...}, connection: :my_redis, attempts: 3, backoff: %{type: exponential, delay: 1000} ) # Delays: 1s, 2s, 4s # Fixed backoff BullMQ.Queue.add(api-calls, call-api, %{url: ...}, connection: :my_redis, attempts: 5, backoff: %{type: fixed, delay: 5000} ) # Delays: 5s, 5s, 5s, 5s内置退避类型类型行为计算公式exponential每次重试延迟翻倍delay * 2^(attempt - 1)fixed每次重试延迟相同delay以上公式可在 backoff.ex 中直接验证def calculate(:exponential, attempt, base_delay, opts) do jitter Keyword.get(opts, :jitter, 0) delay trunc(:math.pow(2, attempt - 1) * base_delay) apply_jitter(delay, jitter) end抖动Jitter避免惊群效应当大量任务同时失败重试时固定/指数延迟可能造成对下游服务的同步冲击。BullMQ Elixir 支持在退避配置中加入jitter取值0 jitter 1为延迟值引入 ±jitter比例内的随机扰动# Exponential with jitter BullMQ.Queue.add(my_queue, job, %{}, connection: :redis, attempts: 5, backoff: %{type: :exponential, delay: 1_000, jitter: 0.2} )apply_jitter/2的实现backoff.ex会在delay * (1 - jitter)到delay * (1 jitter)之间随机取值。自定义退避策略如果内置的两种策略不够用可以注册自定义策略函数函数签名为(attempt, base_delay, error, job) - delay_ms# Register a custom strategy BullMQ.Backoff.register(:linear, fn attempt, delay, _error, _job - attempt * delay end) # Use it BullMQ.Queue.add(my_queue, job, %{}, connection: :redis, attempts: 5, backoff: %{type: :linear, delay: 1_000} )自定义策略还允许根据失败原因做条件判断例如对不可重试的错误直接返回0见 backoff.ex 的文档示例。策略注册表由BullMQ.Backoff模块以 Agent 形式维护BullMQ 应用启动时自动拉起start_link/1。自定义任务 ID默认情况下每个任务由队列的 Redis 计数器#{base}:id见 keys.ex分配一个全局自增唯一 ID。你也可以通过:job_id显式指定# Using custom job ID BullMQ.Queue.add(users, process-user, %{user_id: 123}, connection: :my_redis, job_id: user-123-process ) # Adding the same job ID again will return the existing job重复添加相同job_id的任务时队列会返回已存在的任务幂等添加。这一特性天然适合确保某个业务操作只入队一次的场景——但要注意它与下方去重Deduplication的区别job_id去重是永久性的任务不删除ID 就始终占用而 Deduplication 有 TTL / 完成即释放的机制。去重Deduplication防止重复任务入队去重用于阻止相同业务逻辑的任务被重复加入队列。完整机制见 Deduplication Guide这里给出三种模式的用法与取舍。三种去重模式# Simple mode: deduplicate until job completes BullMQ.Queue.add(tasks, process, %{}, connection: :my_redis, deduplication: %{id: unique-task-id} ) # Throttle mode: deduplicate for 5 seconds BullMQ.Queue.add(tasks, process, %{}, connection: :my_redis, deduplication: %{id: unique-task-id, ttl: 5_000} ) # Debounce mode: replace and extend TTL BullMQ.Queue.add(tasks, process, %{data: latest}, connection: :my_redis, delay: 5_000, deduplication: %{id: unique-task-id, ttl: 5_000, extend: true, replace: true} )模式配置行为适用场景Simple仅id任务完成或失败前相同id的重复添加被忽略长任务防并发重入Throttleidttl在 TTL毫秒窗口内去重限制高频触发Debounceidttlextendreplace每次新任务都刷新 TTL并替换任务数据把多次快速更新合并为一次如搜索索引从类型定义types.ex看去重配置还支持keep_last_if_active等字段去重键在 Redis 中以#{base}:de:#{dedup_id}形式存储见 keys.ex 的dedup/2因此get_deduplication_job_id/3、remove_deduplication_key/3等管理接口也都是围绕这一键实现的。相关测试在 postgres_test.exs 与 redis_test.exs 中验证了只有任务所有者owner能删除去重键这一行为。管理去重状态# 查看当前去重由哪个任务触发 {:ok, job_id} BullMQ.Queue.get_deduplication_job_id(my-queue, dedup-id, connection: :redis ) # 提前移除去重键允许新任务入队 {:ok, 1} BullMQ.Queue.remove_deduplication_key(my-queue, dedup-id, connection: :redis )一个常见模式是任务开始处理时就移除去重键这样当前任务运行期间允许下一个任务排队实现串行执行、不丢事件的语义。完整的 Worker 内移除示例可参考 Deduplication Guide。最佳实践提醒去重id应代表逻辑操作本身如sync-user-#{user_id}而非随机值Simple 模式适合绝不允许并发双跑的关键操作Throttle 模式适合节流Debounce 模式适合高频更新合并。LIFO 处理让新任务插队默认队列是FIFO先进先出。设置lifo: true后任务会被插入等待列表的头部表现为后进先出的栈式行为# This job will be processed before older jobs BullMQ.Queue.add(urgent, urgent-task, %{}, connection: :my_redis, lifo: true )适合最新状态优先处理的场景例如只关心最新配置的任务。注意 LIFO 仅影响普通等待列表延迟、优先任务仍走各自的集合结构。任务清理控制完成/失败后的保留策略默认情况下BullMQ不自动删除已完成或失败的任务remove_on_complete、remove_on_fail默认false任务会一直保留在#{base}:completed、#{base}:failed有序集合中供查询审计。你可以逐任务配置自动清理# Remove immediately when completed BullMQ.Queue.add(temporary, temp-job, %{}, connection: :my_redis, remove_on_complete: true ) # Keep last 100 completed jobs BullMQ.Queue.add(with-history, job, %{}, connection: :my_redis, remove_on_complete: %{count: 100} ) # Remove completed jobs older than 1 hour (in ms) BullMQ.Queue.add(time-limited, job, %{}, connection: :my_redis, remove_on_complete: %{age: 3_600_000} ) # Remove failed jobs after keeping 50 BullMQ.Queue.add(cleanup-failures, job, %{}, connection: :my_redis, remove_on_fail: %{count: 50} ) # Keep completed jobs but remove failed ones BullMQ.Queue.add(success-matters, job, %{}, connection: :my_redis, remove_on_complete: false, remove_on_fail: true )配置取值语义对应 types.ex 中的keep_jobs类型取值行为true任务一完成/失败立即删除false永不自动删除默认%{count: n}只保留最近n个任务超出部分删除%{age: ms}删除超过指定毫秒数的任务在 scripts.ex 的任务完成脚本中可以看到默认清理值的实现未配置时传入%{count -1}表示不清理。另外Worker 端也支持队列级的remove_on_complete/remove_on_fail默认值见 types.ex 的worker_opts可对整条队列统一生效。批量添加add_bulk 的高吞吐之路逐个调用add/4会产生大量网络往返。BullMQ.Queue.add_bulk/3将多个任务打包提交通过 Redispipeline MULTI/EXEC 事务保证批量原子性实测吞吐可达逐条添加的 10 倍以上。基本用法任务以三元组列表{name, data, opts}传入jobs [ {email, %{to: user1example.com}, [priority: 1]}, {email, %{to: user2example.com}, []}, {email, %{to: user3example.com}, [delay: 60_000]} ] # All jobs are added atomically - either all succeed or none do {:ok, added_jobs} BullMQ.Queue.add_bulk(emails, jobs, connection: :my_redis)从 queue.ex 的实现可以看到add_bulk/3的分流策略标准任务无delay、无priority即delay 0 and priority 0走后端优化的批量命令路径pipelined Lua 脚本见 scripts.ex每条命令通过 Msgpax 打包参数后一次execute_pipeline提交延迟/优先任务自动回退到逐条添加路径因为它们需要写入各自的集合结构若部分任务失败返回{:error, {:partial_failure, results}}其中results保留了每个任务的成败明细。高性能批量添加连接池并行当需要一次性灌入大批量任务1 万级以上时可创建多个 Redis 连接并行分发# Create a pool of 8 connections pool for i - 1..8 do name :redis_pool_#{i} {:ok, _} BullMQ.RedisConnection.start_link(name: name, host: localhost) name end # Add 100,000 jobs at ~60,000 jobs/sec # Each chunk is added atomically jobs for i - 1..100_000, do: {job, %{index: i}, []} {:ok, added} BullMQ.Queue.add_bulk(my-queue, jobs, connection: :redis, connection_pool: pool )注意示例中为演示手工启动了连接生产环境应把池内连接加入应用的 supervision tree{BullMQ.RedisConnection, name: ..., host: ...}以保证生命周期管理与断线重连。使用connection_pool时每个连接各自的批次是原子的但跨连接的整体操作不保证全有或全无见 queue.ex 文档说明。批量选项OptionDefaultDescriptionpipelinetrue使用流水线提升效率atomictrue将批次包裹在 MULTI/EXEC 事务中。设为false时使用普通 pipeline略快但非原子。配合connection_pool时每个批次独立原子connection_poolnil用于并行处理的连接列表max_pipeline_size10_000每个 pipeline 批次的最大任务数性能基准当前仓库实测根据 Benchmarks测试环境Apple M2 Pro、Redis 7.x、Elixir 1.18以 10 万任务为样本的批量添加吞吐如下方式连接数吞吐加速比逐条顺序15,700 j/s1.0x原子批量124,000 j/s4.2x原子批量239,000 j/s6.8x原子批量454,000 j/s9.5x原子批量858,000 j/s10.2x原子批量1656,000 j/s9.8x关键结论MULTI/EXEC 事务带来约 4 倍提升连接数在 4~8 之间达到饱和峰值默认配置atomic: true、max_pipeline_size: 10_000即处于最佳区间。注意这些数据是仓库在特定硬件环境下的实测实际吞吐会受网络延迟、Redis 配置、任务数据体积影响。你可以用mix run benchmark/add_job_benchmark.exs复现见 benchmark/add_job_benchmark.exs。全部选项速查表以下是BullMQ.Queue.add/4支持的全部任务选项合并 job_options.md 与 types.ex 中job_opts类型定义OptionTypeDefaultDescriptionconnectionatom/pidrequiredRedis 连接prefixstringbull队列键前缀priorityinteger0数值越小优先级越高delayinteger0延迟执行毫秒数attemptsinteger1总尝试次数含首次backoffmapnil重试退避配置lifobooleanfalse插入等待列表头部后进先出job_idstringauto自定义任务 IDdeduplicationmapnil去重配置见 guideremove_on_completebool/mapfalse完成后清理配置remove_on_failbool/mapfalse失败后清理配置keep_logsintegernil保留的最大日志条目数配合Job.log/3timestampintegernow任务创建时间戳毫秒telemetry_metadatastringnil序列化的 trace 上下文telemetry 自动设置omit_contextbooleanfalse跳过 trace 上下文传播此外job_opts类型还包含 Flow 相关选项parent父任务引用%{id: ..., queue: ...}、fail_parent_on_failure、ignore_dependency_on_failure、remove_dependency_on_failure以及timeout任务超时毫秒数、repeat定时任务配置见 Job Schedulers——这些面向流程编排与定时调度的选项在对应指南中有更完整的讲解。从源码看配置如何落盘理解选项的底层存储方式有助于排查配置是否真正生效类问题Job Hash 持久化任务创建后name、data、opts、timestamp、delay、priority等字段被写入 Redis 的任务 Hash见 job.ex 的to_redis/1其中opts以短键 JSON 编码存储保证与 Node.js 生态互操作。状态集合分流普通任务进:wait列表、延迟任务进:delayed有序集合、优先任务进:prioritized有序集合、完成/失败任务分别进:completed/:failed集合keys.ex。Worker 通过moveToActive等 Lua 脚本按优先 → 等待 → 延迟的顺序取任务。去重键独立存储去重键#{base}:de:#{id}独立于任务 Hash因此remove_deduplication_key/3可以在不触碰任务本体的前提下解除去重。下一步学习 Workers 如何消费这些配置过的任务通过 Rate Limiting 控制队列吞吐上限用 Job Schedulers 创建定时循环任务深入 Deduplication 掌握去重全貌配置 Telemetry 实现可观测性与分布式追踪从 Getting Started 回顾完整的队列与 Worker 搭建流程。赞分享后端消息队列任务调度【免费下载链接】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 Job Scheduler 重复选项Repeat Options全面指南startDate、endDate、tz、limit 与 immediately 实战详解BullMQ Job Scheduler 重复选项Repeat Options全面指南startDate、endDate、tz、limit 与 immed后端消息队列任务调度Got 请求选项Options完全指南从 Options 类到全部配置项解析Got 请求选项Options完全指南从 Options 类到全部配置项解析 导读 Got 是 Node.js 生态中广受欢迎的 HTTP 请求库其强大后端网络osquery 部署配置完全指南从 options、Schedule 到 Query Packs 与配置规范详解osquery 部署配置完全指南从 options、Schedule 到 Query Packs 与配置规范详解 本指南以 osquery 官方部署文档 do观测代理网络安全上一篇如何永久保存微信聊天记录WeChatMsg完整解决方案指南下一篇Teleport Windows CA 拆分RFD 0239深度解读Desktop Access 从 User CA 到独立 Windows CA 的演进创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考