Ruby 应用从 Sidekiq 迁移到 River 完整实战指南:生产者、Worker、配置、测试与切流方案

发布时间:2026/10/12 1:44:29

Ruby 应用从 Sidekiq 迁移到 River 完整实战指南:生产者、Worker、配置、测试与切流方案 任务调度后端【免费下载链接】riverThe polyglot queue: Fast and reliable background jobs in Go, Ruby, Rust, and JS/TS on Postgres or SQLite.项目地址https://gitcode.com/gh_mirrors/river/river点击查看免费下载导读本文以本仓库 Ruby 版 River 客户端riverqueuegem 及其riverqueue-activerecord/riverqueue-sequel驱动为核心系统讲解如何将一个正在使用 Sidekiq 的 Ruby 应用迁移到 River涵盖生产者producers、Worker工作者、配置、测试与存量任务处理全流程。读完本文你将掌握 Sidekiq 与 River 在作业模型上的本质差异、逐项迁移的操作清单、事务边界与重试语义的翻译规则、独立 Worker 进程的启动与优雅停机、中间件/上下文替换、Active Job 与 Action Mailer 的处理、周期任务与唯一性任务映射以及一套可回滚的生产切流方案。River 将作业存储在 Postgres 或 SQLite 中并可以在与业务记录相同的数据库事务内插入作业Worker 在 Ruby 线程中运行。注意一个根本前提Redis 中的 Sidekiq 负载并不是 River 的数据行只更换 gem 并不会迁移既有作业存量数据需要单独处理。为什么从 Sidekiq 迁移到 River对于已经在使用 Postgres 或 SQLite 的应用River 提供以下优势原子入队Atomic enqueueing应用数据变更与作业在同一数据库事务中一起提交或回滚。Sidekiq 的 transactional push 需要等到提交之后但随后的 Redis 写入是独立操作二者之间存在间隙River 直接消除了这个间隙。更少的基础设施复用应用数据库作为作业队列无需为作业队列单独运维 Redis。如果应用的其他功能确实需要 Redis则仍可保留。SQL 可见性可以用普通 SQL 直接查看作业参数args、尝试次数attempts、错误历史errors和保留的结果retained results并与业务记录做关联查询。核心客户端内置能力更多唯一作业unique jobs、周期任务periodic intervals、不消耗尝试次数的 snooze、可续传检查点resumable checkpoints与记录输出recorded output减少应用侧胶水代码与额外扩展。Ruby 与 Go 之间的通路各语言客户端共享 River 的同一套 schema可以用兼容的 kind 和 JSON 参数互相收发作业从而支持 Worker 跨语言逐步迁移。需要明确的是这并不意味着更高的吞吐量或精确一次执行——作业仍然需要幂等设计而且队列流量会给数据库增加负载。迁移前请先核对下文列出的兼容性缺口。自动化迁移的总体流程按编号小节依序推进。以本仓库的 Ruby 源码作为 API 权威依据Go 文档描述的是概念其方法名并非 Ruby 方法。代码示例中使用了应用自有命名如FulfillOrder、AppJobs请显式创建或改造为你的应用对应实现。动手编辑之前先为每个作业类生成一张迁移清单现有类名参数结构队列重试策略生产者中间件/上下文目标 kind切换策略FulfillOrderJoborder_idintegerorders5 次重试Checkout 服务Rails executorfulfill_order先排空旧作业将清单中的每一项分类为直接改写direct rewrite、语义变化semantic change、需要 River Pro、或需要应用自行实现application implementation required。对无法一一对应的行为应记录在案而不是臆造一个 River API。核心 River 没有 Sidekiq 兼容模块也没有perform_async、Sidekiq 风格的 YAML 加载器、fake/inline 测试模式独立的riverqueue-railsgem 提供 Active Job 适配器与 Worker 入口见第 8 节。下文原生 Worker 示例不依赖该 Rails 集成 gem。迁移期间应保证现有系统持续可运行把代码转换、生产切流和Redis 数据搬运分开处理。1. 盘点应用现状Inventory在源码、测试、initializer、脚本、部署清单与依赖中全面搜索rg -n Sidekiq|sidekiq|perform_async|perform_in\b|perform_at\b|perform_bulk|push_bulk . rg -n perform_later|deliver_later|queue_adapter|retry_on|discard_on|queue_as . rg -n unique_for|unique_until|sidekiq_options|sidekiq_retry_in|sidekiq_retries_exhausted . rg -n sidekiq-cron|sidekiq-scheduler|sidekiq-unique-jobs|REDIS_URL .同时盘点已调度scheduled、重试中retrying、死信dead与执行中in-flight的作业动态选择的队列/类cron 计划与时区批处理回调batch callbacks速率限制租户/locale/链路追踪上下文以及存储 Sidekiq JID 的代码。主仓库之外的生产者也要纳入盘点。删除 Redis 基础设施之前先确认 Redis 是否还被其他应用功能使用。Sidekiq 的功能列表区分 OSS、Pro、Enterprise 三个版本。请记录应用实际使用的功能包括第三方扩展而不是按购买的版本整体映射。2. 安装驱动并准备 Schema以使用 Postgres 的 Active Record 应用为例# Gemfile: 排空存量作业期间继续保留 Sidekiq gem riverqueue-activerecord gem pg使用 Sequel 则装riverqueue-sequel使用 SQLite 则把pg换成sqlite3。每个驱动都依赖riverqueue但数据库 gem 由应用自行选择。需要原子入队时务必使用与业务记录同一个数据库。无需安装 Go直接用river命令应用随 gem 打包的规范 River 迁移bundle exec river migrate-up --database-url postgres://localhost/my_app迁移命令会自动检测 bundle 中已安装的 River 驱动 gem两者都装时优先 Sequel这个选择不必与应用自身驱动一致因为两者执行的是同一套迁移。开发库与隔离的测试库都要在跑示例前提前准备。不要把spec/support下的 schema 测试夹具复制到生产环境也不要通过 Rails 模型重建 River 表。Pro 额外需要其独立的规范迁移线--line pro与另行分发的 gem。关于迁移的更多细节migrate-status、--dry-run、--target、--steps、migrate-down、Ruby 迁移 API、Go 兼容性、Pro 迁移、SQLite 迁移 008 的AUTOINCREMENT重建等参见 Schema migrations 文档。3. 转换一个作业及其生产者Sidekiq 用位置参数调用perform参数是 JSON。River 通过 Worker 注册表解析显式 kind并把一个River::Job交给work。Worker 的具体约定见 worker.rb 源码Workers注册表把 kind 映射到 worker 对象add时若不显式传 kind则从worker.kind或worker.class.kind读取。迁移前class FulfillOrderJob include Sidekiq::Job sidekiq_options queue: orders, retry: 5 def perform(order_id) FulfillOrder.call(order_id) end end jid FulfillOrderJob.perform_async(42)迁移后把每个类放到与类名一致、且被应用加载的文件中class FulfillOrderArgs def initialize(order_id:) order_id order_id end def insert_opts River::InsertOpts.new(max_attempts: 6, queue: :orders) def kind fulfill_order def to_json JSON.generate(order_id: order_id) end class FulfillOrderWorker def self.kind fulfill_order def work(job) FulfillOrder.call(job.args.fetch(order_id)) end end保持FulfillOrder.call作为应用幂等的业务操作。希望每次尝试都拿到新实例时注册 worker类而不是实例注册的实例与插件实例会在线程间共享。在 Web 进程配置好数据库连接之后创建插入客户端。这个 Rails initializer 定义了一个应用自有的访问器# config/initializers/river.rb require riverqueue-activerecord module AppJobs def self.client client || River::Client.new(River::Driver::ActiveRecord.new) end end这个访问器面向普通 Web 进程启动场景fork 之后要重新初始化客户端不要跨进程或跨 Ractor 共享客户端/连接池。Worker 进程使用第 5 节中单独配置的客户端。result AppJobs.client.insert(FulfillOrderArgs.new(order_id: 42)) job_id result.job.id返回值是River::JobInsertResultjob.id是数据库整数Postgres 序列生成、大体递增、可能有空洞不是Sidekiq 的 JID 字符串。相应地更新 API 契约、存储的引用、取消接口与日志。如果不想定义参数类可以用更简单的生产方式result AppJobs.client.insert( River::JobArgsHash.new(:fulfill_order, {order_id: 42}), insert_opts: River::InsertOpts.new(max_attempts: 6, queue: :orders) )注意JobArgsHash定义见 ruby/lib/job.rb不会继承FulfillOrderArgs#insert_opts用它时务必显式传入所需选项Worker 类本身也不提供插入默认值。Client#insert的解析逻辑见 client.rb确认max_attempts同时参考调用方 insert_opts 与参数类insert_opts默认值为 25MAX_ATTEMPTS_DEFAULT上限为 32767MAX_ATTEMPTS_LIMITqueue默认defaultpriority默认 1。读取解码后的 JSON 时使用字符串键。把每一个旧的定位参数显式翻译为具名字段不要把完整的 Sidekiq 信封塞进 River 的args。保持 ID 与简单 JSON 值而不要存模型实例、GlobalID 包装或任意 Ruby 对象。保持幂等两种系统在失败后都可能不止执行一次。4. 翻译入队方式与事务边界Sidekiq 操作Ruby River 操作perform_async(...)client.insert(args)perform_in(seconds, ...)InsertOpts.new(scheduled_at: Time.now.utc seconds)perform_at(time, ...)InsertOpts.new(scheduled_at: time.getutc)数字 epoch 用Time.at转换.set(queue: ...).perform_async(...)client.insert(args, insert_opts: River::InsertOpts.new(queue: ...))perform_bulk/push_bulkclient.insert_many参数数组或River::InsertManyParams在perform内入队子作业在work内用job.client.insert(...)client AppJobs.client client.insert( FulfillOrderArgs.new(order_id: 42), insert_opts: River::InsertOpts.new(scheduled_at: Time.now.utc 300) ) results client.insert_many([42, 43].map do |order_id| River::InsertManyParams.new( FulfillOrderArgs.new(order_id: order_id), insert_opts: River::InsertOpts.new(queue: :orders) ) end) job_ids results.map { |result| result.job.id }调度作业需要正在运行的维护 leader 以及该队列的消费者。调度设置的是最早执行时间而非精确截止时间作业可能因队列繁忙而略晚执行scheduled_at保证不早于该时刻执行见 insert_opts.rb 注释。大批量导入要分块单次调用是原子的但多次调用之间除非显式包裹否则不是同一事务。把应用写入与 River 插入放进同一个事务ActiveRecord::Base.transaction do order Order.create!(status: pending) AppJobs.client.insert FulfillOrderArgs.new(order_id: order.id) endSequel 的等价写法是DB.transaction并且 River 驱动必须由同一个DB构造——仅共享 URL 是不够的插入必须使用同一事务连接。多数据库与分片应用要显式审计。当意图是单次原子写入时不要把入队移进after_commit。远程服务调用无法加入该事务此外这个 Ruby 客户端不会把业务写入与作业完成做原子提交。5. 启动与停止 Worker 进程Sidekiq 的队列权重与严格队列顺序并不能翻译成 River 的队列 worker 数量。River 为每个队列分别分配并发度优先级1到4决定队列内作业的取用顺序1 最高、4 最低高优先级始终先于低优先级被取出。三个队列各配 10 个 worker意味着每个客户端最多 30 个并发作业多进程会将容量成倍放大。使用river worker运行专用进程它负责应用启动、信号处理与停机截止时间。配合riverqueue-rails时在config.river.configure中配置队列与原生 worker见第 8 节然后执行RAILS_ENVproduction bundle exec river worker --rails --stop-timeout 30也可以只用核心 gem 配置单个队列。将上面的 worker 与 args 类放入会被 eager-load 的应用路径。直接使用核心 gem 时插件会把应用工作包裹在 Rails executor 中可选的riverqueue-rails包会自动提供这一执行包裹和 worker 入口见第 8 节。此配置要求 eager-load 的代码开发环境的自动重载需要单独的集成。# config/river.rb # frozen_string_literal: true require_relative ../config/environment require riverqueue-activerecord Rails.application.eager_load! class RailsExecutionPlugin def work(_job, operation) Rails.application.executor.wrap { operation.call } end end River::Client.new( River::Driver::ActiveRecord.new, config: River::Config.new( job_timeout: 300, logger: Rails.logger, plugins: [RailsExecutionPlugin.new], queues: {orders: 10}, workers: River::Workers.new.add(FulfillOrderWorker) ) )在进程管理器中运行它RAILS_ENVproduction bundle exec river worker --config config/river.rb --stop-timeout 30--config文件按 cli.rb 的实现会被当作 Ruby 求值其最后一个表达式必须返回一个未启动的 client命令负责start它并保持主线程存活。不要在该文件中调用.start也不要自行安装信号处理器。信号语义如下TERM/INT停止取件并排空正在执行的尝试TSTP请求stop(wait: false)但不退出进程之后再发TERM完成停止。优雅截止时间到达后会触发中断随后有 5 秒收尾窗口超时才强制退出——因此外部 supervisor 的终止窗口要留得更长。更多细节退出状态、恢复、滚动替换、SIGTSTP行为对照表见 Dedicated worker processes 文档。注意停止期间被成功终断的尝试会重新变为available且不消耗尝试次数若进程在最终化前死亡运行中的作业会交给正常的 stuck-job 恢复逻辑处理。Worker 必须容忍重试——停机无法回滚外部副作用。要为 worker、队列生产者、维护服务以及共享同一连接池的 Web 流量统一预算数据库连接数。对目标进程数与并发度做压测不要把连接池大小直接设成作业并发数而不留开销。监控数据库争用尤其使用 SQLite 时。在 prefork 步骤之后创建客户端如果在每个 Web initializer 里都启动客户端会在每个 Web 进程里启动消费者——除非有意为之否则那里应使用仅插入insert-only的客户端。6. 翻译重试、取消与超时Sidekiq 的整数retry: n统计的是首次执行之后的重试次数River 的max_attempts统计的是全部尝试次数。对于全新作业使用n 1Sidekiq 默认 25 次重试对应 26 次尝试而 River 默认是 25 次尝试MAX_ATTEMPTS_DEFAULT见 client.rb。两者的重试调度表也不同。现有行为迁移决策retry: 5插入时max_attempts: 6retry: 0或retry: falsemax_attempts: 1阻止重试River 会把失败的作业保留为discarded状态删除/死信行为需要单独决策sidekiq_retry_in实现next_retry(job, error)返回绝对 Time而不是延迟秒数sidekiq_retries_exhausted、death handlers在应用代码中实现持久化的终态失败处理没有对等的回调注册 API永久取消在 work 中raise River.job_cancel(reason)依赖未就绪raise River.job_snooze(seconds)重新调度且不消耗尝试次数固定重试延迟的写法class FulfillOrderWorker def next_retry(_job, _error) Time.now.utc 60 end end运行时负责记录抛出的错误。不要为了静默而 rescue 错误后直接返回——正常返回即代表成功。River 的默认超时是 60 秒JOB_TIMEOUT_DEFAULT见 config.rb。显式审查现有作业耗时要么配置job_timeout要么实现timeout(job)返回秒数、nil表示禁用、0表示继承客户端默认值。网络客户端超时也要一并配置。错误上报器要避免意外触发取消error_handler lambda do |error, job| ErrorReporter.capture(error, job_id: job.id, kind: job.kind) nil # 返回 true 或 :cancel 会让 River 取消该作业 end config River::Config.new(error_handler: error_handler)核心 River 会把discarded作业保留在river_job表中请有意配置保留策略。默认值completed/cancelled保留 1 天86400 秒discarded保留 7 天604800 秒传nil或-1表示无限期保留见 config.rb 的 retention 校验。River Pro 的死信存储是独立的需要 Pro 配置。这两种保留模型都不能假设等于 Sidekiq 的 Dead set。7. 替换中间件与执行上下文Sidekiq 有独立的 client 与 server 中间件链River 通过Config.new(plugins: [...])接收插件实例。插件签名详见 Ruby 功能指南。用途River 插件方法修改每次插入的元数据Hookinsert_begin(params)观察每次插入结果Hookinsert_end(result)包裹插入Middlewareinsert_many(params, operation)作业执行前恢复上下文Hookwork_begin(job)观察返回/错误Hookwork_end(job, error)包裹执行并清理上下文Middlewarework(job, operation)class LocalePlugin def insert_begin(params) params.metadata[locale] I18n.locale.to_s end def work(job, operation) I18n.with_locale(job.metadata.fetch(locale, I18n.default_locale)) do operation.call end end end在每个生产者上注册插入插件包括 worker 用于入队子作业的客户端在消费者上注册 work 插件。靠前的插件包裹靠后的插件对 Rails 应用把RailsExecutionPlugin放在应用上下文插件之前。不要把与尝试相关的上下文放进共享实例变量即使发生异常也要恢复线程局部上下文。中间件必须调用operation.call才能继续。不要把 Sidekiq 中间件的否决veto习惯带过来——在 River 的 work 中间件里提前返回会直接标记作业完成而根本不执行 worker。请使用显式取消或应用侧校验。insert_end不是 after-commit 通知它可能在外层事务之后回滚的事务内运行。订阅subscriptions适合本地遥测但事件是有限、可能被丢弃的不是持久的跨进程回调系统。重要的后续动作请使用持久化应用记录或工作流任务。8. 显式处理 Active Job 与 Action MailerActive Record 驱动本身不是Active Job 适配器。要保留既有 Active Job 与 Action Mailer 作业安装riverqueue-rails并执行bin/rails generate river:install bin/rails river:migrate bin/jobs start安装器会把config.active_job.queue_adapter配置为:river既有的perform_later与deliver_later随即改用 River。切流前请审查生成的队列配置与 Rails 集成指南。几个要点Active Job 拥有应用层重试未处理的异常会让当前 River 投递行被 discard而不会再触发一轮后端重试。retry_on会用相同的 Active Job UUID 创建新的 River 行。依赖同连接原子入队时需显式禁用 after-commit 延迟Rails 8.x 为self.enqueue_after_transaction_commit false。集成保留 Active Job 的序列化、回调、GlobalID、locale 与执行上下文它不会导入既有的 Redis 作业。优先级必须是nil对应 River 优先级 1或 14 的整数不支持的优先级会直接抛错而不是静默改变语义。另一种做法是迁移期间让 Active Job 留在原后端或把其业务操作抽取为原生 River worker。回调、retry_on、discard_on、GlobalID 反序列化、locale 与队列命名都要显式翻译。对邮件入队一个包含收件人/模型 ID 的原生 River 作业并在该 worker 中调用deliver_now——若在那里调用deliver_later只会再入队另一个 Active Job邮件不会在 River 中完成。9. 映射周期、唯一性与商业版功能周期任务Recurring schedules小时级间隔可以这样在 worker 客户端上配置periodic River::PeriodicJob.new( id: :hourly_order_reconciliation, constructor: - { [River::JobArgsHash.new(reconcile_orders, {}), River::InsertOpts.new(queue: :orders)] }, schedule: River::PeriodicInterval.new(3600) ) config River::Config.new( periodic_jobs: [periodic], queues: {orders: 10}, workers: workers # 同时为 reconcile_orders 注册一个 worker )把该 config 传给消费客户端。间隔不是 cron 表达式需要日历调度时加gem fugit, ~ 1.13并用River::PeriodicCron.new(0 9 * * 1-5, timezone: America/New_York)作为schedule:。Fugit 是可选的、只在构造该 helper 时加载默认时区是 UTC源码见 periodic_cron.rb支持CRON_TZ/TZ前缀与 Go 的every时长语法。自定义 schedule 可以响应next(time)也可以是返回下一个Time的可调用对象。测试夏令时切换与错过运行的情况。核心版 schedule 存在内存中运行时可经client.periodic_jobs动态增删见 periodic_job.rbPro 提供持久化调度。在可获得 leader 选举资格的客户端上配置兼容的 schedule。启用替代调度器时要停用旧周期生产者否则两个调度器会入队同一发生点。唯一性UniquenessSidekiq Enterprise 的unique_for是锁 TTLRiver 的by_period用的是由调度时间推导的时间桶不是滑动的锁 TTL。unique_until: :start也没有直接映射因为 River 要求自定义唯一状态集合中必须包含runningREQUIRED_UNIQUE_STATES强制要求 available/pending/running/scheduled见 client.rb。当等价的作业还未完成时可以做如下唯一性约束unique River::UniqueOpts.new( by_args: true, by_queue: true, by_state: %w[available pending running scheduled retryable] ) result AppJobs.client.insert( FulfillOrderArgs.new(order_id: 42), unique_opts: unique ) duplicate result.unique_skipped_as_duplicate?这里刻意排除了completedRiver 的默认唯一状态集合DEFAULT_UNIQUE_STATES包含 completed直到保留期清理为止。命中重复时返回已存在的行调用方绝不能把它当作新插入。唯一性不会让外部副作用变成精确一次执行也不会在 Redis 与 SQL 之间去重。关于唯一键的生成算法kind、args、period、queue拼接后 SHA-256 哈希与 Go 客户端必须完全一致见 client.rb 的 make_unique_key_and_bitmask 实现。Sidekiq Pro 批处理BatchesSidekiq batch 协调一组作业与回调River Pro 的 workflow 是更接近的抽象而River::Pro::BatchWorker则是在一次work_many调用中处理多个作业。在安装了私有分发的riverqueue-pro并应用其迁移后一次扇出 成功任务的写法如下require riverqueue-pro # 启动该客户端前先注册 import_row 与 finish_import worker pro_client River::Pro::Client.new( River::Driver::ActiveRecord.new, config: River::Pro::Config.new( core: River::Config.new(queues: {imports: 10}, workers: workers) ) ) workflow River::Pro.workflow name: import do |flow| tasks [101, 102].map do |row_id| flow.add( row_#{row_id}, River::JobArgsHash.new(:import_row, {row_id: row_id}), queue: :imports ) end flow.add( :finish, River::JobArgsHash.new(:finish_import, {import_id: 7}), after: tasks, queue: :imports ) end pro_client.insert_many workflow.jobs必须有一个正在运行的 Pro 消费客户端。默认情况下被取消/丢弃的依赖会阻止成功任务运行。Sidekiq 的complete回调意味着所有作业都执行过一次这与依赖最终化不同其death回调也需要显式重新设计。在替换生产 batch 前请测试重试、取消、动态任务追加、空输入与终态失败行为。其他功能决策Sidekiq 功能或扩展River 迁移路径Enterprise 速率限制应用侧限流器Pro 并发限制约束的是同时作业数而不是每秒请求数Enterprise 周期调度核心版周期任务或 Pro 持久化周期任务审查日历语义Enterprise 加密ProEncryptPlugin配合应用加密器旧负载需要显式解码再重新编码过期作业应用侧截止检查配合显式取消保留与执行超时不会让已入队作业过期长时运行的迭代/检查点核心版 resumable steps/cursors检查点格式需显式翻译顺序执行Pro sequences需显式分组与失败策略Web UI 与指标单独部署 River UI配合插件/订阅集成本 gem 没有Sidekiq::Web的 Rack 挂载替代品多进程监督、滚动重启应用部署 supervisor 配合 River 停止策略Pro 的确切 API 与分发要求见私有riverqueue-ruby-pro仓库的 README它是独立包公共检出中可能不存在。不要假设 Sidekiq 的商业功能在 River Pro 中有完全一致的行为。10. 针对持久化行为重写测试River 没有 Sidekiq fake/inline 测试工具的对应物包括新版Sidekiq.testing!API。正确做法是先对业务操作做单元测试然后对入队、回滚与执行在隔离的、已迁移且与生产适配器一致的数据库上做测试。核心 gem 自带数据库后端的测试辅助工具见 testing.md 文档并有可选的 RSpec 与 Minitest 集成require riverqueue/testing/rspec或riverqueue/testing/minitest。这些断言针对新持久化的行唯一性返回的既有行不算插入并可以在调用线程上执行一次真实尝试require riverqueue/testing/rspec RSpec.configure { |config| config.include River::Testing::RSpec } expect { enqueue_order(42) }.to insert_job( AppJobs.client, args: {order_id 42}, kind: fulfill_order ) row AppJobs.client.insert(FulfillOrderArgs.new(order_id: 42)).job result River::Testing.perform_job(AppJobs.client, row.id) expect(result).to have_attributes(error: nil, outcome: :completed)使用已停止的 client、已注册 worker 与隔离数据库。与线程执行不同当驱动与测试事务使用同一连接时同步 helper 能看到调用方事务内的作业。它们绕过队列容量/暂停也不运行维护或周期生产者——这些运行时行为要保留线程化冒烟测试。一个 RSpec 插入测试使用第 3 节的应用访问器it enqueues the expected payload and options do result AppJobs.client.insert(FulfillOrderArgs.new(order_id: 42)) row AppJobs.client.job_get(result.job.id) expect(row).to have_attributes( args: {order_id 42}, kind: fulfill_order, max_attempts: 6, queue: orders ) end一个不依赖业务依赖、真实运行时的执行冒烟测试class MigrationProbeWorker def self.kind migration_probe def work(job) job.output {seen: job.args.fetch(value)} end it works a committed job do client River::Client.new( River::Driver::ActiveRecord.new, config: River::Config.new( queues: {migration_test: 1}, workers: River::Workers.new.add(MigrationProbeWorker) ) ) result client.insert( River::JobArgsHash.new(:migration_probe, {value: 42}), insert_opts: River::InsertOpts.new(queue: :migration_test) ) begin client.start deadline Process.clock_gettime(Process::CLOCK_MONOTONIC) 5 loop do row client.job_get(result.job.id) break if row.state River::JOB_STATE_COMPLETED raise job did not complete: #{row.state} if Process.clock_gettime(Process::CLOCK_MONOTONIC) deadline sleep(0.01) end expect(client.job_get(result.job.id).metadata.fetch(output)).to eq(seen 42) ensure client.stop_and_cancel end end线程化执行测试要禁用事务性 fixtures其他连接看不到未提交的插入。停止客户端之后只清理隔离测试库。SQLite 请使用共享给连接池各连接的一个临时文件而不要给每个连接单独的 in-memory 数据库。此外还应断言事务回滚会移除作业失败消耗预期的尝试预算snooze 不消耗调度作业不提前执行队列隔离与唯一性生效出错时上下文复位停止允许恢复。用两次调用同一逻辑操作来测试业务幂等性。直接worker.work(job)的单元测试无法验证运行时的重试与最终化。11. 存量作业切流与回滚保障推荐的做法是排空drainSidekiq同时把新的逻辑作业按 kind 逐一路由到 River。过渡期两个消费者可以并存但每一次入队决策必须只选一个后端。旧的 Redis 负载要保留旧类以兼容。若顺序性很重要先结束旧流再启用新流。部署 schema、River 消费者、转换后的 worker以及一个初始仍选择 Sidekiq 的生产者路由开关。先用少量 canary 负载验证 River再切换选定生产者。worker 创建的子作业与周期生产者也要纳入该开关。排空旧的 ready 与 in-flight 作业。单独核算 scheduled、retries、dead 与商业版 batch 状态——ready 队列为空并不能证明 Redis 中已无相关作业。在分配的工作完成或被显式转移之前保留 Sidekiq 消费者与必要调度。按后端与作业 kind 跟踪数量与失败。只有当清单核对无误后才移除 Sidekiq 依赖、路由、initializer、部署进程与仅用于 Redis 的作业基础设施。若必须转移长寿命的 scheduled/retry 作业请用 Sidekiq 公共 API 和River::Client#insert构建一个独立的、可重启的导入器。不存在跨越 Redis 与 SQL 的原子事务。采用如下传输协议停止会修改所选源作业的生产者、调度器与消费者并核对 in-flight 作业。删除任何内容前先导出稳定清单——对活跃队列做 API 枚举可能与变更竞争。把每个白名单内的 Sidekiq 类映射为显式的 River kind、参数转换、队列与尝试策略。在专门的转换方案出现前拒绝未知类以及包装过的 Active Job、加密、batch 负载。在同一个 SQL 事务中插入 River 作业与应用自有的传输回执receipt回执以源身份上的唯一约束保护包含 Redis namespace/cluster 加 JID并把 JID 保留在 metadata 中以供关联。回执的保留独立于 River 作业的保留期。先提交 SQL再确认/移除对应的那个 Redis 作业。重启时查回执而不是重复入队。显式处理唯一性结果——返回的既有 River 作业并不能自动证明目标源作业已转移。用源清单与目标作业双向核对回执包括提交后移除的源作业。不要盲目重放导出。保留未来的执行时间对重试作业的剩余尝试次数要单独决策不要把全新作业的n 1规则套到已部分执行的作业上。死信作业要么归档要么移入经过审查的人工重试流程而不是让每条死信立即可运行。不要把 Sidekiq 的错误历史直接写进 River 的 schema。执行中的商业版 batch 通常需要在 Sidekiq 中跑完或重建为经审查的 workflow。回滚方案把新生产者切回 Sidekiq同时保留一个 River 消费者处理已提交到那里的作业或显式暂停这些作业并规划恢复。把所有 River 作业重新入队到 Redis 可能造成已完成副作用的重复。跨两个后端保持共享业务幂等键稳定。完成标准当满足以下条件时迁移才算完成清单每一行都有已验证的目标与行为生产者路径选择了预期后端已部署的消费者覆盖所有目标队列与 kind生产适配器上的集成测试通过旧的 queued/scheduled/retrying 作业已核对。在移除旧系统前还要在代表性负载下确认重试预算、超时、保留策略、连接池容量、周期调度与停机行为。更多 Ruby API 细节请参考Ruby 主指南、client.rb 源码、config.rb 源码、insert_opts.rb 源码 与 client_runtime.rb 源码。赞分享任务调度后端【免费下载链接】riverThe polyglot queue: Fast and reliable background jobs in Go, Ruby, Rust, and JS/TS on Postgres or SQLite.项目地址https://gitcode.com/gh_mirrors/river/river点击查看免费下载相关推荐River JS/TS 完整测试指南用 riverqueue/test 与内存 SQLite 覆盖生产者、Worker 与完整运行时River JS/TS 完整测试指南用 riverqueue/test 与内存 SQLite 覆盖生产者、Worker 与完整运行时 River 是运行在任务调度后端k0s生产环境部署终极指南从开发测试到企业级应用的完整迁移方案k0s生产环境部署终极指南从开发测试到企业级应用的完整迁移方案 在当今云原生技术快速发展的时代k0s作为零摩擦Kubernetes发行版为生产环境部署提供云原生容器编排边缘计算3分钟快速上手iptv-checker智能检测你的IPTV播放源3分钟快速上手iptv checker智能检测你的IPTV播放源 还在为IPTV播放源频繁失效而烦恼吗面对几百个频道列表你是否还在手动逐个测试今天我要为后端任务调度音视频上一篇Langfuse 多语言支持完全指南3步切换到中文文档下一篇DoctrineFixturesBundle扩展开发3步创建自定义清除器工厂的完整指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
延伸阅读

更多相关文章

2026/10/12 1:44:29

Linux C进程管理:fork/exec/wait与僵尸进程实战解析

最近在整理这几年写Linux C的代码笔记,第一个想聊透的就是进程管理。很多读者留言说fork会用,但每次跑多进程程序都出奇奇怪怪的问题——父进程退出了子进程还在跑、ps里冒出一堆Z状态进程、fork之后printf的输出重复了。这些都是进程管理没形成体系的表…

2026/10/12 1:44:29

YOLOv8结合SAM实现开集实例分割的工程实践

简介:一套面向计算机视觉研究与工程实践的资源,将Meta推出的SAM分割模型与YOLOv8检测框架相结合,专为需要实现开集实例分割与目标检测的场景而设计,适合算法工程师、科研人员和有一定基础的视觉学习者。压缩包共6个文件&#xff0…

2026/10/12 1:44:29

药品包装盒数据集VOC+YOLO双格式:从标注解读到YOLO训练避坑指南

简介:这份药品包装盒目标检测数据集,面向计算机视觉入门者、算法研究及工业质检项目开发者,提供统一格式的真实场景图像,解决训练检测模型时缺乏高质量标注数据的痛点。压缩包共2000个文件,包含1032张药品包装盒图片、…

2026/10/12 2:59:32

STM32驱动DS1302实时时钟:GPIO模拟时序与寄存器配置详解

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/12 2:59:32

虚拟电厂云端功率预测:坐标代替气象站,降本30%-50%

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/12 2:59:32

STM32C5与CubeMX2实战:从选型到避坑的嵌入式开发指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/12 2:59:32

超市收银系统设计说明书:数据模型、事务边界与离线对账全解析

简介:超市收银系统设计说明书是一份面向计算机相关专业毕业设计或课程设计的参考范文,系统讲解超市收银系统的完整设计流程。内容从需求分析入手,覆盖数据流图、数据字典和实体联系图,再到系统概要设计、数据库概念与逻辑结构设计…

2026/10/11 0:02:13

Python调用Gemini Structured Outputs实现工单路由门禁

客服工单最怕的不是模型“答错一句话”,而是它给出一段看起来合理的说明,程序却从中猜错优先级。通俗做法是:要求模型只交 JSON(JavaScript Object Notation,轻量数据格式),再让代码验证它。Gem…

2026/10/11 0:02:13

Spring Boot超市进销存系统毕设实战:从需求拆解到答辩通关

最近带的一个学生项目组里,有A同学跑来问我:选什么毕设题目最稳妥,既能让评审老师觉得工作量够,又不会在答辩时被问到语无伦次。我第一反应就是推荐基于Spring Boot的超市仓库管理系统——也就是超市进销存系统。这个题目乍一看平…

2026/10/11 0:02:13

Flutter StatefulWidget 生命周期核心解析

很多刚开始接触 Flutter 的朋友,在看完一堆“Hello World”和基础组件之后,大概率都会撞上同一堵墙:StatefulWidget 里那堆 initState、build、dispose 方法,到底什么时候被调用?为什么顺序是那样?在里面到…

2026/10/12 0:04:22

绝缘子缺陷检测数据集清洗与工业级训练实战指南

简介:本资源是面向电力AI研发人员、工业视觉工程师及智能巡检系统开发者的绝缘子缺陷检测专用YOLO格式数据集,解决无人机航拍场景下绝缘子破损、污闪、积雪等9类典型缺陷的精准识别与定位难题。数据集共2139张真实巡检图像(含训练/验证/测试集…

还想了解更多?直接咨询顾问

免费诊断 + 免费方案 + 透明报价。

全国咨询热线400-8866-253
免费获取方案
☎咨询二维码 ☎ ↑