新闻详情

新闻详情

首页 / 资讯中心 / 详情

BullMQ Elixir 队列事件订阅实战指南:基于 Redis Streams 的实时任务生命周期监控

发布时间:2026/9/25 3:00:07来源:尧图网络
BullMQ Elixir 队列事件订阅实战指南:基于 Redis Streams 的实时任务生命周期监控
后端消息队列任务调度【免费下载链接】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 客户端的队列事件Queue Events机制展开讲解如何通过 Worker 回调与BullMQ.QueueEvents两种方式实时监听任务从入队、等待、执行到完成/失败的完整生命周期。阅读完本文后你将掌握事件类型与消息格式、QueueEvents 的全部配置项、订阅/退订模式、结构化 Handler 模块写法以及如何在监督树中集成事件监听从而搭建监控面板、告警通知等实时响应系统。事件机制总览Redis Streams 驱动的工作流BullMQ Elixir 客户端的队列事件完全基于Redis Streams实现。每一条队列都对应一个以{prefix}:{queue_name}:events命名的 Stream事件流的 key 生成逻辑见 keys.ex当任务状态发生变化时队列内部会向该 Stream 追加一条事件记录事件监听方则通过阻塞式读取持续消费这些记录。从底层实现看事件流读写被抽象为BullMQ.Backendbehaviour 的两个回调见 backend.expublish_event/3向事件流发布一条事件Redis 后端实际执行的命令是XADD {events} MAXLEN ~ {max_events} * {field...}见 backends/redis.ex其中MAXLEN ~表示对事件流进行近似裁剪防止 Stream 无限增长read_events/3从指定事件 ID 之后阻塞读取新事件Redis 后端使用XREAD BLOCK {block_ms} STREAMS {events} {id}见 backends/redis.exblock_ms默认 5000ms定义在 queue_events.ex。也就是说事件订阅是跨进程、跨语言可见的任何向事件流写入的任务生命周期变化无论由 Elixir Worker、Node.js Worker 还是其他语言客户端产生都能被监听者实时消费这也是本文末尾Node.js 兼容性一节的底层依据。事件流长度可以通过队列的max_len_events配置控制默认 10000 条见 queue.ex 中Queue.update_meta/2的说明并写入队列 meta 的opts.maxLenEvents字段。两种事件接收方式的选择在 BullMQ Elixir 中接收任务事件有两种途径各有适用场景方式监听范围适用场景Worker Callbacks仅当前 Worker 自己处理的任务在任务处理逻辑旁直接响应完成/失败/进度等事件QueueEvents队列中全部任务含其他 Worker、Node.js Worker 处理的任务集中式监控、告警、统计面板Worker Callbacks 适合处理即响应例如发送邮件任务完成后立即记日志、失败时马上告警QueueEvents 适合全局观测例如统计队列整体吞吐、检测停滞任务、做统一审计。二者可以同时使用互不冲突——Worker 回调面向单个消费进程QueueEvents 面向整条队列。Worker Callbacks最直接的本地事件响应当只需要响应当前 Worker 处理的任务事件时最简单的方式是在BullMQ.Worker.start_link/1中直接传入回调函数{:ok, worker} BullMQ.Worker.start_link( queue: emails, connection: :my_redis, processor: process/1, on_completed: fn job, result - Logger.info(Job #{job.id} completed: #{inspect(result)}) end, on_failed: fn job, reason - Logger.error(Job #{job.id} failed: #{reason}) end, on_active: fn job - Logger.debug(Job #{job.id} started) end )从 worker.ex 的选项 schema 可以看到Worker 完整支持以下事件回调选项回调签名触发时机on_completedfn job, result - ... end任务成功完成on_failedfn job, reason - ... end任务失败on_errorfn error - ... endWorker 内部错误on_activefn job - ... end任务进入 active开始处理on_progressfn job, progress - ... end任务进度更新on_stalledfn job_id - ... end任务被判定为停滞on_lock_renewal_failedfn job_ids - ... end任务锁续期失败受影响任务会以{:lock_lost, job_id}被取消防止重复处理其中on_progress的触发链路值得注意处理器内部调用Worker.update_progress/2时会先通过 Backend 更新 Redis 中的进度并同时向事件流发布 progress 事件供 QueueEvents 消费随后再以GenServer.cast通知 Worker 进程触发本地on_progress回调见 worker.ex。因此Worker 回调与QueueEvents 事件在进度场景下是并行生效的。更完整的 Worker 使用与配置说明可参阅 Workers 指南。QueueEvents全队列集中式事件监听当需要监听整条队列的事件——包括由其他进程、其他语言 Worker 处理的任务——请使用BullMQ.QueueEvents。它是基于 GenServer 实现的常驻监听进程源码位于 queue_events.ex。启动与订阅{:ok, events} BullMQ.QueueEvents.start_link( queue: emails, connection: :my_redis ) # 订阅当前进程默认 pid 即 self() BullMQ.QueueEvents.subscribe(events) # 接收事件 receive do {:bullmq_event, :completed, data} - IO.puts(Job #{data[jobId]} completed!) {:bullmq_event, :failed, data} - IO.puts(Job #{data[jobId]} failed: #{data[failedReason]}) {:bullmq_event, :waiting, data} - IO.puts(Job #{data[jobId]} waiting) end内部消费循环源码视角QueueEvents的监听核心是一个阻塞读取 → 分发 → 继续阻塞的循环启动时autorun: true通过Backend.reconnect_blocking/1建立专用阻塞连接见 queue_events.ex连接失败时每 5 秒自动重试每次以Task.async发起一次Backend.read_events/3阻塞读取超时 5000ms见 queue_events.ex读到事件批次后由process_events/2逐条解析将 Stream 字段拆成 KV map、把字符串事件名转换为 atom然后向所有订阅进程广播消息、并调用可选的 Handler 模块见 queue_events.ex读取超时返回nil或处理完毕后立即调度下一轮消费消费者任务崩溃时延迟 1 秒后重试见 queue_events.ex。这套设计保证了事件监听不会阻塞 GenServer 主循环且 Redis 连接意外断开后能自动恢复监听。事件类型与数据字段QueueEvents 会发出下表全部事件。这些事件名同时也出现在 queue_events.ex 的解析函数中事件说明数据字段:added任务已加入队列jobId,name:waiting任务等待处理jobId:active任务开始处理jobId,prev:progress任务进度更新jobId,data:completed任务成功完成jobId,returnvalue,prev:failed任务失败jobId,failedReason,prev:delayed任务被延迟含退避重试jobId,delay:stalled任务被检测为停滞jobId:removed任务被移除jobId,prev:drained队列无等待任务无数据:paused队列被暂停无数据:resumed队列恢复无数据此外从 queue_events.ex 的解析实现还可以看到Elixir 客户端实际还支持:duplicated、:deduplicated、:retries_exhausted、:waiting_children、:cleaned等扩展事件其中:cleaned对应队列清理操作未在表中列出的字符串事件名会直接转换为 atom。消息格式统一的 3 元组所有事件都以如下 3 元组形式投递到订阅进程的邮箱{:bullmq_event, event_type, event_data}其中event_type是 atom:completed、:failed等event_data是字符串键的 map键名与 Node.js BullMQ 保持一致jobId、failedReason等便于跨语言复用。一个完整的 completed 事件示例{:bullmq_event, :completed, %{ event completed, jobId abc123, returnvalue null, prev active }}注意returnvalue是任务返回值序列化后的字符串prev表示任务进入该状态前的前一状态如active→completed、waiting→active集成测试正是依据event_data[prev] waiting来断言 active 事件的正确性见 queue_events_integration_test.exs。QueueEvents 完整配置项QueueEvents.start_link/1的全部选项由 NimbleOptions schema 校验见 queue_events.ex{:ok, events} BullMQ.QueueEvents.start_link( queue: my_queue, # 队列名必填 connection: :my_redis, # 后端连接引用必填可为 atom 名、pid 或 {:via, registry} 元组 prefix: bull, # Redis key 前缀默认 bull与 Node.js 兼容 backend: BullMQ.Backends.Redis, # 后端模块默认取应用配置或 Redis autorun: true, # 是否立即开始监听默认 true last_event_id: $, # 起始事件 ID默认 $ 仅接收新事件 handler: MyEventHandler, # Handler 模块可选实现 BullMQ.QueueEvents.Handler handler_state: %{}, # 传给 Handler 的初始状态可选 name: nil # 可选用于注册 GenServer 进程名 )各选项要点prefix所有 Redis key 的公共前缀默认bull。必须与生产方Queue / Worker使用的前缀一致否则监听不到事件。Elixir 客户端的 key 生成统一由 keys.ex 的Keys.new/2完成。last_event_id事件流游标。$默认表示只接收启动之后的新事件传入数字 ID 或0可从历史事件开始回放。Redis 后端直接用该值作为XREAD的游标见 backends/redis.ex。autorun: false先启动但不监听之后调用BullMQ.QueueEvents.run/1手动开始。handlerhandler_state结构化事件处理见下一节。多订阅者广播模式同一个 QueueEvents 实例可被多个进程同时订阅事件会广播给每一个订阅者。订阅/退订 API 默认作用于当前进程self()也支持显式指定其他 pid{:ok, events} BullMQ.QueueEvents.start_link( queue: tasks, connection: :my_redis ) # 订阅多个进程 BullMQ.QueueEvents.subscribe(events) # 订阅 self() BullMQ.QueueEvents.subscribe(events, other_pid) # 结束订阅 BullMQ.QueueEvents.unsubscribe(events) BullMQ.QueueEvents.unsubscribe(events, other_pid)从实现上看subscribe/2会通过Process.monitor/1监控订阅进程订阅进程退出后自动从订阅列表中移除见 queue_events.ex 与 queue_events.ex无需手动清理。集成测试也验证了多订阅者广播行为主进程与另一个 spawn 的进程都能收到同一事件见 queue_events_integration_test.exs。Handler 模块模式结构化事件处理对生产级应用推荐实现BullMQ.QueueEvents.Handlerbehaviour把事件处理逻辑封装为独立模块并通过返回值维护增量状态defmodule MyApp.QueueHandler do behaviour BullMQ.QueueEvents.Handler require Logger impl true def handle_event(:completed, %{jobId id, returnvalue value}, state) do Logger.info(Job #{id} completed with: #{value}) {:ok, state} end impl true def handle_event(:failed, %{jobId id, failedReason reason}, state) do Logger.error(Job #{id} failed: #{reason}) MyApp.Alerts.notify_failure(id, reason) {:ok, state} end impl true def handle_event(:drained, _data, state) do Logger.info(Queue drained - no more waiting jobs) {:ok, state} end impl true def handle_event(_event, _data, state) do # 忽略其他事件 {:ok, state} end end # 使用 handler {:ok, events} BullMQ.QueueEvents.start_link( queue: tasks, connection: :my_redis, handler: MyApp.QueueHandler, handler_state: %{notifications_sent: 0} )关于该 behaviour 的几点实现细节见 queue_events.ex回调签名统一为handle_event(event_atom, data_map, state)必须返回{:ok, new_state}或{:error, reason}返回的new_state会作为下一次调用的state传入即状态随事件流累积事件 ID 也会随处理进度推进除handle_event/3外还定义了可选的init/1回调use BullMQ.QueueEvents.Handler宏会提供默认实现init/1返回{:ok, nil}handle_event/3直接透传 state你只需覆写需要的事件分支集成测试中验证了 handler 会被依次调用并累积状态见 queue_events_integration_test.exs。延迟启动与优雅关闭延迟启动autorun: false某些场景下需要先完成其他初始化、再开始监听事件{:ok, events} BullMQ.QueueEvents.start_link( queue: tasks, connection: :my_redis, autorun: false # 先不监听 ) # 稍后手动开始监听 BullMQ.QueueEvents.run(events)集成测试确认autorun: false时即使队列有任务事件也收不到任何消息见 queue_events_integration_test.exs。优雅关闭# 关闭事件监听取消进行中的消费任务并关闭阻塞连接 BullMQ.QueueEvents.close(events)close/1会取消进行中的消费者任务、清理后端连接并把进程标记为closing此后晚到的消费结果会被安全忽略见 queue_events.ex 与 queue_events.ex。terminate/2回调中也会执行同样的清理确保进程退出时连接被正确释放。实战示例QueueEvents 驱动的监控面板综合前面的 API可以构建一个实时统计队列完成/失败/活跃任务数的监控 GenServerdefmodule MyApp.QueueMonitor do use GenServer def start_link(queue_name) do GenServer.start_link(__MODULE__, queue_name, name: __MODULE__) end def init(queue_name) do {:ok, events} BullMQ.QueueEvents.start_link( queue: queue_name, connection: :my_redis ) BullMQ.QueueEvents.subscribe(events) {:ok, %{ events: events, completed: 0, failed: 0, active: 0 }} end def handle_info({:bullmq_event, :completed, _data}, state) do {:noreply, %{state | completed: state.completed 1, active: state.active - 1}} end def handle_info({:bullmq_event, :failed, _data}, state) do {:noreply, %{state | failed: state.failed 1, active: state.active - 1}} end def handle_info({:bullmq_event, :active, _data}, state) do {:noreply, %{state | active: state.active 1}} end def handle_info({:bullmq_event, _event, _data}, state) do {:noreply, state} end def get_stats do GenServer.call(__MODULE__, :get_stats) end def handle_call(:get_stats, _from, state) do {:reply, %{ completed: state.completed, failed: state.failed, active: state.active }, state} end end该模式可直接扩展为失败告警推送结合:failed事件、队列空闲感知:drained事件、延迟任务趋势统计:delayed事件等。监督树集成将QueueEvents纳入应用监督树可获得崩溃自动重启等 OTP 保障。由于QueueEvents是基于 GenServer 的子进程规范可以与其他子进程并列声明children [ {Redix, name: :my_redis, host: localhost}, {BullMQ.QueueEvents, queue: important-queue, connection: :my_redis, handler: MyApp.ImportantQueueHandler } ]监督树模式下建议配合handler模块使用这样进程重启后监听逻辑与状态初始化完全由 Handler 负责应用自身无需关心事件消费细节。Node.js 兼容性Elixir 的 QueueEvents 与 Node.js BullMQ完全兼容由于双方共享同一套 Redis key 约定默认前缀bull事件流 key 为bull:{queue}:events与同一事件数据格式jobId、failedReason等字符串键Node.js Worker 产生的事件可被 Elixir QueueEvents 接收反之亦然。这使得混合技术栈的团队可以用 Elixir 编写全局监控/告警服务监听由 Node.js 消费的队列在 Elixir 服务与 Node.js 服务之间共享同一条队列的事件流而无需额外的消息中间件。下一步学习路径深入了解 Worker 的全部回调与并发模型Workers 指南为队列事件接入指标与分布式追踪Telemetry 指南结合定时任务理解:delayed、:added事件的来源Job Schedulers 指南深入事件流底层实现阅读 queue_events.ex、backend.ex 与 keys.ex以及集成测试 queue_events_integration_test.exs赞分享后端消息队列任务调度【免费下载链接】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 for Elixir 实战指南基于 Redis 与 PostgreSQL 的分布式任务队列BullMQ for Elixir 实战指南基于 Redis 与 PostgreSQL 的分布式任务队列 BullMQ for Elixir 是开源多语言任务后端消息队列任务调度BullMQ for Elixir 快速上手指南从 Redis 连接到队列、Worker 与任务生命周期BullMQ for Elixir 快速上手指南从 Redis 连接到队列、Worker 与任务生命周期 本指南以 BullMQ Elixir 端口为背景手后端消息队列任务调度BullMQ for Elixir 入门指南基于 OTP 与 Redis 的全功能任务队列实战BullMQ for Elixir 入门指南基于 OTP 与 Redis 的全功能任务队列实战 本篇技术指南以 docs/gitbook/elixir/int后端消息队列任务调度上一篇GitFS实战如何用版本控制文件系统提升开发效率下一篇从源码到自定义绑定node-libcurl 编译与扩展开发完整指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
网站建设高端定制企业官网
RELATED

相关资讯

更多精彩内容,欢迎继续阅读

较早相关资讯

最新相关资讯

基于 MongoDB Atlas 搭建 Mage AI 数据源测试环境的实践指南 2026/9/25 3:43:42

基于 MongoDB Atlas 搭建 Mage AI 数据源测试环境的实践指南

数据工程数据编排ETL任务调度批处理流处理数据集成后端 【免费下载链接】mage-ai 🧙 Build, run, and manage data pipelines for integrating and transforming data. 项目地址: https://gitcode.com/gh_mirrors/ma/mage-ai 点击查看 免费下载 MongoDB…

阅读更多 →
OptiScaler 完整指南:5 种方法在 DLSS、FSR、XeSS 之间自由切换并给游戏补帧 2026/9/25 3:43:42

OptiScaler 完整指南:5 种方法在 DLSS、FSR、XeSS 之间自由切换并给游戏补帧

OptiScaler 完整指南:5 种方法在 DLSS、FSR、XeSS 之间自由切换并给游戏补帧 【免费下载链接】OptiScaler OptiScaler bridges upscaling/frame gen across GPUs. Supports DLSS2/XeSS/FSR2 inputs, replaces native upscalers, enables FSR-FG/XeFG on non-FG titl…

阅读更多 →
cuDF Java JAR 构建与发布流水线:基于 ci-wheel 容器的一站式打包实践 2026/9/25 3:43:42

cuDF Java JAR 构建与发布流水线:基于 ci-wheel 容器的一站式打包实践

数据分析数据工程机器学习 【免费下载链接】cudf cuDF - GPU DataFrame Library 项目地址: https://gitcode.com/gh_mirrors/cu/cudf 点击查看 免费下载 cuDF 的 Java API 依赖一个内嵌静态 libcudf 的 JNI 层,其打包过程涉及静态库编译、Maven classi…

阅读更多 →
深入解析 Basic Computer Games 之 Poker:规则、AI 下注策略与移植注意事项 2026/9/25 3:43:35

深入解析 Basic Computer Games 之 Poker:规则、AI 下注策略与移植注意事项

示例工程 【免费下载链接】basic-computer-games An updated version of the classic "Basic Computer Games" book, with well-written examples in a variety of common MEMORY SAFE, SCRIPTING programming languages. See https://coding-horror.github.io/basic…

阅读更多 →
Tekton Pipeline 控制器启动参数(Controller Flags)完全指南:从注册原理到生产配置 2026/9/25 3:43:35

Tekton Pipeline 控制器启动参数(Controller Flags)完全指南:从注册原理到生产配置

云原生CI/CDDevOps后端 【免费下载链接】pipeline A cloud-native Pipeline resource. 项目地址: https://gitcode.com/gh_mirrors/pipelin/pipeline 点击查看 免费下载 tektoncd/pipeline(本项目 pipeline)随发行版内置了多个控制器二进制&…

阅读更多 →
SendGrid Go SDK 使用指南:基于 sendgrid-go 全面调用 Twilio SendGrid v3 API 2026/9/25 3:43:35

SendGrid Go SDK 使用指南:基于 sendgrid-go 全面调用 Twilio SendGrid v3 API

网络安全 【免费下载链接】sliver Adversary Emulation Framework 项目地址: https://gitcode.com/gh_mirrors/sl/sliver 点击查看 免费下载 本指南以本仓库 vendor/github.com/sendgrid/sendgrid-go/USAGE.md 为骨架,完整讲解用 Go 语言通过 sendgrid-…

阅读更多 →

今日资讯

本周资讯

本月资讯

看完文章仍有疑问?

联系尧图顾问,获取一对一建站咨询

立即免费咨询 📞 400-888-8888
📞 ✉