Erlang 并发编程与 Actor 模型深度解析

深入剖析 Erlang 并发编程的核心机制,包括轻量级进程、消息传递、链接与监控、选择性接收以及 'let it crash' 容错哲学。

Erlang 的并发模型是这门语言最引以为豪的特性之一。与大多数编程语言依赖操作系统线程实现并发不同,Erlang 在虚拟机层面实现了轻量级进程(Lightweight Processes),每个进程仅占用几百字节到几千字节的内存,单个 Erlang 节点可以轻松支撑数十万个并发进程。这种设计不仅带来了极高的并发密度,更重要的是通过消息传递进程隔离从根本上消除了共享内存带来的数据竞争问题。本文将深入 Erlang 并发编程的全部核心概念,从进程创建到容错设计,展示如何构建可靠的并发系统。

一、Actor 模型概述

Erlang 的并发模型通常被称为 Actor 模型的典型实现。Actor 模型由 Carl Hewitt 于 1973 年提出,其核心思想是:系统中的所有计算实体都是独立的 Actor,Actor 之间通过异步消息传递进行通信,不存在共享状态

1.1 Actor 模型的三要素

要素说明Erlang 对应
Actor独立的计算单元pid() 进程标识符
消息传递异步发送消息! 操作符
行为切换接收消息后改变状态receive 表达式 + 尾递归
% Actor 的基本形态:一个循环接收消息的进程
-module(actor_demo).
-export([start/0, loop/1]).

start() ->
    spawn(?MODULE, loop, [initial_state]).

loop(State) ->
    receive
        {From, Message} ->
            NewState = handle_message(Message, State),
            From ! {self(), reply, process_message(Message)},
            loop(NewState)  % 尾递归切换状态
    end.

1.2 与共享内存并发模型的对比

维度共享内存模型(Java/Go)消息传递模型(Erlang/Elixir)
数据同步锁、信号量、原子变量消息传递天然隔离
竞态条件需要仔细加锁避免不存在(无共享状态)
死锁风险极低
扩展性受限于线程数单节点百万级进程
容错性一个线程崩溃影响全局单个进程崩溃不影响其他

二、轻量级进程

2.1 创建进程

Erlang 进程通过 spawn 系列函数创建,创建成本极低:

-module(process_demo).
-export([main/0, worker/0]).

worker() ->
    receive
        {compute, Task} ->
            Result = do_work(Task),
            io:format("Worker ~p computed: ~p~n", [self(), Result])
    end.

main() ->
    % spawn/3: 创建新进程,执行指定模块的函数
    Pid1 = spawn(?MODULE, worker, []),
    Pid2 = spawn(fun worker/0),  % 使用匿名函数引用
    
    % spawn/1: 直接执行匿名函数
    Pid3 = spawn(fun() ->
        io:format("Anonymous worker running~n")
    end),
    
    % 向进程发送消息
    Pid1 ! {compute, some_task},
    Pid2 ! {compute, another_task},
    
    % 获取当前进程 ID
    io:format("Main process: ~p~n", [self()]).

2.2 进程标识符(PID)

每个 Erlang 进程都有一个全局唯一的进程标识符(PID),格式如下:

% PID 格式:<NodeID.ProcessID.Serial>
% 例如:<0.84.0>
% - 第一个 0 表示本地节点
% - 84 是进程在进程表中的索引
% - 最后一个 0 是序列号,防止 PID 重用导致的混淆

% 进程信息查询
process_info(Pid) ->
    [{registered_name, []},
     {memory, 2680},           % 进程使用内存(字节)
     {message_queue_len, 0},   % 消息队列长度
     {status, waiting},        % 进程状态
     {current_function, {process_demo, worker, 0}}].

2.3 进程内存模型

Erlang 进程拥有独立的堆内存空间,进程间数据通过消息复制进行传递:

-module(memory_demo).
-export([show_memory/0]).

show_memory() ->
    % 创建一个 1000 元素的列表
    BigList = lists:seq(1, 1000),
    
    % 向进程发送消息时,数据会被复制
    Pid = spawn(fun receiver/0),
    Pid ! {data, BigList},  % BigList 被复制到新进程的堆中
    
    % 现代 Erlang(20+)引入 off-heap message passing
    % 大消息使用共享堆引用,减少复制开销
    ok.

receiver() ->
    receive
        {data, List} ->
            io:format("Received ~p elements~n", [length(List)])
    end.

三、消息传递机制

3.1 异步发送与邮箱

Erlang 的消息传递是异步的,发送方将消息放入接收方的邮箱后立即返回,无需等待:

% 发送消息语法:Pid ! Message
Pid ! hello.
Pid ! {request, self(), "get_user", 123}.

% 发送后是否需要等待回复取决于业务逻辑
Pid ! {request, self(), "do_something"},
    % 发送方继续执行其他操作...
    receive
        {response, Result} -> Result
    after 5000 -> timeout  % 5 秒超时
    end.

每个进程都有一个邮箱(Mailbox),按消息到达顺序存储消息。消息不会丢失,即使接收方正在处理其他消息。

3.2 选择性接收

receive 表达式支持模式匹配和超时,这是 Erlang 并发编程最优雅的设计之一:

-module(selective_receive).
-export([server/0]).

server() ->
    receive
        % 按优先级处理不同类型的消息
        {urgent, Msg} ->
            handle_urgent(Msg);
            
        {normal, Msg} when is_binary(Msg) ->
            handle_normal(Msg);
            
        {normal, Msg} ->
            io:format("Unexpected normal format: ~p~n", [Msg])
    after 
        10000 ->  % 10 秒超时
            io:format("No message received in 10 seconds~n")
    end.

选择性接收的实现是 Erlang 虚拟机的一个重要优化点。早期实现中,receive 每次扫描整个邮箱,复杂度为 O(N)。现代 Erlang 引入了引用计数邮箱机制,为每个消息存储一个「最后检查位置」指针,将每次接收的平均复杂度降低到 O(1)。

3.3 消息协议设计

在并发系统中,消息协议的设计直接影响系统的可维护性。推荐采用 {Tag, From, Request} 的通用格式:

% 请求-响应协议示例
-module(kv_store).
-export([start/0, loop/1]).
-define(TIMEOUT, 5000).

start() ->
    spawn(?MODULE, loop, [#{}]).

loop(State) ->
    receive
        {get, From, Key} ->
            Value = maps:get(Key, State, undefined),
            From ! {get_response, Key, Value},
            loop(State);
            
        {put, From, Key, Value} ->
            NewState = State#{Key => Value},
            From ! {put_response, Key, ok},
            loop(NewState);
            
        {delete, From, Key} ->
            NewState = maps:remove(Key, State),
            From ! {delete_response, Key, ok},
            loop(NewState);
            
        stop ->
            io:format("KV store shutting down~n")
    end.

% 客户端 API
get(Pid, Key) ->
    Pid ! {get, self(), Key},
    receive
        {get_response, Key, Value} -> Value
    after ?TIMEOUT -> {error, timeout}
    end.

put(Pid, Key, Value) ->
    Pid ! {put, self(), Key, Value},
    receive
        {put_response, Key, ok} -> ok
    after ?TIMEOUT -> {error, timeout}
    end.

四、链接与监控机制

4.1 进程链接(link/1)

Erlang 的容错机制建立在进程间的**链接(Link)**关系上。当两个进程链接后,任何一方异常退出都会导致另一方也退出:

-module(link_demo).
-export([spawn_linked_worker/0, worker/0]).

spawn_linked_worker() ->
    % spawn_link 创建进程并立即建立双向链接
    spawn_link(?MODULE, worker, []).

worker() ->
    receive
        crash ->
            exit(some_error)  % 触发异常退出
    end.

% 测试:
% > Pid = link_demo:spawn_linked_worker().
% > link(Pid).  % 手动建立链接
% > Pid ! crash.
% 此时,链接到 Pid 的所有进程都会收到退出信号

4.2 进程监控(monitor/2)

与链接不同,**监控(Monitor)**是单向的。监控者可以感知被监控者的状态变化,但不会因为被监控者退出而自动退出:

-module(monitor_demo).
-export([watch_worker/0]).

watch_worker() ->
    Pid = spawn(fun worker/0),
    
    % 建立单向监控
    Ref = monitor(process, Pid),
    
    % 等待结果或监控触发
    receive
        {'DOWN', Ref, process, Pid, Reason} ->
            io:format("Worker ~p exited with reason: ~p~n", [Pid, Reason]),
            % 可以在这里重启 worker
            watch_worker();
            
        {result, Data} ->
            io:format("Got result: ~p~n", [Data])
    end.

worker() ->
    % 模拟工作
    timer:sleep(2000),
    exit(worker_crashed).

4.3 退出信号捕获

进程可以通过 process_flag(trap_exit, true) 捕获链接进程的退出信号,将原本致命的退出信号转换为普通消息:

-module(supervisor_demo).
-export([start_worker/1, worker_supervisor/1]).

worker_supervisor(WorkerModule) ->
    process_flag(trap_exit, true),  % 捕获退出信号
    start_child(WorkerModule).

start_child(WorkerModule) ->
    Pid = spawn_link(WorkerModule, start, []),
    io:format("Started worker ~p~n", [Pid]),
    wait_for_exit(Pid, WorkerModule).

wait_for_exit(Pid, WorkerModule) ->
    receive
        {'EXIT', Pid, normal} ->
            io:format("Worker ~p exited normally~n", [Pid]);
            
        {'EXIT', Pid, Reason} ->
            io:format("Worker ~p crashed: ~p. Restarting...~n", [Pid, Reason]),
            timer:sleep(1000),  % 延迟重启避免 CPU 飙升
            start_child(WorkerModule)
    end.

这是 OTP Supervisor 行为的核心思想:如果一个进程崩溃了,就重启它。这种「let it crash」的哲学与防御式编程形成鲜明对比,它承认错误不可避免,与其在每个函数中处理所有边界情况,不如将错误隔离到单个进程中并快速重启。

五、并发模式

5.1 客户端-服务器模式

-module(gen_server_pattern).
-export([start/1, rpc/2, loop/1]).

start(InitialState) ->
    spawn(?MODULE, loop, [InitialState]).

% 同步远程过程调用
rpc(Server, Request) ->
    Server ! {call, self(), Request},
    receive
        {Server, Response} -> Response
    end.

loop(State) ->
    receive
        {call, From, Request} ->
            {Response, NewState} = handle_call(Request, State),
            From ! {self(), Response},
            loop(NewState)
    end.

handle_call({get, Key}, State) ->
    {maps:get(Key, State, undefined), State};
handle_call({put, Key, Value}, State) ->
    {ok, State#{Key => Value}}.

5.2 工作者池(Worker Pool)

-module(worker_pool).
-export([start/2, dispatch/2]).

start(WorkerModule, Count) ->
    Workers = [spawn(WorkerModule, start, []) || _ <- lists:seq(1, Count)],
    spawn(fun() -> scheduler(Workers, queue:new()) end).

% 轮询调度 + 请求队列缓冲
scheduler(Workers, Queue) ->
    receive
        {dispatch, Task} when length(Workers) > 0 ->
            [Worker | Rest] = Workers,
            Worker ! {work, Task},
            scheduler(Rest, Queue);
            
        {dispatch, Task} ->
            % 所有worker都忙,放入队列
            scheduler(Workers, queue:in(Task, Queue));
            
        {done, Worker} ->
            case queue:out(Queue) of
                {empty, _} ->
                    scheduler([Worker | Workers], Queue);
                {{value, Task}, NewQueue} ->
                    Worker ! {work, Task},
                    scheduler(Workers, NewQueue)
            end
    end.

dispatch(Pool, Task) ->
    Pool ! {dispatch, Task}.

5.3 Pub/Sub 发布订阅

-module(pubsub).
-export([start/0, subscribe/2, publish/2]).

start() ->
    spawn(fun() -> broker(#{}) end).

broker(Subscribers) ->
    receive
        {subscribe, Pid, Topic} ->
            Current = maps:get(Topic, Subscribers, []),
            monitor(process, Pid),  % 自动清理退订
            broker(Subscribers#{Topic => [Pid | Current]});
            
        {publish, Topic, Message} ->
            Pids = maps:get(Topic, Subscribers, []),
            [P ! {published, Topic, Message} || P <- Pids],
            broker(Subscribers);
            
        {'DOWN', _Ref, process, Pid, _Reason} ->
            % 清理退出的进程
            NewSubs = maps:map(fun(_, Pids) -> lists:delete(Pid, Pids) end, Subscribers),
            broker(NewSubs)
    end.

subscribe(Broker, Topic) ->
    Broker ! {subscribe, self(), Topic}.

publish(Broker, Topic, Message) ->
    Broker ! {publish, Topic, Message}.

六、性能与调优

6.1 进程数与系统限制

% 查看当前进程数
> erlang:system_info(process_count).
78

% 查看最大进程限制
> erlang:system_info(process_limit).
262144

% 调整最大进程数(启动时设置 +P 参数)
% erl +P 1000000

6.2 消息队列监控

消息队列堆积是 Erlang 系统最常见的性能问题。应监控以下指标:

-module(monitor_utils).
-export([check_mailbox/1, find_large_mailboxes/0]).

check_mailbox(Pid) ->
    case process_info(Pid, message_queue_len) of
        {message_queue_len, Len} when Len > 1000 ->
            io:format("WARNING: Pid ~p has ~p messages queued~n", [Pid, Len]);
        {message_queue_len, Len} ->
            io:format("Pid ~p mailbox: ~p~n", [Pid, Len])
    end.

find_large_mailboxes() ->
    Pids = processes(),
    [
        {Pid, Len} 
        || Pid <- Pids,
           {message_queue_len, Len} <- [process_info(Pid, message_queue_len)],
           Len > 1000
    ].

6.3 调度器与核亲和性

Erlang VM 默认每个 CPU 核心运行一个调度器线程,通过 +S 参数控制:

# 使用 8 个调度器线程
erl +S 8:8

# 绑定调度器到 CPU 核心,减少缓存失效
erl +S 8:8 +sbt db

七、总结

Erlang 的并发模型不仅是 Actor 模型的教科书级实现,更通过「let it crash」的容错哲学将并发编程的错误处理提升到了新的层次。轻量级进程让并发成为一种编程常态而非特殊能力,消息传递的不可变性从根本上消除了数据竞争,链接和监控机制让进程树成为自恢复的有机体。

理解 Erlang 并发编程的关键不在于记住 API,而在于建立一种思维方式:系统由大量独立进程组成,进程通过消息传递协作,进程必然会崩溃,重要的是让系统在崩溃后自行愈合。这种思维方式不仅适用于 Erlang,也深刻影响了 Akka、Elixir、Go 的并发设计以及 Kubernetes 等云原生系统的编排哲学。

继续阅读

探索更多技术文章

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

全部文章 返回首页

「erlang」更多文章

  1. Erlang/OTP 生产案例与性能调优
  2. Erlang 分布式编程:节点互联与集群部署
  3. Elixir 入门与 Phoenix Web 框架实战