Erlang gen_statem 状态机与 gen_event 事件管理深度解析

深入 OTP 的 gen_statem 状态机行为,掌握 state functions 与 handle_event 两种回调模式、延迟事件与状态超时机制,以及 gen_event 事件管理的设计与应用对比。

状态机是表达复杂业务流程最自然的抽象:连接建立与断开、支付状态流转、协议握手、任务调度,都能用一组有限状态和状态间转移精确描述。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调用参数,可同步回复
castgen_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_functionshandle_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_servergen_statemgen_event
核心模型客户端-服务器有限状态机发布-订阅
状态维度单一状态 StateStateName + Data每个 Handler 独立状态
同步调用handle_call/3handle_event({call, _}, ...)handle_call/2
异步消息handle_cast/2handle_event(cast, ...)handle_event/2
定时器手动 send_aftertimeout / 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 上构建复杂事件驱动系统的关键一步。

继续阅读

探索更多技术文章

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

全部文章 返回首页

「erlang」更多文章

  1. 自定义 OTP Behaviour 实战:Callback 规范与行为封装
  2. Phoenix Channels 实时通信实战:WebSocket 与 PubSub 深入
  3. Mix 工具链与 Elixir 工程化实战