ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

BullMQ 批处理实战指南:addBulk、FlowProducer.addBulk 与单任务批量的三种选型

BullMQ 批处理实战指南:addBulk、FlowProducer.addBulk 与单任务批量的三种选型 后端消息队列任务调度【免费下载链接】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 中批处理Batches一词可以指三种完全不同的场景批量入队相互独立的普通任务Queue.addBulk、一次性创建多棵父子依赖的任务树FlowProducer.addBulk以及让单个 Worker 回调一次拿到多条任务处理的 Pro 原生批处理。三种 API 解决的是不同的问题选错会导致重试语义、事件语义和资源开销都与预期不符。读完本文你将掌握三种批处理模式各自的适用场景、完整可运行的代码示例、底层实现原理Redis 流水线与 Lua 脚本、PostgreSQL 后端的独立快路径以及如何通过QueueEvents实时跟踪批量任务的进度与完成情况。Batches 的三种含义先选对 API目标所属版本API高效、原子地批量入队大量相互独立的任务开源版Queue.addBulk一次性创建多棵父子任务树开源版FlowProducer.addBulk在一个Worker 回调里处理多个等待中的任务BullMQ ProPro batchesWorkerProbatch选项需要特别澄清一个常见误区Worker 的concurrency选项确实会让多个任务并行执行但每次处理器调用仍然只收到一个job——这与 Pro 批处理中job.getBatch()一次返回最多batch.size个任务是完全不同的机制详见 Workers: concurrency。并行度解决的是吞吐批处理解决的是单次回调处理多条数据。批量添加独立任务Queue.addBulk当你已经拥有一份工作清单且每一项都应该保持为独立任务拥有独立的重试次数、独立的事件、独立的完成状态时使用Queue.addBulkimport { Queue } from bullmq; const queue new Queue(invoices, { connection }); const jobs await queue.addBulk([ { name: send, data: { invoiceId: A-100 } }, { name: send, data: { invoiceId: A-101 } }, { name: send, data: { invoiceId: A-102 } }, ]);与在循环里逐个调用add相比addBulk显著减少了与 Redis 之间的往返次数round-trips。语言示例与原子性说明可参考 Adding jobs in bulk该文档还给出了 Python 与 Rust 的等价写法from bullmq import Queue queue Queue(paint) jobs await queue.addBulk([ { name: jobName, data: { paint: car } }, { name: jobName, data: { paint: house } }, { name: jobName, data: { paint: boat } } ])use bullmq::{Queue, QueueOptions}; let queue Queue::new(paint, QueueOptions::default()).await?; let jobs queue.add_bulk(vec![ (jobName.into(), serde_json::json!({paint: car}), None), (jobName.into(), serde_json::json!({paint: house}), None), (jobName.into(), serde_json::json!({paint: boat}), None), ]).await?;原子性要么全部入队要么一个都不入addBulk的调用结果只有两种成功或失败且成功时所有任务都会被加入队列失败时则一个都不会——不存在加了一半的中间状态。这一点从源码实现可以清楚印证Queue.addBulk会把队列级默认选项与每个任务的opts合并{ ...this.jobsOpts, ...job.opts, jobId }并携带遥测telemetry传播元数据随后委托给Job.createBulk见 src/classes/queue.tsJob.createBulk先为每个条目构造Job实例、序列化数据并调用validateOptions做校验再一次性调用后端的addJobs见 src/classes/job.ts在 Redis 后端addJobs将每个任务的入队命令按序压入同一条 pipeline最后通过pipeline.exec()一次性执行若任一命令出错则整体抛出异常见 src/classes/redis-queue-backend.ts。每条命令最终都会落到 src/commands/addStandardJob-9.lua 这类 Lua 脚本中执行脚本内部负责递增任务计数器、写入任务键、按 LIFO/FIFO/优先级入队、处理去重与父子依赖等完整逻辑。PostgreSQL 后端的独立快路径在 PostgreSQL 后端src/postgres/postgres-queue-backend.ts中addJobs还有一个专门的性能优化当所有任务都不带父任务、不涉及去重parentId null parentQueue null dedupId null时会走快路径使用集合化的add_jobs_bulkSQL一次INSERT 一次事件INSERT替代逐行执行的流引擎对于普通addBulk场景明显更快。这也解释了为什么文档推荐先判断任务是否独立再决定用addBulk还是逐条add。传参细节每个元素与Queue.add签名一致addBulk接收的数组元素由name、data、opts三个属性组成签名与Queue.add完全一致因此opts里可以携带jobId、delay、priority、attempts、removeOnComplete等所有标准任务选项对应源码中的BulkJobOptions类型。需要注意data的 JSON 序列化约束类实例在 Worker 端会丢失原型方法源码注释中明确指出了这一 caveat。批量入队时如果opts未指定jobId任务 ID 将由 Redis 中的计数器KEYS[4]见addStandardJob-9.lua中的INCR调用自动生成。批量创建多个流程FlowProducer.addBulk当每个工作单元本身是一个流程一个父任务 若干子任务且希望一次调用创建多棵任务树时使用FlowProducer.addBulkimport { FlowProducer } from bullmq; const flowProducer new FlowProducer({ connection }); await flowProducer.addBulk([ { name: import-file, queueName: imports, data: { fileId: first.csv }, children: [ { name: parse, queueName: import-steps, data: { fileId: first.csv }, }, { name: index, queueName: import-steps, data: { fileId: first.csv }, }, ], }, { name: import-file, queueName: imports, data: { fileId: second.csv }, children: [ { name: parse, queueName: import-steps, data: { fileId: second.csv }, }, { name: index, queueName: import-steps, data: { fileId: second.csv }, }, ], }, ]);这个调用同样是原子的要么全部流程创建成功要么全部失败。在 src/classes/flow-producer.ts 的实现中addBulk首先执行validateFlowJobs校验整棵树的合法性然后等待后端就绪再以单次addFlow事务Redis 适配器为单个MULTI事务见 src/classes/redis-queue-backend.ts 的注释一次性写入横跨多个队列的整批任务树。详细示例见 Adding flows in bulk。把一批数据建模成单个任务当这一批数据必须共享同一个重试次数、超时时间或完成结果时正确的做法是把它们放进一个任务的 payload而不是自造一个 Worker 端的批处理器await queue.add(send-batch, { invoiceIds: [A-100, A-101, A-102], });import { Worker } from bullmq; const worker new Worker( invoices, async job { const invoiceIds: string[] job.data.invoiceIds; for (const [index, invoiceId] of invoiceIds.entries()) { await sendInvoice(invoiceId); await job.updateProgress({ completed: index 1, total: invoiceIds.length, }); } return { sent: invoiceIds.length }; }, { connection }, );这个模式的核心价值在于单一结局整批数据要么整体重试、要么整体超时完成状态只有一个可观测进度通过job.updateProgress({ completed, total })在循环内逐步上报进度配合下文提到的QueueEvents即可在 UI 中实时渲染进度条返回聚合结果处理器返回{ sent: invoiceIds.length }作为该任务的最终结果存储。但要注意当每一项都需要独立重试、独立失败追踪时请优先使用addBulk或 flows否则一个子项失败会导致整批重跑。Pro 原生批处理一个回调处理多条任务开源版的 Worker 每次只把一个任务交给处理器。原生Worker 端批处理是 BullMQ Pro 的特性其用法如下import { WorkerPro } from taskforcesh/bullmq-pro; const worker new WorkerPro( invoices, async job { const batch job.getBatch(); // 示例对整个 batch 发一次批量 DB/API 调用 await sendInvoices(batch.map(batchedJob batchedJob.data)); }, { connection, batch: { size: 10 }, }, );启用 Pro 批处理前必须了解它与普通任务的差异详见 BullMQ Pro: Batches失败与事件语义不同批处理引入了一个dummy包装任务失败标记通过setAsFailed完成事件需要通过QueueEventsPro监听批聚合成型条件batch配置支持size目标批量大小以及minSize/timeout最小规模与最长等待时间来控制攒批行为组亲和性group affinity支持让相关任务优先聚合到同一批。Pro 批处理适合一次批量 DB/API 调用即可完成多条记录处理的场景例如批量写库、批量发通知能显著降低每条记录分别建立连接的开销但正因如此它的失败粒度与事件模型与开源版不同选型时需结合业务对单条失败追踪的诉求慎重评估。批量任务的实时进度用 QueueEvents 监听无论你选择了addBulk、单个批量任务还是 Pro 批处理实时 UI 都应该通过QueueEvents来监听让一个进程就能观察到所有 Worker 上报的完成事件与progress事件import { QueueEvents } from bullmq; const queueEvents new QueueEvents(invoices, { connection }); queueEvents.on(progress, ({ jobId, data }) { // data 即 job.updateProgress 传入的 { completed, total } renderProgress(jobId, data); }); queueEvents.on(completed, ({ jobId, returnvalue }) { renderDone(jobId, returnvalue); });事件流Redis Streams由 Worker 侧在任务状态迁移时写入QueueEvents以发布/订阅方式消费因此即使多个 Worker 进程并行处理同一批任务UI 侧也只需一个监听实例即可获得全局视图。选型决策速查你的需求推荐方案关键特性大量相互独立的任务高效/原子入队Queue.addBulk减少 Redis 往返全有或全无每项独立重试与事件多棵父子依赖树一次性创建FlowProducer.addBulk原子创建跨队列任务树子任务先于父任务处理一批数据共享重试/超时/完成结果单任务 数组 payload单一结局updateProgress上报进度单个 Worker 回调批量处理多条WorkerProbatchProjob.getBatch()size/minSize/timeout不同事件语义相关阅读Adding jobs in bulkTypeScript / Python / Rust 示例与原子性说明Adding flows in bulkEvents / real-time updatesWorkers: concurrencyBullMQ Pro: Batches实现与测试佐证src/classes/queue.ts、src/classes/job.ts、src/classes/redis-queue-backend.ts、src/commands/addStandardJob-9.lua、tests/bulk.test.ts、tests/flow.test.ts赞分享后端消息队列任务调度【免费下载链接】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 跨队列批量原子添加任务FlowProducer.addBulk 原理与实战指南BullMQ 跨队列批量原子添加任务FlowProducer.addBulk 原理与实战指南 导读 本文讲解 BullMQ 中一种特殊的生产者模式 一次性后端消息队列任务调度BullMQ 批量添加任务addBulk原子性、性能优化与多语言实战指南BullMQ 批量添加任务addBulk原子性、性能优化与多语言实战指南 导读 在 BullMQ 中当需要一次性向队列写入多个任务且要求要么全部成功后端消息队列任务调度BullMQ FlowProducer.addBulk 深度指南跨队列原子批量添加 Flow 任务树BullMQ FlowProducer.addBulk 深度指南跨队列原子批量添加 Flow 任务树 导读 在 BullMQ 中Flow 是一棵由存在父子依后端消息队列任务调度上一篇obfuscator-llvm 混淆编译器插件使用教程下一篇DeepHermes-Egregore-v1-RLAIF-8b-Atropos-i1-GGUF高级配置多GPU部署与推理优化终极指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表