Oban 任务队列实战
· · ·
把一次 HTTP 请求,变成可靠、可重试、可观测的后台工作。
题目进度 0 / 0 ✓ 0
Oban 是 Elixir 生态里基于 PostgreSQL 的后台任务系统。它不是简单的内存队列:任务写入数据库后,即使应用进程重启,任务仍然可以被恢复和处理。
本课目标
学完后,你应该能解释 Oban 的队列模型,写出一个基本 Worker,正确处理重试、幂等、唯一约束,并知道生产环境如何监控失败任务。
一、为什么需要后台任务
发送邮件、生成报告、拉取公告、调用 LLM、同步第三方 API 都可能慢、会失败,或者不适合阻塞 HTTP 请求。
# 不推荐:请求一直等到外部 API 完成
post "/reports" do
report = Report.generate(params)
ExternalApi.upload(report)
send_resp(conn, 201, "ok")
end
# 推荐:请求只负责入队
post "/reports" do
%{report_id: params["id"]}
|> GenerateReport.new()
|> Oban.insert!()
send_resp(conn, 202, "accepted")
end
202 Accepted 表示请求已接受,实际工作稍后执行。客户端可以通过任务状态接口或通知获知结果。
二、Oban 的核心模型
| 概念 | 含义 | 关键问题 |
|---|---|---|
| Job | 一条待执行或已执行的任务记录 | 参数是否可序列化? |
| Worker | 定义任务如何执行的模块 | 失败后是否安全重试? |
| Queue | 按业务能力隔离并发度 | 这个队列允许同时跑多少个? |
| Plugin | 调度、清理、限流等后台组件 | 失败、过期、卡住的 Job 怎么处理? |
三、写一个 Worker
defmodule Pipeline.GenerateReport do
use Oban.Worker,
queue: :reports,
max_attempts: 3
@impl Oban.Worker
def perform(%Oban.Job{args: %{"report_id" => report_id}}) do
with {:ok, report} <- Reports.load(report_id),
{:ok, file} <- Reports.render(report),
:ok <- Reports.publish(file) do
:ok
end
end
end
use Oban.Worker注入 Worker 所需的行为和配置。args是持久化 JSON,尽量传 ID,不要把巨大对象塞进 Job。- 成功返回
:ok。 - 返回错误元组或抛出异常时,Oban 会根据配置安排重试。
四、入队:不要绕过数据库
job = Pipeline.GenerateReport.new(%{report_id: report.id})
case Oban.insert(job) do
{:ok, job} -> {:ok, job.id}
{:error, changeset} -> {:error, changeset}
end
# 多个任务可以放在同一个事务里
Repo.transaction(fn ->
{:ok, order} = Orders.create(attrs)
{:ok, _job} = Oban.insert(Pipeline.SendReceipt.new(%{order_id: order.id}))
end)
把业务写入和任务入队放进同一事务,可以避免“订单已创建但任务没入队”的不一致。
五、重试不是魔法:必须保证幂等
同一个 Job 可能执行多次。网络超时尤其危险:外部系统可能已经成功,但本地没收到响应,于是任务被重试。
# 不安全:每次重试都可能重复扣款
Payment.charge(card, amount)
# 更安全:使用业务幂等键
Payment.charge(card, amount, idempotency_key: "order:#{order.id}")
写入数据库时使用唯一约束、insert_all(..., on_conflict: :nothing) 或状态机,确保重复执行不会产生重复副作用。
六、唯一任务、限流与队列隔离
use Oban.Worker,
queue: :announcements,
max_attempts: 5,
unique: [period: 300, fields: [:args, :worker], states: [:available, :scheduled]]
- unique 防止相同任务被重复排队,但它不是幂等性的替代品。
- CPU 密集任务、外部 API 任务、邮件任务应使用不同队列。
- 队列并发度要服从下游限制,不能只按本机 CPU 配置。
七、失败处理与可观测性
生产环境至少要观察:
- available / scheduled / executing / retryable / discarded 的 Job 数量。
- 每个队列的执行延迟、成功率、平均耗时。
- 连续失败的 Worker、最终 discarded 的任务和异常原因。
- 外部 API 的限流、超时和响应码。
失败任务不能只写日志。日志用于定位,数据库中的 Job 状态用于审计,指标和告警用于及时发现问题。
八、启动配置与生产清单
config :pipeline, Oban,
repo: Pipeline.Repo,
queues: [
reports: 5,
announcements: 2,
default: 10
],
plugins: [
{Oban.Plugins.Pruner, max_age: 60 * 60 * 24 * 7},
{Oban.Plugins.Lifeline, rescue_after: :timer.minutes(30)}
]
- 确认 Oban migration 已执行,且生产数据库有正确索引。
- 为慢任务设置合理超时,避免执行进程无限占用队列。
- 部署新版本时兼容旧 Job 的 args;不要随意删除正在排队的 Worker 模块。
- 数据库是队列的一部分,必须纳入备份、容量和连接池规划。
九、测验
1
Oban 的 Job 主要持久化在哪里?
2
后台任务接口通常返回 HTTP 202 的含义是?
3
Worker 的 perform/1 成功执行应返回?
4
为什么 Job 应尽量只保存业务 ID?
5
重试任务最重要的业务前提是?
6
unique 配置主要解决什么问题?
7
业务写入和 Oban.insert 放在同一事务里的主要好处是?
8
为什么不同类型任务要隔离到不同队列?
9
哪一项最适合作为失败任务告警信号?
10
部署新版本时,正在排队的旧 Job 需要特别注意什么?
**下一步:**把 Oban Worker 接入真实业务时,先画出状态、重试和幂等边界,再写代码。