Elixir OTP 监督树与任务并发:Supervisor、GenServer 与 Task

Elixir OTP 监督树与任务并发:Supervisor 策略(one_for_one/one_for_all/rest_for_one/simple_one_for_one)、GenServer 状态机(init/handle_call/handle_cast/handle_info)、Task 异步并发(async/await/yield)、Agent 轻量状态、DynamicSupervisor 动态子进程、Registry 进程发现、OTP 应用结构(Application 回调/规格)、错误隔离与重启强度(MaxR/MaxT)、分布式 Erlang 节点通信、与 Phoenix 的集成模式。

引言

Elixir 站在 Erlang VM(BEAM)的肩膀上,继承了「Let it crash」的容错哲学,但通过更现代的语法和更丰富的标准库降低了入门门槛。OTP(Open Telecom Platform)是 Erlang 的运行时框架——Supervisor 监督进程树、GenServer 通用服务器、Task 轻量并发任务。掌握这三个核心抽象,你就能构建「故障隔离、自动恢复、水平扩展」的分布式系统。本文聚焦 Elixir 中的 OTP 实践:从 Supervisor 策略到 GenServer 设计模式,从 Task 并发到分布式节点通信。

前置:/erlang-language-basics/(Erlang 基础)、/erlang-bit-syntax-binaries/(位语法与二进制)。


目录


1. OTP 核心三元组:Supervisor、GenServer、Task

1.1 为什么 OTP

# OTP 不是库,是「构建容错系统的框架」
# 核心设计:监督树(Supervision Tree)
#   叶子节点:Worker(GenServer/Agent/Task)做实际工作
#   中间节点:Supervisor 监督子进程,崩溃时按策略重启
#   根节点:Application 管理整个 supervision tree
# 目标:局部故障不扩散,自动恢复不人工

1.2 三元组分工

组件职责类比
Supervisor监督子进程,崩溃时重启项目经理
GenServer维护状态,处理同步/异步消息有状态服务
Task一次性异步计算临时工

1.3 启动链路

Application.start → Supervisor.init → 启动子进程(GenServer/Task/其他 Supervisor)
→ 子进程注册名字 → 系统进入运行状态
# 全部失败则 Application 启动失败 → 无半启动状态

记忆 OTP 是容错框架不是库——监督树(Worker 干活、Supervisor 监工、Application 管整体);三元组 Supervisor(重启策略)+ GenServer(有状态服务)+ Task(一次性异步);启动链路 Application→Supervisor→Workers,全成功才启动。


2. GenServer 状态机与消息处理

2.1 最简 GenServer

defmodule Counter do
  use GenServer

  # 客户端 API
  def start_link(init_val), do: GenServer.start_link(__MODULE__, init_val, name: __MODULE__)
  def increment, do: GenServer.cast(__MODULE__, :increment)
  def get, do: GenServer.call(__MODULE__, :get)

  # 服务端回调
  def init(init_val), do: {:ok, init_val}

  def handle_cast(:increment, state) do
    {:noreply, state + 1}
  end

  def handle_call(:get, _from, state) do
    {:reply, state, state}
  end
end

2.2 回调对照表

回调触发用途
init/1启动初始化状态
handle_call/3call 请求同步请求,需回复
handle_cast/2cast 请求异步请求,不回复
handle_info/2任意消息超时、DOWN 消息、外部信号
terminate/2停止清理资源(尽量不用)
code_change/3热升级状态迁移

2.3 call vs cast

call:同步 → 等回复 → 用 handle_call(需 from 参数,{:reply, reply, state})
cast:异步 → 发完即走 → 用 handle_cast({:noreply, state})
info:其他进程发的裸消息、定时器超时、监控 DOWN 信号 → handle_info

记忆 GenServer = 客户端 API(封装 call/cast)+ 服务端回调(handle_call/handle_cast/handle_info);call 同步等回复、cast 异步即发、info 处理其他消息;回调返回 {:reply, _, _} / {:noreply, _} / {:stop, _, _}。


3. Supervisor 监督策略与重启语义

3.1 四种策略

:one_for_one      — 一个子进程崩溃,只重启它
:one_for_all      — 一个子进程崩溃,重启全部子进程
:rest_for_one     — 一个子进程崩溃,重启它和它后面的所有子进程
:simple_one_for_one — 动态添加同类型子进程( Worker 池)

3.2 配置示例

defmodule MyApp.Supervisor do
  use Supervisor

  def init(_args) do
    children = [
      {Counter, 0},
      {MyApp.Repo, []},
      {Task.Supervisor, name: MyApp.TaskSupervisor}
    ]

    Supervisor.init(children,
      strategy: :one_for_one,
      max_restarts: 5,      # 5 次重启
      max_seconds: 10       # 10 秒内
    )
  end
end

3.3 重启语义

transient:只在异常退出时重启(正常 exit 不重启)
temporary:永不重启
permanent:总是重启(默认)

# 生命周期:
# init 失败 → Supervisor 标记子进程失败 → 根据策略重启
# 超过 MaxR/MaxT → Supervisor 自己也崩溃,向上传播

记忆 四种策略——one_for_one(各管各)、one_for_all(一损俱损)、rest_for_one(级联重启)、simple_one_for_one(动态池);重启语义 transient(异常才重启)/temporary(不重启)/permanent(总重启);MaxR/MaxT 防重启风暴。


4. Task 异步并发与超时控制

4.1 基础用法

# 启动异步任务
task = Task.async(fn -> HTTPoison.get!("https://api.example.com") end)

# 等待结果
result = Task.await(task, 5000)  # 5 秒超时

# 多个并发
tasks = [
  Task.async(fn -> fetch_user(1) end),
  Task.async(fn -> fetch_user(2) end),
  Task.async(fn -> fetch_user(3) end)
]
results = Task.yield_many(tasks, 5000)

4.2 受监督的 Task

# Task.Supervisor 下启动 → 崩溃被 Supervisor 处理
task = Task.Supervisor.async(MyApp.TaskSupervisor, fn ->
  risky_operation()
end)

Task.await(task)

# Task.Supervisor.start_child 不链接(fire-and-forget)
Task.Supervisor.start_child(MyApp.TaskSupervisor, fn ->
  send_notification()
end)

4.3 超时与取消

# 超时处理
result = Task.yield(task, 1000) || Task.shutdown(task)

# Task.yield 返回 nil=任务仍在跑,{:ok, val}=完成
# Task.shutdown 发退出信号终止任务

记忆 Task.async 启动异步计算、Task.await 等结果带超时、Task.yield_many 批量并发;受监督 Task 用 Task.Supervisor.async(崩溃被重启);超时用 yield 检查 + shutdown 终止。


5. Agent 轻量状态与简易并发原语

5.1 什么时候用 Agent

# 只需要存一个值 + 读/写操作
# 不需要复杂消息处理 → 比 GenServer 轻量
# 例:配置缓存、计数器、简单状态机

5.2 用法

{:ok, pid} = Agent.start_link(fn -> %{} end, name: :config)

# 获取
Agent.get(:config, &Map.get(&1, :key))

# 更新
Agent.update(:config, &Map.put(&1, :key, "value"))

# 获取并更新(原子操作)
Agent.get_and_update(:config, fn state ->
  {state[:key], Map.put(state, :key, "new")}
end)

记忆 Agent 是轻量 GenServer——只存一个值 + 读/写/get_and_update;不需要复杂消息处理时代码比 GenServer 少很多;配置缓存/计数器/简单状态适合 Agent。


6. DynamicSupervisor 与动态进程管理

6.1 为什么需要 DynamicSupervisor

# Supervisor 的 children 在 init 时固定
# DynamicSupervisor:运行时动态添加/删除子进程
# 用例:WebSocket 连接池、长轮询客户端、临时 Worker

6.2 配置与使用

defmodule MyApp.DynamicSupervisor do
  use DynamicSupervisor

  def start_link(init_arg) do
    DynamicSupervisor.start_link(__MODULE__, init_arg, name: __MODULE__)
  end

  def init(_init_arg) do
    DynamicSupervisor.init(strategy: :one_for_one)
  end

  def start_child(spec) do
    DynamicSupervisor.start_child(__MODULE__, spec)
  end
end

# 使用
MyApp.DynamicSupervisor.start_child(
  {MyWorker, arg}
)

记忆 DynamicSupervisor 运行时动态增删子进程(init 时为空);WebSocket 连接/临时 Worker/客户端池用 DynamicSupervisor;start_child 动态添加、terminate_child 删除。


7. Registry 进程发现与名称服务

7.1 本地 Registry

# 启动
Registry.start_link(keys: :unique, name: MyApp.Registry)

# 注册进程
Registry.register(MyApp.Registry, "user_123", %{})

# 查找进程
Registry.lookup(MyApp.Registry, "user_123")
# => [{pid, value}]

# 通过 Registry 发送消息
[{pid, _}] = Registry.lookup(MyApp.Registry, "user_123")
send(pid, :ping)

7.2 分区 Registry(:duplicate)

# :duplicate 允许一个 key 对应多个进程
# 用例:Pub/Sub 订阅者列表、房间在线用户
Registry.lookup(MyApp.PubSub, "room:general")
# => [{pid1, sub1}, {pid2, sub2}, ...]

记忆 Registry 是本地进程发现——:unique 一对一(用户会话)、:duplicate 一对多(Pub/Sub 订阅);register 注册、lookup 查找、通过查到的 pid 发消息。


8. OTP Application 结构与生命周期

8.1 Application 结构

my_app/
├── lib/
│   ├── my_app.ex           # Application 回调模块
│   ├── my_app/
│   │   ├── supervisor.ex   # 顶级 Supervisor
│   │   ├── worker.ex       # GenServer Worker
│   │   └── ...
├── mix.exs                 # 项目配置 + application 规格

8.2 Application 回调

defmodule MyApp.Application do
  use Application

  def start(_type, _args) do
    children = [
      MyApp.Repo,
      MyAppWeb.Endpoint,
      {Phoenix.PubSub, name: MyApp.PubSub}
    ]
    Supervisor.start_link(children, strategy: :one_for_one, name: MyApp.Supervisor)
  end
end

# mix.exs
  def application do
    [
      mod: {MyApp.Application, []},
      extra_applications: [:logger, :runtime_tools]
    ]
  end

记忆 Application = 入口(start 返回 Supervisor)+ mix.exs(mod 指定回调);OTP app 有标准目录结构;启动失败则整个 app 不启动,不存在半启动状态。


9. 分布式 Erlang 节点通信

9.1 节点命名与连接

# 启动带名字的节点
iex --sname node1@localhost

# 在 iex 中连接
Node.connect(:node2@localhost)   # => true/false
Node.list()                       # => [:node2@localhost]

9.2 跨节点调用

# 在 node1 上调用 node2 的模块
GenServer.call({Counter, :node2@localhost}, :get)

# 跨节点 spawn
Node.spawn(:node2@localhost, fn -> IO.puts("Hello from node2") end)

9.3 分布式注意事项

# 安全:默认 cookie 认证,生产环境用 TLS(-proto_dist inet_tls)
# 网络分区:脑裂时用 net_kernel 监测,设计分区容忍
# 不适用于:跨地域高延迟(用 GenStage/CQRS 替代)

记忆 分布式——iex –sname 启动命名节点、Node.connect 连接、跨节点调用 {Module, node};安全用 TLS、防脑裂、跨地域高延迟不适合直接分布式。


10. 速查表与一句话记忆

概念一句话
Supervisor监督子进程,崩溃重启
one_for_one各管各
one_for_all一损俱损
rest_for_one级联
simple_one_for_one动态池
GenServer有状态服务
call同步等回复
cast异步即发
Task一次性异步
Agent轻量存取值
DynamicSupervisor运行时增删子进程
Registry本地进程发现
ApplicationOTP 应用入口
Node.connect分布式节点连接

一句话记忆:OTP 是容错框架——Supervisor(监工重启)+ GenServer(有状态服务:call 同步/cast 异步/info 其他消息)+ Task(一次性异步:async/await/yield_many);四种策略 one_for_one(各管各)/one_for_all(全重启)/rest_for_one(级联)/simple_one_for_one(动态池);Agent 轻量存取值、DynamicSupervisor 运行时增删子进程(WebSocket 池)、Registry 本地进程发现(unique 一对一/duplicate 一对多);Application.start 是整个监督树的根,全成功才启动;分布式用 –sname + Node.connect + {Module, node} 跨节点调用——「OTP 教会 Erlang/Elixir 的不是如何不崩溃,而是崩溃后如何优雅地站起来」。


延伸阅读

继续阅读

探索更多技术文章

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

全部文章 返回首页

「erlang」更多文章

  1. OTP 应用设计模式:监督树结构、release 打包与热升级
  2. Erlang/Elixir 安全加固:加密、认证与分布式信任
  3. Ecto 高级查询与数据库工程:关联、多态与性能优化