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'
1.2 Cookie 安全模型
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 的,它为理解分布式系统的本质提供了一个极简而优雅的模型。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。