Elixir 后台任务与调度:Oban 架构、唯一性与生产运维

深入 Elixir 后台任务框架 Oban:Postgres 队列的架构与表结构、Worker 定义与返回值语义、唯一性约束与指数退避重试、Cron 定时插件、队列隔离与并发控制,以及 Telemetry 可观测、Oban Web UI 与 OTP 监督树集成。

几乎每一个 Web 应用都会在某个时刻遇到「这个操作不该在请求里做」的场景:发送邮件、生成报表、调用第三方 API、清理过期数据。把这些工作塞进 HTTP 请求的同步路径,最直接的后果是响应时间被外部依赖绑架——第三方 API 抖动一下,用户就要等上十秒;更糟的是,请求进程一旦超时被杀,工作可能执行了一半,状态永久不一致。

后台任务框架的价值就在于把「副作用」从请求路径中剥离,交给一个可持久化、可重试、可观测的独立执行器。Elixir 生态中,Oban 是这一角色的绝对主流:它用 PostgreSQL 当队列,用 SKIP LOCKED 做无锁抢占,用 OTP 进程树做执行器,把「可靠投递」这件事建立在数据库事务之上。

一、后台任务的问题域与 Oban 定位

1.1 为什么需要独立的任务系统

维度请求内同步执行后台任务系统
响应时间被最慢的外部依赖拖累恒定(只写一条任务记录)
失败处理用户看到 500,无重试自动重试 + 死信
可观测性淹没在请求日志里独立队列、独立指标
资源隔离与 Web 请求争抢连接池独立并发上限
崩溃恢复状态可能半完成任务记录仍在表中

1.2 Oban 的定位

Oban 的设计哲学非常克制:不自建队列存储,直接复用已有的 PostgreSQL。这带来三个直接好处:

  • 任务与业务数据在同一事务中写入,天然获得「业务成功则任务必存在」的原子性;
  • 不需要额外运维 Redis/RabbitMQ,部署复杂度不变;
  • 可以用 SQL 直接查询、统计、修复任务,调试成本极低。

代价是吞吐上限受 Postgres 写入能力约束。实践数据是:单表 oban_jobs 在合理索引下可以支撑每秒数千次入队,对绝大多数业务系统绰绰有余;真正需要每秒十万级任务时,才该考虑 Kafka 之类的专用流平台(见 https://plumephp.com/erlang-streaming-genstage-flow/ 中的 Broadway)。

1.3 与 GenStage/Broadway 的分工

框架数据来源触发方式状态存储适用
ObanPostgres 表入队 + 定时数据库业务副作用、定时任务
BroadwayKafka/SQS/RabbitMQ外部消息到达无(无状态)高吞吐流式处理
GenStage自定义需求驱动内存中间层背压编排

一句话区分:Oban 管「将来要做的某件事」,Broadway 管「源源不断到来的数据」。

二、Oban 架构与核心概念

2.1 进程拓扑

Oban Supervisor
 ├─ Plugin (Cron / Pruner / Lifeline)
 └─ Queue Supervisor
     ├─ Producer ── FOR UPDATE SKIP LOCKED ──▶ PostgreSQL (oban_jobs)
     │   ├─ Consumer
     │   ├─ Consumer
     │   └─ Consumer
  • Producer:每个队列一个,负责从数据库「取任务」,使用 FOR UPDATE SKIP LOCKED 保证多节点不会抢到同一条;
  • Consumer:实际执行任务的工作进程,数量即 limit;
  • Plugin:按周期运行的后台逻辑,如 Cron 定时、Pruner 清理、Lifeline 抢救孤儿任务。

2.2 安装与配置

# config/config.exs
config :my_app, Oban,
  repo: MyApp.Repo,
  queues: [default: 10, mailers: 20, reports: [limit: 2]],
  plugins: [
    Oban.Plugins.Pruner,
    {Oban.Plugins.Cron, crontab: [{"0 3 * * *", MyApp.Workers.Cleanup}]},
    {Oban.Plugins.Lifeline, rescue_after: :timer.minutes(30)}
  ]
# lib/my_app/application.ex
children = [
  MyApp.Repo,
  {Oban, Application.fetch_env!(:my_app, Oban)},
  MyAppWeb.Endpoint
]

数据库迁移通过 Oban 提供的 helper 完成,注意 up/down 都要写:

defmodule MyApp.Repo.Migrations.AddOban do
  use Ecto.Migration

  def up, do: Oban.Migrations.up(version: 12)
  def down, do: Oban.Migrations.down(version: 12)
end

2.3 oban_jobs 表结构

字段类型作用
idbigint任务 ID,插入时即返回
statetextavailable/executing/completed/retryable/discarded/cancelled
queuetext队列名,决定由哪个 Producer 拾取
workertext执行模块名
argsjsonb任务参数(必须是 JSON 可序列化)
prioritysmallint数值越小优先级越高,默认 0
attemptsmallint已尝试次数
max_attemptssmallint最大尝试次数,默认 20
scheduled_attimestamptz计划执行时间
inserted_at / completed_attimestamptz生命周期时间戳
errorsjsonb[]每次失败的错误记录数组

索引建在 (state, queue, priority, scheduled_at, id) 上,这正是 Producer 拉取任务的查询条件。

三、Worker 定义、参数与返回值

3.1 定义一个 Worker

defmodule MyApp.Workers.Mailer do
  use Oban.Worker,
    queue: :mailers,
    max_attempts: 5,
    priority: 1

  @impl Oban.Worker
  def perform(%Oban.Job{args: %{"email" => email, "template" => tpl}}) do
    case MyApp.Mail.deliver(email, tpl) do
      {:ok, _} -> :ok
      {:error, :rate_limited} -> {:snooze, 60}
      {:error, reason} -> {:error, reason}
    end
  end
end

use Oban.Worker 会注入 new/2、new/3 等构造器,并校验参数是否可 JSON 序列化。

3.2 perform 的返回值语义

返回值语义任务状态
:ok成功completed
{:ok, value}成功并记录返回值completed
{:cancel, reason}主动取消,不再重试cancelled
{:discard, reason}丢弃,记入 errors 但不重试discarded
{:error, reason}失败,进入重试队列retryable
{:snooze, seconds}稍后重新调度,不消耗 attemptscheduled
抛异常视为 {:error, exception}retryable

{:snooze, seconds} 是限流场景的利器:它不增加重试次数,适合「等外部配额恢复」这类非错误性的延迟。

3.3 入队

%{email: "user@example.com", template: "welcome"}
|> MyApp.Workers.Mailer.new()          # 立即执行
|> Oban.insert()

MyApp.Workers.Mailer.new(%{...}, schedule_in: 300)                    # 延迟 5 分钟
MyApp.Workers.Mailer.new(%{...}, scheduled_at: ~U[2026-10-06 09:00:00Z])

# 与业务事务同生共死
MyApp.Repo.transaction(fn ->
  user = MyApp.Repo.insert!(changeset)
  Oban.insert!(MyApp.Workers.Mailer.new(%{user_id: user.id}))
  user
end)

最后一段是 Oban 相对外部队列的核心优势:任务入队与业务写入共享同一个事务。若事务回滚,任务也一并消失,不存在「业务失败但邮件已发出」的不一致。

3.4 批量插入

Oban.insert_all/2 接受一组 Oban.Job 结构,走单条 INSERT ... VALUES,比循环 insert/1 快一个数量级:Oban.insert_all(for id <- user_ids, do: MyApp.Workers.Reindex.new(%{user_id: id}))。但它不触发唯一性冲突处理(unique 选项在批量插入中被忽略),需要唯一约束时要改用逐条插入或自行去重。

四、唯一性、重试与优先级

4.1 唯一性约束

Oban 用「部分唯一索引」实现唯一性,配置在 Worker 的 new/2 参数中:

%{user_id: 42}
|> MyApp.Workers.Sync.new(unique: [period: 300, fields: [:user_id], states: [:available, :scheduled]])
|> Oban.insert()
选项默认值含义
period60唯一性窗口(秒);传 :infinity 表示永久唯一
fields[:queue, :worker, :args]参与比较的字段
states[:available, :scheduled, :executing, :retryable]哪些状态算「已存在」
keys无只比较 args 中的指定键,比 fields: [:args] 更宽松

典型用法是防抖:用户连续点击「同步」按钮,300 秒内只会真正入队一次。

4.2 重试退避算法

Oban 默认使用指数退避,第 attempt 次失败后的等待时间为:

delay = attempt^4 + 15 + rand(30) * (attempt + 1)   # 单位:秒
尝试次数大约延迟
1约 15~75 秒
2约 30~150 秒
3约 100~280 秒
5约 640~1400 秒
10约 2.8 小时

可以按 Worker 覆写:

@impl Oban.Worker
def backoff(%Oban.Job{attempt: attempt}) do
  # 固定 30 秒重试,适合依赖快速恢复的场景
  trunc(:math.pow(attempt, 2)) + 30
end

4.3 优先级与队列隔离

优先级只在同一队列内部生效,跨队列比较无意义:

# 高优先级:数值小
MyApp.Workers.Critical.new(%{}, priority: 0)
# 低优先级:数值大
MyApp.Workers.Bulk.new(%{}, priority: 3)

真正的资源隔离靠队列划分。把「用户可感知的快速任务」与「大批量离线任务」放进不同队列,各给独立并发:

queues: [
  interactive: [limit: 50],   # 用户等待的:高并发
  batch: [limit: 5],          # 离线批处理:低并发,避免抢占数据库
  mailers: [limit: 20]
]

4.4 取消与重试的运行时控制

Oban 提供了运行时干预接口:Oban.cancel_job(job_id) 取消尚未执行的任务,Oban.retry_job(job_id) 立即重试一个已 discarded 的任务,Oban.retry_all_jobs(states: [:discarded], queue: :mailers) 批量重试某个队列下的所有丢弃任务。在 perform/1 内部也可以读取 %Oban.Job{attempt: n, max_attempts: m} 来做差异化处理——比如最后一次尝试时改用「降级路径」。

五、Cron 定时、队列隔离与并发

5.1 Cron 插件

{Oban.Plugins.Cron,
 crontab: [
   {"* * * * *", MyApp.Workers.Heartbeat},
   {"*/5 * * * *", MyApp.Workers.SyncInventory},
   {"0 * * * *", MyApp.Workers.HourlyReport},
   {"0 3 * * *", MyApp.Workers.Cleanup, args: %{"days" => 30}},
   {"0 9 * * 1", MyApp.Workers.WeeklyDigest}
 ]}
表达式含义
* * * * *每分钟
*/5 * * * *每 5 分钟
0 * * * *每小时整点
0 3 * * *每天 03:00
0 9 * * 1每周一 09:00
0 0 1 * *每月 1 日 00:00

Cron 插件在每个节点都会运行,因此它内部使用了 Oban 的唯一性机制保证同一时刻只有一个节点插入任务——但前提是你给定时任务配置了唯一性:

defmodule MyApp.Workers.SyncInventory do
  use Oban.Worker,
    queue: :batch,
    unique: [period: 60, states: [:available, :scheduled, :executing]]
end

漏掉 unique 是 Cron 任务最常见的错误:多节点部署下同一分钟会插入 N 条重复任务。

5.2 Quantum 与 Oban Cron 的取舍

维度Oban.Plugins.CronQuantum
存储Postgres(可审计、可重放)内存(重启即丢)
多节点唯一性去重需配置 :global 或单节点运行
执行记录有(oban_jobs 表)无
依赖需要 Oban独立库
适用已有 Postgres 的应用轻量、无 Oban 的场景

只要应用已经用了 Oban,就没有理由再引入 Quantum——Oban Cron 的执行历史可直接用 SQL 查询,排障成本低得多。

5.3 并发与限流

Oban 的并发控制粒度是「队列」,但真实系统往往需要更细的限制,比如「对同一个第三方 API 全局每秒不超过 10 次」。做法是用 {:snooze, n} 配合队列 limit:

defmodule MyApp.Workers.ExternalAPI do
  use Oban.Worker, queue: :external, max_attempts: 3

  @impl Oban.Worker
  def perform(%Oban.Job{args: %{"url" => url}}) do
    case RateLimiter.acquire(:external_api) do
      :ok -> do_request(url)
      :denied -> {:snooze, 1}
    end
  end
end

把队列 limit 设为 10 即可让全局并发不超过 10;配合 :snooze 实现平滑的令牌桶。

5.4 Lifeline 与孤儿任务

进程崩溃(比如节点被 kill -9)时,正在 executing 的任务不会被正常标记,会永远卡住。Oban.Plugins.Lifeline 负责抢救:

{Oban.Plugins.Lifeline, rescue_after: :timer.minutes(30)}

它把超过 rescue_after 仍处于 executing 状态的任务重置为 available 并重新入队。注意:这意味着任务可能被执行两次,所以 perform/1 必须幂等——用唯一约束、幂等键或「先检查状态再操作」来保证。

六、可观测、UI 与 OTP 监督集成

6.1 Telemetry 事件

Oban 全程发出 Telemetry 事件,可以零侵入接入指标:

事件触发时机关键测量值
[:oban, :job, :start]任务开始执行system_time
[:oban, :job, :stop]任务成功duration、queue、worker
[:oban, :job, :exception]任务抛异常kind、reason、duration
[:oban, :engine, :insert, :stop]入队完成duration
[:oban, :plugin, :stop]插件周期运行duration、plugin
:telemetry.attach_many("oban-metrics",
  [[:oban, :job, :stop], [:oban, :job, :exception]],
  fn event, measurements, metadata, _cfg ->
    MyApp.Metrics.record(event, measurements, metadata.queue)
  end, nil)

6.2 关键运维指标

  • 队列积压:SELECT queue, count(*) FROM oban_jobs WHERE state = 'available' GROUP BY queue,持续增长说明消费能力不足;
  • 执行延迟:now() - scheduled_at 对 available 任务取分位数,反映任务从「该执行」到「真执行」的等待;
  • 失败率:exception 事件数除以 stop 事件数;
  • 重试分布:SELECT attempt, count(*) FROM oban_jobs WHERE state = 'retryable' GROUP BY attempt,大量集中在 attempt = max_attempts 说明存在系统性故障;
  • discarded 数量:需要人工介入的任务,应设置告警阈值。

6.3 Oban Web 与测试

oban_web 提供一个 Phoenix LiveView 管理界面,可以查看队列、检索任务、重试失败任务:

# mix.exs
{:oban_web, "~> 2.11", only: [:dev, :prod]}

# router.ex
oban_dashboard("/oban", pipe_through: [:browser, :require_admin])

测试时使用 Oban.Testing,它提供 perform_job/2、assert_enqueued/1 等断言:

# config/test.exs
config :my_app, Oban, testing: :manual

# 测试
assert_enqueued worker: MyApp.Workers.Mailer, args: %{email: "a@b.c"}

# 直接同步执行,断言副作用
assert :ok = perform_job(MyApp.Workers.Mailer, %{email: "a@b.c", template: "welcome"})

testing: :manual 会禁用队列执行,避免测试中真的发邮件。

6.4 与监督树的关系

Oban 自身就是一个完整的 OTP 应用,通过 {Oban, config} 作为子进程挂到你的 supervisor 下。这意味着:

  • 某个 Consumer 崩溃只会导致该任务重试,队列其余部分不受影响;
  • Oban Supervisor 崩溃会被应用级 supervisor 重启,重启后 Producer 会重新从数据库拉取任务,不会丢任务;
  • 多节点部署时,每个节点各自运行一套 Oban 进程,通过数据库的 SKIP LOCKED 协调,天然水平扩展。

这正是 OTP「let it crash」与「状态外置」的经典组合:进程可以随便崩,因为真相在数据库里。

七、最佳实践与总结

  • 任务参数只放 ID,不放业务对象:args 要序列化进数据库,传 user_id 让 Worker 自己去查,避免参数过期与体积膨胀;
  • perform/1 必须幂等:Lifeline 抢救、网络超时重投、人工重试都会导致重复执行,用唯一约束或状态机保证;
  • 入队与业务写同一事务:Repo.transaction 里同时写业务数据和 Oban.insert,杜绝不一致;
  • 定时任务必须配 unique:多节点 + Cron 的组合下,没有唯一性就会重复执行;
  • 队列按「资源特征」划分:慢任务与快任务分队列,避免离线批处理拖垮用户可感知的交互任务;
  • 优先用 {:snooze, n} 而非 {:error, reason}:限流与等待场景下 snooze 不消耗重试次数,语义更准确;
  • 可观测先行:把积压量、执行延迟、失败率画成看板,接入方式与 https://plumephp.com/erlang-logging-telemetry-observability/ 一致;
  • 重试与外部调用的组合:调用外部 HTTP 时的超时、熔断策略见 https://plumephp.com/erlang-http-client-pooling/。

Oban 的成功在于它没有发明新的基础设施,而是把「可靠队列」这件难事建立在一个已经足够可靠的系统之上。事务保证入队原子性,SKIP LOCKED 保证多节点安全,唯一索引保证去重,OTP 保证执行器容错。理解了这套组合,你就能用最小的运维成本,把同步请求里的副作用全部搬到后台,让 Web 进程只做它最擅长的事——快速响应。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「erlang」更多文章

  1. Erlang/Elixir 容器化与集群部署:Release、Docker 与 libcluster
  2. Elixir 认证授权实战:JWT、Guardian 与 Phoenix.Token
  3. Elixir HTTP 客户端与连接池:Mint、Finch 与 Req 实战