首页
/ Rails Active Job 实战详解:作业声明、延迟调度、队列适配器与 Continuations 断点续跑

Rails Active Job 实战详解:作业声明、延迟调度、队列适配器与 Continuations 断点续跑

2026-09-05 22:40:58作者:殷蕙予

Active Job 是 Rails 的异步作业框架,负责把"稍后要执行的活儿"从 HTTP 请求-响应周期中剥离出来,交给各种队列后端(Queuing Backend)执行。本文基于 Rails 仓库中的 activejob/README.md 及其配套源码展开,完整覆盖作业声明与入队、set 调度选项、GlobalID 参数序列化、Action Mailer 的 deliver_later 集成、各队列适配器的能力对比,以及较新的 Continuations(可中断、可恢复作业)机制,并给出每个结论对应的源码与测试路径,方便你对照阅读。

一、Active Job 的定位与设计目标

README 对 Active Job 的定义是:"一个用于声明作业、并让其在多种队列后端上运行的框架"。这些作业可以是定时清理、计费扣款、邮件群发——任何可以拆成小块、可并行执行的工作单元都适用。

它承担两个层面的职责:

  1. 统一作业接口。README 指出,Active Job 的核心目标是"确保每个 Rails 应用都有一套作业基础设施,哪怕只是一个 immediate runner(立即执行器)"。有了这层抽象,框架特性和第三方 gem 就无需关心 Delayed Job 与 Resque 之间的 API 差异,可以直接构建在 perform_later 之上。
  2. 后端可替换。选择哪个队列后端变成纯粹的运维决策(operational concern),切换后端时"不需要重写你的作业"(without having to rewrite your jobs)。

从源码结构看,这一抽象的核心是 ActiveJob::Base:它通过组合多个模块(CoreQueueAdapterQueueNameQueuePriorityEnqueuingExecutionCallbacksExceptionsInstrumentationLogging 等)拼装出完整的作业行为,作业类只需继承它并实现 perform 方法。README 给出的标准用法即:

class MyJob < ActiveJob::Base
  queue_as :my_jobs

  def perform(record)
    record.do_work
  end
end

其中 queue_as :my_jobs 指定该作业进入名为 my_jobs 的队列,其实现见 QueueName 模块:除静态名称外还支持传入 block 做动态队列名(例如根据 arguments.first 的付费状态路由到 :paid_feeds:feeds),并可通过 queue_name_prefix / queue_name_delimiter(默认 "_")为所有队列加统一前缀,便于在多应用共享同一队列中间件时做隔离。

二、入队与调度:perform_later 及 set 选项

README 给出了三种最典型的入队方式,这里完整保留并逐一说明:

MyJob.perform_later record  # 入队,队列系统空闲时尽快执行
MyJob.set(wait_until: Date.tomorrow.noon).perform_later(record)  # 明天中午 12 点执行
MyJob.set(wait: 1.week).perform_later(record)  # 从现在起 1 周后执行

2.1 set 支持的全部选项

结合 ActiveJob::Core::ClassMethods#set 的 RDoc,set 接受四个选项:

选项 含义
:wait 延迟指定时长后执行,如 VideoJob.set(wait: 5.minutes).perform_later(Video.last)
:wait_until 在指定时刻执行,如 VideoJob.set(wait_until: Time.now.tomorrow)
:queue 覆盖类级队列名,如 VideoJob.set(queue: :some_queue)
:priority 指定优先级(数值越小优先级越高)

多个选项可以组合,例如 VideoJob.set(queue: :some_queue, wait: 5.minutes, priority: 10).perform_later(Video.last)set 返回一个 ConfiguredJob 预配置对象;对实例而言,Core#set 的实现很直白::wait 换算成 scheduled_at = options[:wait].seconds.from_now:wait_until 直接赋值 scheduled_at:queuequeue_name_from_part:priority 转成整数。

2.2 入队链路:从 perform_later 到适配器

从源码结构看,perform_later 到真正入队的调用链如下(见 Enqueuing 模块):

  1. 类方法 perform_later(...) 通过 job_or_instantiate 拿到作业实例(若传入的已是本类实例则直接复用),随后调用 job.enqueue,支持 yield 作业给可选 block,返回入队结果;
  2. 实例方法 enqueue(options = {}) 先执行 set(options) 应用调度选项,再调用 raw_enqueue
  3. raw_enqueue 包裹 :enqueue 回调链执行 _raw_enqueue:若设置了 scheduled_at 则走 queue_adapter.enqueue_at(self, scheduled_at.to_f)(延迟执行),否则走 queue_adapter.enqueue(self)(立即执行);
  4. 若适配器抛出 EnqueueError,不会冒泡,而是记录到 job.enqueue_error 并把返回值置为 false,调用方可据此判断入队是否成功。

一个容易被忽略但很实用的细节:Enqueuing 模块定义了类属性 enqueue_after_transaction_commit(默认 false)。当 Active Job 与 Active Record 联合使用时,在数据库事务内调用 perform_later 会隐式把入队推迟到事务提交之后(回滚则丢弃作业),避免"作业先于数据可见"的竞态。可在全局或单个作业类上设置 self.enqueue_after_transaction_commit = true/false,相关行为有专门的测试 enqueue_after_transaction_commit_test.rb 覆盖。

作业执行完毕后,ActiveJob::Core#serialize 决定了交给队列后端的完整数据结构:job_classjob_id(默认 SecureRandom.uuid)、queue_namepriorityarguments(序列化后的参数)、executionslocaletimezoneenqueued_atscheduled_at。这意味着作业的执行环境(时区、语言)会随作业数据一起跨进程传递,执行端通过 deserialize 恢复。

三、GlobalID 参数支持:直接传 Active Record 对象

README 专门用一小节强调:Active Job 支持对参数做 [GlobalID 序列化],这使得"直接把活的 Active Record 对象传给作业"成为可能,而不必像过去那样传 类名 + id 再手动 constantize.find。README 给出的对比示例值得完整保留:

改造前:

class TrashableCleanupJob
  def perform(trashable_class, trashable_id, depth)
    trashable = trashable_class.constantize.find(trashable_id)
    trashable.cleanup(depth)
  end
end

改造后:

class TrashableCleanupJob
  def perform(trashable, depth)
    trashable.cleanup(depth)
  end
end

其适用范围是"任何混入了 GlobalID::Identification 的类,默认包括所有 Active Record 模型"。默认情况下,ActiveJob::Arguments 接受的原生类型包括 StringIntegerFloatNilClassTrueClassFalseClassBigDecimalSymbolDateTimeDateTimeActiveSupport::TimeWithZoneActiveSupport::DurationHashActiveSupport::HashWithIndifferentAccessArrayRange,以及 GlobalID 实例;这一白名单在 Enqueuing 的 perform_later 文档 中有明确说明,且可以通过注册自定义序列化器扩展(序列化逻辑见 Serializers 模块,参数级测试见 argument_serialization_test.rb)。

四、Action Mailer 的 deliver_later:把邮件变成作业

README 指出 Active Job 还充当 Action Mailer #deliver_later 的后端:"这让任何邮件都能轻松变成一个稍后执行的作业"。README 认为这是现代 Web 应用中最常见的作业类型之一——把发信移到请求-响应周期之外,用户就不必为它等待。

在仓库中可以印证这条链路:ActionMailer::Base 的文档明确写着 NotifierMailer.welcome(User.first).deliver_later # enqueue the email sending to Active Job,并提供了 delivery_job(默认作业类)与 deliver_later_queue_name 两个类属性用于定制。默认的 ActionMailer::MailDeliveryJob 接收邮件序列化数据,在后台重新构建并投递,从而完全复用 Active Job 的队列、重试与调度能力。

五、支持的队列后端与适配器能力对比

README 说明 Active Job 内置了多个队列后端的适配器(Resque、Delayed Job 等),并给出了一条重要的治理声明:Rails 不再接收新适配器的 pull request,正在把现有适配器向外抽取(actively extracting the current adapters),鼓励库作者在自家 gem 中或作为独立 gem 提供 Active Job 适配器。

从当前仓库 QueueAdapters 来看,随附的适配器包括外部队列后端:Backburner、Delayed Job、queue_classic、Resque、Sneakers,以及三个用于测试和开发的内置适配器:AsyncAdapter(线程池异步执行)、InlineAdapter(进程内立即执行,即 README 所说的 "immediate runner")、TestAdapter(仅记录不执行,供测试断言)。每个适配器位于 activejob/lib/active_job/queue_adapters/ 目录下,如 async_adapter.rbinline_adapter.rbtest_adapter.rbresque_adapter.rb 等。

选择后端时,README 指向 ActiveJob::QueueAdapters 的文档;该文件中的能力对比表(Backends Features)值得直接引用,列出了各后端在异步执行(Async)、多队列(Queues)、延迟执行(Delayed)、优先级(Priorities)、超时(Timeout)、重试(Retries)六个维度的支持情况:

后端 Async Queues Delayed Priorities Timeout Retries
Backburner Yes Yes Yes Yes Job Global
Delayed Job Yes Yes Yes Job Global Global
queue_classic Yes Yes Yes* No No No
Resque Yes Yes Yes (Gem) Queue Global Yes
Sneakers Yes Yes No Queue Queue No
Active Job Async Yes Yes Yes No No No
Active Job Inline No Yes N/A N/A N/A N/A
Active Job Test No Yes N/A N/A N/A N/A

表中各维度的含义(摘自同文件注释):

  • Async:作业能否以非阻塞方式运行(独立/分叉进程或不同线程);"No" 表示作业在请求进程内同步运行;
  • Queues:能否用 queue_asset 指定作业所在队列;
  • Delayed:能否通过 perform_later 在未来执行;标注 "(Gem)" 表示需额外 gem 支持,"No" 表示只能"有机会就跑";
  • Priorities:作业处理顺序的控制粒度——"Job"(作业级)、"Queue"(队列级)、"Global"(全局配置)、"No";
  • Timeout:作业运行时限的配置粒度(作业级 / 队列级 / 全局 / 不支持);
  • Retries:重试次数配置能力(作业级 / 全局 / 不支持)。

补充说明:queue_classic 自 3.1 版起支持作业调度(Delayed 列中的 "Yes*"),旧版本可借助 queue_classic-later gem。适配器的具体实现测试可参见 adapter_test.rbasync_adapter_test.rb 等用例。

六、Continuations:可中断、可恢复的长作业

README 在 "Continuations" 一节指出:Continuations 允许作业被中断并恢复(interrupted and resumed),详见 ActiveJob::Continuation。这是当前仓库中一个内容相当丰富的新特性,continuation.rb 的 RDoc 给出了完整的设计说明,其价值在于让长作业能在应用重启(部署)时保留进度。核心机制如下:

6.1 基本用法:step 定义步骤

作业类 include ActiveJob::Continuable 后即启用该能力,被中断的作业会自动重试恢复。用 step 方法定义步骤,步骤可带可选 cursor(游标)跟踪进度。官方示例(摘自源码 RDoc):

class ProcessImportJob < ApplicationJob
  include ActiveJob::Continuable

  def perform(import_id)
    # 每次执行(含恢复)都会运行
    @import = Import.find(import_id)

    step :validate do
      @import.validate!
    end

    step(:process_records) do |step|
      @import.records.find_each(start: step.cursor) do |record|
        record.process
        step.advance! from: record.id
      end
    end

    step :reprocess_records
    step :finalize
  end

  def reprocess_records(step)
    @import.records.find_each(start: step.cursor) do |record|
      record.reprocess
      step.advance! from: record.id
    end
  end

  def finalize
    @import.finalize!
  end
end

执行语义是:步骤按遇到的顺序立即执行;作业被中断后,已完成的步骤会被跳过(skip),进行中的步骤会从最后记录的 cursor 处恢复;不属于任何 step 的代码在每次运行时都会执行,因此恢复场景下要注意幂等。step 既可传 block(以 step 对象为参数),也可传方法名(方法可不带参数或接收 step 对象)。

6.2 Cursor(游标)

Cursor 用于在步骤内跟踪进度,可以是任何能经 ActiveJob::Base.serialize 序列化的对象,默认 nil;恢复时自动还原最后值,由步骤代码负责"从正确位置继续":

  • step.set! 把 cursor 设为指定值(示例:items[step.cursor..].each { ...; step.set! (step.cursor || 0) + 1 });
  • 定义步骤时可用 start: 指定初始游标,如 step :iterate_items, start: 0
  • step.advance! 调用当前游标的 succ 前进;游标不支持 succ 时抛 ActiveJob::Continuation::UnadvanceableCursorError
  • advance! from: record.id 适用于 ID 不连续的场景(如 find_each 按主键迭代);
  • 游标可以是数组以遍历嵌套集合,例如 start: [0, 0] 同时记录 (account_id, record_id),在两层 find_each 中分别 step.set!

6.3 Checkpoint(检查点)与中断时机

设置/前进游标即自动创建检查点;也可在无需更新游标时手动调用 step.checkpoint!(如循环中逐条 destroy! 并检查点)。检查点的语义是"作业可被中断的位置":届时作业会调用 queue_adapter.stopping?,若返回 true 则以 :stopping 为原因抛 ActiveJob::Continuation::Interrupt(继承自 Exception 而非 StandardError,因此不会被常规异常处理捕获);返回其他真值则作为中断原因。每次作业执行中,除第一个步骤外,每个步骤开始前都有自动检查点;作业不会在适配器标记 stopping 的瞬间被打断,而是继续运行到下一个检查点或进程停止——这意味着作业应比关闭超时(shutdown timeout)更频繁地打检查点,以保证优雅重启。

被中断后,作业自动重试,进度序列化在作业数据的 continuation 键下,内容为:已完成步骤列表 + 当前步骤及其游标(若步骤进行中)。

6.4 Isolated Steps、Attributes 与错误处理

  • isolated: true:让某一步骤总是在独立的一次执行中运行(step :slow_step, isolated: true),适合"在作业宽限期内无法打检查点"的超长步骤——它确保步骤开始前进度已序列化回作业数据;
  • Attributes:步骤只保存序列化进度、不保存其他状态。跨步骤复用中间结果时,用 ActiveJob::Attributes 声明作业属性(该模块已包含在 Continuable 中),中断时自动序列化、恢复时自动还原;
  • 错误自动重试:若作业在报错后本应经 Active Job 重试而未重试,进度会交还给底层队列后端而丢失。为缓解这一点,只要作业"已经取得进度"(完成过步骤或推进过游标),报错时就会自动重试。

6.5 配置项

ActiveJob::Continuable 提供三个类属性配置(均可在作业类中覆盖):

配置 默认值 说明
max_resumptions nil(无限次) 作业最多可恢复的次数;从源码看,超限会抛 Continuation::ResumeLimitError
resume_options { wait: 5.seconds } 恢复时传给 retry_job 的选项,如 { wait: 1.seconds, queue: :resumed }
resume_errors_after_advancing true 推进游标之后发生错误时是否仍恢复

示例:

class ProcessImportJob < ApplicationJob
  self.max_resumptions = 3
  self.resume_options = { wait: 1.seconds, queue: :resumed }
  self.resume_errors_after_advancing = false
end

行为测试可对照 continuation_test.rbattributes_test.rb

七、安装、许可证与延伸阅读

README 的安装说明:最新版本可通过 RubyGems 安装 gem install activejob;源码位于 Rails 项目的 activejob 目录下(本仓库即 activejob/)。Active Job 以 MIT 许可证发布(见 activejob/MIT-LICENSE)。

更完整的入门内容(创建与入队作业、配置后端、后台执行、异步发信、部署时暂停/恢复作业)见仓库内指南源文件 guides/source/active_job_basics.md——README 中指向的 "Active Job Basics" 指南在本仓库中对应此文件;适配器 API 的权威参考为 ActiveJob::QueueAdapters 的 RDoc。作业序列化的完整行为有 job_serialization_test.rbserializers_test.rb 覆盖,入队行为见 queuing_test.rb

小结:Active Job 通过 perform_later + set(wait:/wait_until:/queue:/priority:) 提供统一的入队接口,用 GlobalID 让 Active Record 对象可作参数,用适配器层把队列后端的选型降维为运维决策,并以 deliver_later 承载了邮件发送这一最高频的异步场景;新增的 Continuations 则把"长作业跨部署存活"这一运维痛点纳入了框架能力。理解了上述源码路径(enqueuing.rbcore.rb → 各 queue_adapters/*continuation.rb),即可在本仓库中继续深入任何一环的实现细节。

登录后查看全文
热门项目推荐
相关项目推荐