Learning
VOL. XIII · NO. 16 · Elixir · 01 JAN 1970

Oban 任务队列实战

Elixir 编程 · 01 JAN 1970 · 6 min read · 1,334 words
· · ·

把一次 HTTP 请求,变成可靠、可重试、可观测的后台工作。

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 接入真实业务时,先画出状态、重试和幂等边界,再写代码。