如何实现Sidekiq父批次等待子批次任务全部完成后执行回调
实现Sidekiq父批次+子Worker的完成回调工作流
嘿,我完全懂你想要的这种工作流——父批次生成一堆子任务,等所有子任务都跑完再触发父级的完成回调,之前看文档没找到对应场景确实挠头😉 其实Sidekiq Batch本身就支持这种需求,可能是你没get到正确的打开方式,下面给你一步步拆解实现方案:
1. 定义父Worker负责创建批次
先写一个父Worker,它的核心职责就是初始化Batch、批量添加子任务,同时绑定完成回调:
class ParentBatchInitiatorWorker include Sidekiq::Worker def perform # 初始化Batch,指定完成后触发的回调Worker和自定义元数据 batch = Sidekiq::Batch.new batch.on(:complete, ParentBatchFinishWorker, batch_tag: "user_data_sync") # 把所有子任务包裹在batch.jobs块里,自动归到当前批次 batch.jobs do # 这里可以动态生成子任务,比如从数据库拉取任务列表 User.all.each do |user| UserDataSyncWorker.perform_async(user.id) end end end end
2. 定义子Worker处理具体业务
子Worker就是你实际要执行的业务单元,每个任务独立运行:
class UserDataSyncWorker include Sidekiq::Worker def perform(user_id) # 这里写你的业务逻辑:比如同步用户第三方数据、生成报表等 user = User.find(user_id) UserDataSyncService.call(user) puts "完成用户#{user.id}的数据同步" end end
3. 定义完成回调Worker
这个Worker会在批次里所有子任务都执行完毕(包括重试成功的情况)后自动触发,用来处理批次收尾逻辑:
class ParentBatchFinishWorker include Sidekiq::Worker def perform(status, options) # status对象包含批次的执行统计:总任务数、成功数、失败数等 total_tasks = status.total successful_tasks = status.successful failed_tasks = status.failed # options是创建Batch时传入的自定义元数据 batch_tag = options[:batch_tag] # 这里写你的收尾逻辑:比如发送通知、更新批次状态、汇总结果等 AdminNotificationService.send_batch_completion_alert( tag: batch_tag, total: total_tasks, success: successful_tasks, fail: failed_tasks ) puts "批次#{batch_tag}执行完成:#{successful_tasks}/#{total_tasks}成功" end end
关键说明
- 核心是
batch.jobs块:所有在这个块内添加的perform_async任务都会被Sidekiq归到当前批次下,自动跟踪执行状态。 - 回调触发时机:只有当批次内所有子任务都完成(成功或重试后成功),才会调用
on(:complete)绑定的Worker,完全符合你的需求。 - 动态任务支持:子任务数量不需要提前固定,不管是循环生成还是从数据库查询得到,都可以在
batch.jobs块里动态添加。
注意事项
- 如果你用的是Sidekiq开源版,需要额外安装
sidekiq-batchgem;Sidekiq Pro自带Batch功能,直接用就行。 - 若子任务可能失败,记得配置合理的重试策略,避免个别任务一直失败导致回调无法触发。
status对象里的failed属性会记录失败的任务数,你可以在回调里针对性处理失败情况。
内容的提问来源于stack exchange,提问作者adil hussain
相关产品推荐
相关产品推荐

