几乎每一个 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 的分工
| 框架 | 数据来源 | 触发方式 | 状态存储 | 适用 |
|---|---|---|---|---|
| Oban | Postgres 表 | 入队 + 定时 | 数据库 | 业务副作用、定时任务 |
| Broadway | Kafka/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 表结构
| 字段 | 类型 | 作用 |
|---|---|---|
id | bigint | 任务 ID,插入时即返回 |
state | text | available/executing/completed/retryable/discarded/cancelled |
queue | text | 队列名,决定由哪个 Producer 拾取 |
worker | text | 执行模块名 |
args | jsonb | 任务参数(必须是 JSON 可序列化) |
priority | smallint | 数值越小优先级越高,默认 0 |
attempt | smallint | 已尝试次数 |
max_attempts | smallint | 最大尝试次数,默认 20 |
scheduled_at | timestamptz | 计划执行时间 |
inserted_at / completed_at | timestamptz | 生命周期时间戳 |
errors | jsonb[] | 每次失败的错误记录数组 |
索引建在 (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} | 稍后重新调度,不消耗 attempt | scheduled |
| 抛异常 | 视为 {: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()
| 选项 | 默认值 | 含义 |
|---|---|---|
period | 60 | 唯一性窗口(秒);传 :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.Cron | Quantum |
|---|---|---|
| 存储 | 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 进程只做它最擅长的事——快速响应。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。