Erlang 分布式编程:节点互联与集群部署

深入 Erlang 分布式编程的能力边界,探讨节点发现、RPC 通信、分布式 Mnesia 数据库、CAP 权衡与生产集群的部署策略。

Erlang 在设计之初就考虑了分布式场景,其虚拟机内置了节点互联和透明 RPC 能力,使得编写能够跨物理机运行的并发程序几乎与单机编程一样简单。从 WhatsApp 的数十万台服务器集群到 RabbitMQ 的多节点镜像队列,Erlang 分布式模型经受了最严苛的生产环境检验。本文将深入 Erlang 分布式编程的核心机制,包括节点通信模型、分布式 Mnesia 数据库的设计权衡,以及大规模集群的部署策略。

一、Erlang 分布式模型

1.1 基础知识

Erlang 分布式系统的基本单元是节点(Node),一个节点就是一个运行中的 Erlang 运行时实例。节点之间通过 TCP 连接建立通信:

% 启动一个带名称的节点
$ erl -sname mynode@localhost

% 远程启动
$ erl -sname node_a -setcookie mysecret
$ erl -sname node_b -setcookie mysecret
% 在 node_a 上连接 node_b
(node_a@localhost)1> net_adm:ping('node_b@localhost').
pong  % 连接成功(pang 表示失败)

% 查看已连接的节点
(node_a@localhost)2> nodes().
['node_b@localhost']

% 在远程节点上执行代码
(node_a@localhost)3> rpc:call('node_b@localhost', erlang, node, []).  
'node_b@localhost'

Erlang 节点使用魔术 cookie 进行认证,同一集群的所有节点必须共享相同的 cookie:

% 设置 cookie
erlang:set_cookie(node(), mysecretcookie).

% 查看当前 cookie
erlang:get_cookie().

⚠️ 安全警告:Cookie 是明文传输的,且节点间通信默认不加密。在生产环境中应:

  • 使用 -proto_dist inet_tls 启用 TLS
  • 将 cookie 文件权限设为 400(只读)
  • 通过防火墙限制节点间通信端口范围

1.3 分布式进程通信

分布式 Erlang 最优雅的特性是位置透明性——向进程发送消息时,不需要知道进程运行在本地还是远程节点:

-module(dist_demo).
-export([start_server/0, client_task/1]).

start_server() ->
    spawn(fun server_loop/0).

server_loop() ->
    receive
        {From, {compute, X}} ->
            Result = X * X,
            From ! {self(), {result, Result}},
            server_loop();
        stop ->
            io:format("Server shutting down~n")
    end.

client_task(ServerPid) ->
    ServerPid ! {self(), {compute, 42}},
    receive
        {ServerPid, {result, Result}} ->
            io:format("Got result from ~p: ~p~n", [ServerPid, Result])
    after 5000 ->
        io:format("Timeout~n")
    end.

% 使用示例(node_a):
% Pid = dist_demo:start_server().
% REGISTER pid with global or local name

% node_b 可以透明地发送消息:
% global:send(my_server, {self(), {compute, 10}}).

Pid 本身携带了节点信息(<0.84.0> 中的 0 表示本地),因此消息可以自动路由到正确的节点。

二、全局命名与 RPC

2.1 全局注册表

在分布式系统中,使用 PID 直接引用进程不便于管理。Erlang 提供了 global 模块实现跨节点命名:

% 全局注册(整个集群可见)
global:register_name(my_server, Pid).

% 全局查找
global:whereis_name(my_server).

% 全局发送
global:send(my_server, Message).

% 进程退出时自动注销
global:register_name(my_server, Pid, fun global_name_resolution/3).

2.2 RPC 框架

% 远程过程调用
rpc:call(Node, Module, Function, Args).
rpc:cast(Node, Module, Function, Args).  % 异步调用

% 多节点并行调用
rpc:multicall(Nodes, Module, Function, Args).

% 示例:在所有节点上获取负载信息
LoadInfo = rpc:multicall(nodes(), os, cpu_topology, []).

2.3 分布式任务分发

-module(task_distributor).
-export([distribute/2, worker/0]).

% 将任务分发给集群中的工作进程
distribute(Tasks, WorkersPerNode) ->
    Nodes = [node() | nodes()],
    
    % 在每个节点上启动工作进程
    Workers = lists:flatten([
        [rpc:call(N, ?MODULE, spawn_worker, []) || _ <- lists:seq(1, WorkersPerNode)]
        || N <- Nodes
    ]),
    
    % 使用 round-robin 分配任务
    assign_tasks(Tasks, Workers, round_robin(Workers)).

spawn_worker() ->
    spawn(?MODULE, worker, []).

worker() ->
    receive
        {execute, From, Task} ->
            Result = process_task(Task),
            From ! {result, self(), Result},
            worker()
    end.

assign_tasks([], _, _) -> ok;
assign_tasks([Task | Rest], Workers, [Worker | RemainingWorkers]) ->
    Worker ! {execute, self(), Task},
    NextWorkers = case RemainingWorkers of
        [] -> Workers;
        _ -> RemainingWorkers
    end,
    assign_tasks(Rest, Workers, NextWorkers).

round_robin(List) -> lists:duplicate(100, List).  % 简单循环

三、分布式 Mnesia 数据库

Mnesia 是 Erlang 内置的分布式数据库,它将 ETS/DETS 的高性能与事务支持、分布式复制能力结合在一起。

3.1 Schema 与表定义

-module(mnesia_demo).
-export([init/0, create_user/2, get_user/1]).

-record(user, {id, name, email, created_at}).

init() ->
    % 创建 schema(首次在集群上运行)
    mnesia:create_schema([node() | nodes()]),
    mnesia:start(),
    
    % 创建分布式表
    mnesia:create_table(user, [
        {attributes, record_info(fields, user)},
        {disc_copies, [node()]},        % 本地磁盘副本
        {ram_copies, nodes()},           % 所有节点的内存副本
        {type, set},
        {index, [email]}
    ]).

create_user(Id, Name, Email) ->
    F = fun() ->
        mnesia:write(#user{
            id = Id,
            name = Name,
            email = Email,
            created_at = calendar:local_time()
        })
    end,
    mnesia:transaction(F).

get_user(Id) ->
    F = fun() -> mnesia:read({user, Id}) end,
    case mnesia:transaction(F) of
        {atomic, [User]} -> {ok, User};
        {atomic, []} -> {error, not_found}
    end.

3.2 复制模式选择

配置说明一致性性能可用性
{ram_copies, Nodes}内存副本事务保证最高节点宕机丢数据
{disc_copies, Nodes}磁盘副本事务保证中等节点宕机不丢
{disc_only_copies, Nodes}仅磁盘事务保证较低节点宕机不丢
{fragments, N}哈希分片最终一致分布式需要 Quorum

3.3 CAP 权衡

Mnesia 的默认行为牺牲了一定程度的可用性以保证强一致性:

  • 写入:需要所有副本节点确认(同步复制)
  • 读取:从本地副本读取(可能读到旧数据)
  • 网络分区:节点被隔离后,被隔离节点上的写操作会被阻塞

可以通过配置调整 CAP 偏向:

% 主节点配置 - 分区时主节点继续服务
mnesia:set_master_nodes(Tab, [PrimaryNode]).

% 脏读(绕过事务,类似 NoSQL 的 Eventual Consistency)
mnesia:dirty_read({user, Id}).
mnesia:dirty_write(User).

四、集群管理与高可用

4.1 节点动态加入

-module(cluster_manager).
-export([join/1, leave/1, discover/0]).

% 加入集群
join(SeedNode) ->
    case net_adm:ping(SeedNode) of
        pong ->
            % 集群中发现
            AllNodes = rpc:call(SeedNode, erlang, nodes, []),
            [net_adm:ping(N) || N <- AllNodes],
            
            % 复制 Mnesia 表
            mnesia:change_config(extra_db_nodes, AllNodes),
            mnesia:add_table_copy(user, node(), disc_copies),
            
            {ok, nodes()};
        pang ->
            {error, seed_unreachable}
    end.

% 优雅离开
leave(Node) ->
    % 迁移数据副本
    Tabs = mnesia:system_info(tables),
    [mnesia:del_table_copy(T, Node) || T <- Tabs],
    
    % 断开连接
    erlang:disconnect_node(Node),
    ok.

% 自动发现(基于 UDP 多播或 Kubernetes DNS)
discover() ->
    % Kubernetes 环境下通过 Headless Service 发现
    {ok, Hostname} = inet:gethostname(),
    discover_k8s(Hostname).

discover_k8s(BaseHostname) ->
    % 查询 DNS 获取同服务的所有 Pod IP
    {ok, Addrs} = inet_res:getbyname(BaseHostname, a),
    [try_connect(A) || A <- Addrs].

4.2 故障检测与自动切换

-module(node_monitor).
-export([start/0, monitor_nodes/0]).

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

monitor_nodes() ->
    net_kernel:monitor_nodes(true),
    loop().

loop() ->
    receive
        {nodeup, Node} ->
            io:format("Node ~p joined the cluster~n", [Node]),
            % 触发数据再平衡
            rebalance_data(Node),
            loop();
            
        {nodedown, Node} ->
            io:format("Node ~p left the cluster~n", [Node]),
            % 检查 Mnesia 副本是否仍然满足 N+1
            check_replica_safety(),
            loop();
            
        {heartbeat, From} ->
            From ! {heartbeat_ack, node()},
            loop()
    end.

rebalance_data(NewNode) ->
    % 将部分副本迁移到新节点以平衡负载
    Tables = mnesia:system_info(tables),
    lists:foreach(fun(Tab) ->
        Copies = mnesia:table_info(Tab, ram_copies),
        case lists:member(node(), Copies) andalso not lists:member(NewNode, Copies) of
            true ->
                mnesia:add_table_copy(Tab, NewNode, ram_copies);
            false ->
                ok
        end
    end, Tables).

五、生产部署策略

5.1 网络拓扑设计

┌─────────────────────────────────────────────────┐
│                    Load Balancer                  │
│                (HAProxy / Nginx)                │
└──────────────┬────────────────┬─────────────────┘
               │                │
      ┌────────▼────────┐      ▼
      │   Erlang Node 1  │  ┌───▼────┐
      │   (API Server)   │  │ Node 2  │
      │                  │  │ (API)   │
      │  ┌────────┐     │  │         │
      │  │ Cowboy │     │  │ ┌────┐ │
      │  │ WebSrv │     │  │ │Cowboy│ │
      │  └────────┘     │  │ └────┘ │
      │                  │  └────────┘
      │  ┌────────┐     │
      │  │ Mnesia │◄────┼────► 数据同步
      │  │ Primary│     │
      │  └────────┘     │
      └─────────────────┘
               │
      ┌────────▼────────┐
      │   Erlang Node 3  │
      │  (Background Job)│
      │                  │
      │  ┌──────────┐   │
      │  │ Scheduler│   │
      │  │ Worker   │   │
      │  └──────────┘   │
      └─────────────────┘

5.2 配置最佳实践

% vm.args 生产配置

% 节点名称
-name app@192.168.1.10

% Cookie 文件路径(安全)
-setcookie ${ERL_COOKIE}

% 集群通信端口范围
-kernel inet_dist_listen_min 9100 inet_dist_listen_max 9150

% 启用分布式 TLS
-proto_dist inet_tls
-ssl_dist_opt server_certfile /etc/ssl/erl_server.pem
-ssl_dist_opt server_keyfile /etc/ssl/erl_server.key
-ssl_dist_opt server_cacertfile /etc/ssl/ca.pem

% 系统限制
+K true              % 启用内核轮询
+A 64                % 异步线程池大小
+S 8:8               % 8 个调度器
+P 1000000           % 最大进程数

5.3 监控与观测

% 关键指标收集
-module(metrics).
-export([cluster_stats/0]).

cluster_stats() ->
    [
        {node, node()},
        {connected_nodes, length(nodes())},
        {process_count, erlang:system_info(process_count)},
        {memory_total, erlang:memory(total)},
        {memory_processes, erlang:memory(processes)},
        {message_queue, total_message_queue()},
        {reductions, element(1, erlang:statistics(reductions))},
        {ets_tables, length(ets:all())},
        {mnesia_tables, length(mnesia:system_info(tables))}
    ].

total_message_queue() ->
    lists:sum([
        element(2, process_info(Pid, message_queue_len))
        || Pid <- processes(),
           is_process_alive(Pid)
    ]).

六、总结

Erlang 分布式编程的核心优势在于透明性一致性。开发者不需要学习复杂的 RPC 框架或消息队列 API,进程间通信在单机和多机场景下使用完全相同的语义。Mnesia 提供了分布式事务数据库能力,虽然不适合海量数据场景,但对于配置管理、会话状态和元数据存储恰到好处。

在生产环境中使用 Erlang 分布式功能时,需要特别注意:网络分区场景下的行为设计、Cookie 安全管理、TLS 加密配置,以及 Mnesia 的副本策略选择。Erlang/OTP 的分布式抽象能力仍然是业界 unmatched 的,它为理解分布式系统的本质提供了一个极简而优雅的模型。

继续阅读

探索更多技术文章

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

全部文章 返回首页

「erlang」更多文章

  1. Erlang/OTP 生产案例与性能调优
  2. Elixir 入门与 Phoenix Web 框架实战
  3. Erlang OTP 框架:构建工业级并发应用