状态机是表达复杂业务流程最自然的抽象:连接建立与断开、支付状态流转、协议握手、任务调度,都能用一组有限状态和状态间转移精确描述。OTP 从 gen_fsm 演进到 gen_statem,把状态机的表达能力提升到了新的高度——它同时支持 State Functions 与 Handle Event 两种回调模式,内置延迟事件、状态超时、事后延迟(postpone)等机制。与之互补的 gen_event 则提供了多处理器订阅事件流的发布-订阅框架。本文将从设计动机出发,深入这两种行为模式的回调契约、运行时语义与生产实践。
一、OTP 状态机行为演进
1.1 为什么需要状态机抽象
在并发系统中,进程常常需要根据当前「所处阶段」对相同消息做出不同响应。若用 gen_server 的 handle_call/3 把所有分支写在一个函数里,状态爆炸后代码会迅速腐烂。状态机将「状态」提升为一等公民:
| 维度 | 朴素 gen_server | 状态机 |
|---|---|---|
| 状态表示 | State 记录中的某个字段 | 运行时状态名(StateName) |
| 事件分发 | 手动 case 嵌套 | 按状态名自动路由到对应回调 |
| 状态合法转移 | 编译期无法约束 | 由回调返回值显式声明 |
| 超时/延迟 | 手写 handle_info 定时器 | 内建 timeout 与 state_timeout |
1.2 从 gen_fsm 到 gen_statem
gen_statem(OTP 19+)是 gen_fsm 的全面替代:
gen_fsm强制使用 State Functions 模式,回调固定为Module:StateName/2;gen_statem支持 state_functions 与 handle_event_function 两种模式,并且可以叠加 state_enter(进入状态回调);gen_statem用 Actions 列表统一表达回复、延迟事件、状态超时等副作用,不再需要gen_fsm中send_event_after/2的隐式定时器。
-module(gateway).
-behaviour(gen_statem).
-export([start_link/0, connect/1, send/2, close/1]).
-export([callback_mode/0, init/1, terminate/3, code_change/4]).
%% 状态函数
-export([disconnected/3, connecting/3, connected/3]).
-define(SERVER, ?MODULE).
-record(data, {
host,
port,
retries = 0,
socket = undefined
}).
二、State Functions 回调模式
2.1 回调模式声明
callback_mode/0 决定 gen_statem 如何分发事件:
%% 默认:状态函数模式。每个状态对应一个 Module:StateName/3 回调
callback_mode() -> state_functions.
%% 也可以返回列表,叠加进入状态回调:
%% callback_mode() -> [state_functions, state_enter].
在 state_functions 模式下,gen_statem 收到事件后,根据当前状态 StateName 调用 Module:StateName(EventType, EventContent, Data)。同一事件在不同状态下的行为被自然切分。
2.2 事件类型
gen_statem 将输入消息统一抽象为四类事件:
| 事件类型 | 来源 | EventContent 含义 |
|---|---|---|
{call, From} | gen_statem:call/2,3 | 调用参数,可同步回复 |
cast | gen_statem:cast/2 | 异步消息内容 |
info | 直接 ! 发送或端口/定时器消息 | 任意 Erlang term |
{timeout, TimerRef} | 内部超时事件(由 timeout Action 产生) | 预设的事件内容 |
{state_timeout, TimerRef} | 状态超时事件 | 预设的事件内容 |
%% 断开连接状态下:
%% - 收到 connect 请求 → 进入 connecting 状态并启动连接
%% - 其余消息全部忽略
disconnected({call, From}, connect, Data) ->
{keep_state, Data,
[{next_event, internal, do_connect}, % 内部事件,切换状态后处理
{reply, From, ok}]};
disconnected(cast, _Msg, Data) ->
io:format("Ignored cast in disconnected~n"),
{keep_state, Data};
disconnected(info, _Info, Data) ->
{keep_state, Data}.
2.3 状态函数返回值
每个状态函数必须返回一个「状态转移元组」,gen_statem 据此决定下一步:
| 返回值 | 语义 |
|---|---|
{next_state, NewState, Data} | 转移状态并保持数据 |
{keep_state, Data} | 保持当前状态,更新数据 |
{repeat_state, Data} | 保持状态,但触发 state_enter(若启用) |
{keep_state_and_data, ...} | 保持状态与数据(等价 {keep_state, Data} 但无需显式带出 Data) |
{stop, Reason} / {stop, Reason, Data} | 终止状态机 |
repeat_state 与 keep_state 的区别仅在启用了 state_enter 时才有意义——前者会再次调用 Module:enter_State/4。
2.4 完整连接状态机
下面是一个带重试退避的 TCP 连接状态机:
%% 进入 connecting:记录重试时间并启动状态超时
connecting({call, _From}, _Event, Data) ->
{keep_state, Data, state_timeout(Data#data.retries)};
connecting(state_timeout, retry, #data{retries = R} = Data) ->
case tcp_connect(Data#data.host, Data#data.port) of
{ok, Socket} ->
{next_state, connected, Data#data{socket = Socket, retries = 0}};
{error, Reason} when R < 5 ->
NewRetries = R + 1,
io:format("Retry ~p after ~p~n", [NewRetries, backoff(NewRetries)]),
{keep_state, Data#data{retries = NewRetries},
[{state_timeout, backoff(NewRetries), retry}]};
{error, Reason} ->
{stop, {connection_failed, Reason}}
end;
connecting(state_timeout, _Event, Data) ->
%% 收到了未知的 state_timeout 事件,忽略
{keep_state, Data}.
%% 已连接状态:转发数据
connected({call, From}, {send, Payload}, #data{socket = S} = Data) ->
case gen_tcp:send(S, Payload) of
ok ->
{keep_state, Data, [{reply, From, ok}]};
{error, Reason} ->
{next_state, connecting, Data,
[{reply, From, {error, Reason}},
{state_timeout, 0, retry}]} % 立即重连
end;
connected(info, {tcp_closed, _Socket}, Data) ->
{next_state, disconnected, Data};
connected(info, {tcp_data, Socket, Bin}, Data) ->
handle_incoming(Bin),
{keep_state, Data}.
%% 辅助:指数退避
backoff(Attempt) -> 100 * trunc(math:pow(2, Attempt)).
state_timeout(Retries) ->
[{state_timeout, backoff(Retries), retry}].
tcp_connect(Host, Port) ->
gen_tcp:connect(Host, Port, [binary, {active, true}], 5000).
handle_incoming(_Bin) -> ok.
三、Handle Event 回调模式
3.1 单一回调分发
对于状态众多但转移逻辑高度对称的状态机,state_functions 会产生大量近乎空转的 StateName/3 回调。此时可以用 handle_event_function 模式,把所有事件集中到 handle_event/4:
callback_mode() -> handle_event_function.
%% handle_event(EventType, EventContent, StateName, Data)
handle_event({call, From}, get_status, State, Data) ->
{keep_state, Data, [{reply, From, State}]};
handle_event(state_timeout, retry, connecting, #data{retries = R} = Data) ->
%% 重试逻辑与 2.4 相同
{keep_state, Data, ...};
handle_event(cast, _Msg, _State, Data) ->
{keep_state, Data};
handle_event(info, {tcp_closed, _}, _State, Data) ->
{next_state, disconnected, Data};
3.2 两种模式的选择
| 对比维度 | state_functions | handle_event_function |
|---|---|---|
| 代码组织 | 按状态分组,天然隔离 | 按事件分组,集中处理 |
| 状态数量 | 越多越清晰 | 状态多时 case 膨胀 |
| 事件共享逻辑 | 难以复用 | 可在同一函数内统一分支 |
| 适用场景 | 协议栈、支付流转 | 事务、命令路由 |
实际项目中,若各状态共享大量公共处理逻辑(如日志、审计、鉴权),
handle_event_function更合适;若状态间行为差异极大且互不共享,state_functions可读性更好。
四、延迟事件与状态超时
4.1 postpone 事后延迟
gen_statem 最强大的机制之一是 postpone:把当前事件「暂时搁置」,等状态转移完成后再重新投递。这用于处理到达过早的事件:
%% 场景:连接到一半时来了 send 请求,应当等 connected 之后再执行
connecting(cast, {send, _} = Msg, Data) ->
{keep_state, Data, [{postpone, true}]}; % 挂起该事件
connecting(state_timeout, retry, Data) ->
case tcp_connect(Data#data.host, Data#data.port) of
{ok, Socket} ->
{next_state, connected, Data#data{socket = Socket},
[{state_timeout, 5000, idle_timeout}]};
{error, _} ->
{keep_state, Data, [{state_timeout, 1000, retry}]}
end;
%% 转移到 connected 后,之前被 postpone 的 {send, _} 会按投递顺序重新进入
connected(cast, {send, Payload}, #data{socket = S} = Data) ->
gen_tcp:send(S, Payload),
{keep_state, Data};
postpone 的关键语义:事件不会丢失、不会乱序,只是在状态转移完成后重新入队。相比「把事件缓存到 Data 中再手动处理」,它避免了重复逻辑且无需关心队列清理。
4.2 三类定时机制
gen_statem 内置三类定时器,覆盖不同生命周期:
| Action | 生命周期 | 典型用途 |
|---|---|---|
{timeout, Time, Event} | 单次,进程级 | 请求超时、会话超时 |
{state_timeout, Time, Event} | 状态相关,转移即失效 | 心跳、重连退避 |
{event_timeout, Time, Event} | 事件级,新事件重置 | 客户端响应窗口 |
%% 状态级心跳:进入 connected 后每 5 秒触发一次,
%% 一旦转移状态(如 tcp_closed)自动取消
connected(state_timeout, heartbeat, #data{socket = S} = Data) ->
gen_tcp:send(S, <<"ping">>),
{keep_state, Data, [{state_timeout, 5000, heartbeat}]};
%% 进程级请求超时:每次收到 call 都给一个 30 秒截止
{call, From} = Event,
{keep_state, Data, [{timeout, 30000, {request_timeout, From}},
{reply, From, processing}]}
%% 事件级超时:5 秒内没有新事件才触发
{keep_state, Data, [{event_timeout, 5000, idle}]}
state_timeout 与普通 timeout 的最大区别:状态转移(next_state)会自动取消未触发的 state_timeout,非常适合「此状态内必须完成」的约束;而 timeout 与状态无关,只有显式 keep_state_and_data 后取消或重置。
4.3 使用 next_event 串联内部流程
通过 {next_event, Type, Content} 可以编排多步骤内部流程,让状态机如同事件驱动的工作流:
init(_Args) ->
{ok, disconnected, #data{},
[{next_event, cast, bootstrap}]}.
disconnected(cast, bootstrap, Data) ->
{keep_state, Data,
[{next_event, internal, load_config},
{next_event, internal, open_db}]};
disconnected(internal, load_config, Data) ->
Config = read_config(),
{keep_state, Data#data{config = Config}};
disconnected(internal, open_db, Data) ->
case db:connect(Data#data.config) of
ok ->
{next_state, ready, Data};
{error, _} ->
{keep_state, Data, [{state_timeout, 1000, retry_db}]}
end.
五、gen_event 事件管理
5.1 事件管理器架构
gen_event 实现多处理器发布-订阅:一个事件管理器(Event Manager)进程持有若干 Handler,向管理器 notify 的事件会被广播给所有 Handler。Handler 之间互不可见、互不影响,可以动态增删、热替换。
gen_event 事件管理器
│
┌────────────┼────────────┐
▼ ▼ ▼
Handler A Handler B Handler C
(审计日志) (指标上报) (告警通知)
5.2 实现一个 Handler
-module(metrics_handler).
-behaviour(gen_event).
%% 回调接口
-export([init/1, handle_event/2, handle_call/2, handle_info/2,
terminate/2, code_change/3]).
init([]) ->
{ok, #{count => 0}}.
handle_event({http_request, Path, Status, Latency}, State) ->
Count = maps:get(count, State) + 1,
update_metrics(Path, Status, Latency),
{ok, State#{count => Count}};
handle_event(_, State) ->
{ok, State}.
handle_call(get_count, State) ->
{ok, maps:get(count, State), State};
handle_call(_Request, State) ->
{ok, {error, bad_request}, State}.
handle_info(_Info, State) ->
{ok, State}.
terminate(_Reason, _State) ->
ok.
code_change(_OldVsn, State, _Extra) ->
{ok, State}.
update_metrics(_Path, _Status, _Latency) -> ok.
5.3 动态增删与通知
%% 启动事件管理器
{ok, Mgr} = gen_event:start_link().
%% 添加处理器(管理器启动时会调用 Handler:init/1)
gen_event:add_handler(Mgr, metrics_handler, []).
gen_event:add_handler(Mgr, audit_handler, ["./audit.log"]).
%% 异步广播(不等待 Handler 完成)
gen_event:notify(Mgr, {http_request, "/api/users", 200, 42}).
%% 同步广播(等待所有 Handler 处理完)
gen_event:sync_notify(Mgr, {flush, self()}).
%% 带返回值的调用(Handler 的 handle_call/2)
gen_event:call(Mgr, metrics_handler, get_count).
%% 移除 / 替换 / 查询
gen_event:delete_handler(Mgr, audit_handler, stop_reason).
gen_event:swap_handler(Mgr, {old_handler, Args}, {new_handler, Args2}).
gen_event:which_handlers(Mgr).
5.4 故障隔离与监督
Handler 抛异常时,gen_event 默认将其移除并调用 terminate/2,其余 Handler 不受影响。若希望 Handler 崩溃时保留,可以用 gen_event:add_sup_handler/3——它与管理器互相监控,任一终止都会通知对方:
%% 监督型 Handler:Handler 崩溃会通知 Manager,Manager 崩溃会通知 Handler
{ok, _} = gen_event:add_sup_handler(Mgr, critical_handler, Args).
%% 在 Handler 中接收 manager 终止通知
handle_info({'EXIT', Mgr, Reason}, State) ->
%% 管理器挂了,决定自身行为
{ok, State}.
六、gen_server / gen_statem / gen_event 对比
三者都是 OTP 行为,但抽象层次不同:
| 对比维度 | gen_server | gen_statem | gen_event |
|---|---|---|---|
| 核心模型 | 客户端-服务器 | 有限状态机 | 发布-订阅 |
| 状态维度 | 单一状态 State | StateName + Data | 每个 Handler 独立状态 |
| 同步调用 | handle_call/3 | handle_event({call, _}, ...) | handle_call/2 |
| 异步消息 | handle_cast/2 | handle_event(cast, ...) | handle_event/2 |
| 定时器 | 手动 send_after | timeout / state_timeout | 手动 |
| 事件消费者 | 单个 | 单个 | 多个(广播) |
| 动态热替换 | 不支持 | 不支持 | swap_handler/3 |
| 适用场景 | 通用服务、资源管理 | 协议、流程、分布式事务 | 日志、指标、插件系统 |
选择建议:
- 需要对外暴露一组同步 API + 内部状态,用
gen_server; - 业务存在显式状态与合法转移集合,用
gen_statem; - 需要把同一事件派发给多个独立关注者,用
gen_event; - 状态机需要动态安装/卸载处理器(如插件、采集器),组合使用
gen_statem管理生命周期 +gen_event做广播。
七、最佳实践与总结
7.1 状态机建模要点
- 先画状态转移图再写代码:把合法转移枚举出来,非法事件在回调用兜底分支忽略或记日志;
- 状态数据最小化:
Data只放状态转移所需信息,避免塞入业务对象导致状态污染; - 善用 state_timeout 表达「状态内截止」:比手写定时器更安全,状态转移自动清理;
- postpone 替代手动缓冲:解决「事件早到」场景,避免在
Data里维护待处理队列; - 控制台用
gen_statem:start_link观察:sys:get_state/1可直接查看{StateName, Data}。
7.2 事件驱动设计要点
- Handler 必须无状态化或状态可重建,因为
gen_event可能在异常后被移除重建; - 敏感操作用
sync_notify,高频低优先级事件用notify避免阻塞生产者; - 事件结构建议采用
{Tag, Payload}元组,Handler 只关心自己订阅的 Tag; gen_event不是分布式总线——跨节点广播用pg(Process Groups)或global,可参考 https://plumephp.com/posts/distributed-systems/ 中的事件驱动模式。
7.3 结合 OTP 全家桶
gen_statem 与 gen_event 常与监督树协同:状态机作为 Worker 挂在 supervisor 下,事件管理器作为共享基础设施。完整的 OTP 行为协作可参见 https://plumephp.com/erlang-otp-framework/,而状态机的并发与容错根基来自 https://plumephp.com/erlang-concurrency-actors/。
总结:gen_statem 把状态机从「用 if 拼出来的模式」提升为「运行时原生语义」,配合延迟事件与状态超时,几乎可以优雅地表达任何协议与业务流程;gen_event 则为横切关注点(日志、指标、告警)提供了低耦合的广播机制。理解两者的设计差异与互补关系,是在 BEAM 上构建复杂事件驱动系统的关键一步。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。