Erlang 的分布式是「透明」的:Node.connect/1 之后,跨节点发消息与本地发消息用同一套语法,spawn 一个远端进程也只是一行代码。这种透明性让分布式编程的门槛降到极低,但它不提供任何一致性保证——两个节点之间的连接断开时,BEAM 不会告诉你「现在集群分裂了」,它只是让 Node.list/0 少一个节点,让 :global 里的名字各算各的。
于是「用起来简单」和「用对很难」之间的落差,成了 BEAM 分布式系统最主要的故障来源。本文从拓扑与连接成本讲起,覆盖 libcluster 的集群组建、一致性模型的选择、CRDT 的收敛语义,最后重点讨论网络分区(network partition)与脑裂(split-brain)的治理。
分布式 Erlang 的拓扑与连接成本
启动一个具名节点需要指定节点名与 cookie:
iex --name app@10.0.1.5 --cookie secret_cookie
Node.self() # :"app@10.0.1.5"
Node.list() # 当前已连接的节点列表
Node.connect(:"app@10.0.1.6")
Node.ping(:"app@10.0.1.6") # :pong | :pang
cookie 是节点间的共享密钥,必须一致且保密。节点名有两种形式:长名(--name,用 name@host.fqdn)与短名(--sname,用 name@short-hostname)。容器环境下主机名常变,用长名 + IP 更稳定。
默认情况下,Erlang 节点之间是全互联(full mesh):任何两个节点都建立一条直接连接。连接数是 N×(N−1)/2,即 O(N²):
| 节点数 | 连接数 | 每节点连接数 |
|---|---|---|
| 5 | 10 | 4 |
| 20 | 190 | 19 |
| 50 | 1225 | 49 |
| 100 | 4950 | 99 |
每个连接背后是一个独立的 TCP 会话与一组监控进程,心跳、:global 名称同步、monitor_nodes 通知都会随之放大。50 个节点是大多数 BEAM 集群的实际上限,超过之后连接维护本身就消耗可观资源。
要突破这个上限,可以用 -connect_all false 关闭全互联,改为只连必要的对端:
%% vm.args
-connect_all false
代价是 :global 与依赖全互联的机制(如 :global.register_name/2 的自动同步)不再可靠,需要自己维护拓扑视图。只有在确实需要大规模集群、且能接受手动拓扑管理时才这么做;多数业务在 3~10 个节点之间,全互联完全够用。
节点状态变化通过 monitor_nodes 监听:
:net_kernel.monitor_nodes(true, node_type: :all)
def handle_info({:nodeup, node, _info}, state) do
Logger.info("node up: #{node}")
{:noreply, state}
end
def handle_info({:nodedown, node, info}, state) do
Logger.warning("node down: #{node}, reason: #{inspect(info[:reason])}")
{:noreply, state}
end
info[:reason] 是判断「宕机」还是「网络分区」的重要线索::noconnection 通常意味着网络问题,:shutdown 是正常退出。这些回调是集群感知能力的入口,与 分布式编程基础
中的进程注册、spawn 远端进程配合使用。
libcluster:集群的自动组建
手动 Node.connect/1 在生产环境不可行——节点重启后 IP 会变,扩缩容更是无从下手。libcluster 把「发现对端并连接」抽象成可插拔的策略(strategy)。
# mix.exs
{:libcluster, "~> 3.4"}
# config/runtime.exs
config :libcluster,
topologies: [
gossip: [
strategy: Cluster.Strategy.Gossip,
config: [
port: 45892,
if_addr: "0.0.0.0",
multicast_addr: "230.1.1.251",
multicast_ttl: 1
]
],
k8s: [
strategy: Cluster.Strategy.Kubernetes.DNS,
config: [
service: "myapp-headless",
application_name: "myapp",
polling_interval: 5_000
]
]
]
把 Cluster.Supervisor 挂进应用监督树:
children = [
{Cluster.Supervisor, [topologies, [name: MyApp.ClusterSupervisor]]},
MyApp.Repo,
MyAppWeb.Endpoint
]
常用策略的对比:
| 策略 | 发现方式 | 适用环境 | 注意 |
|---|---|---|---|
Cluster.Strategy.Gossip | UDP 组播互探 | 同网段、开发 | 跨网段/云上组播常被禁 |
Cluster.Strategy.Epmd | 本地 epmd 注册表 | 单机多节点 | 仅限同主机 |
Cluster.Strategy.ErlangHosts | .hosts.erlang 文件 | 静态集群 | 节点列表写死 |
Cluster.Strategy.Kubernetes.DNS | headless Service DNS | K8s | 依赖 DNS 轮询 |
Cluster.Strategy.Kubernetes | K8s API 查 Pod | K8s | 需要 RBAC 权限 |
Gossip 是开发与内网环境的默认选择,它靠 UDP 组播互相打招呼,零配置。但云环境的 VPC 默认关闭组播,跨子网也不通,此时必须换成基于 DNS 或 API 的策略。
Kubernetes 环境用 headless Service 的 DNS 记录做发现最省事:headless Service(clusterIP: None)会为每个 Pod 返回一条 A 记录,libcluster 轮询这个记录即可得到全部节点 IP。polling_interval 决定了扩缩容后集群收敛的速度——设得太小会增加 DNS 压力,太大则新节点长时间不被感知,5 秒是个平衡点。
libcluster 只负责建立连接,不负责一致性。它连上之后,跨节点的状态同步仍需要应用层自己解决。这是理解后续所有内容的前提。
一致性模型与 CAP 取舍
CAP 定理说:网络分区(Partition)发生时,一致性(Consistency)与可用性(Availability)只能选一个。BEAM 的默认行为是选可用性——分区后两边各自继续服务,接受数据分叉。要选一致性,必须显式引入多数派机制。
实际系统的一致性是一个谱系,不是二选一:
| 模型 | 保证 | BEAM 侧实现 |
|---|---|---|
| 线性一致(Linearizable) | 所有操作看起来在单一时间点生效 | :ra(Raft)、单点写 |
| 顺序一致(Sequential) | 所有节点看到相同操作顺序 | 中心化序列号服务 |
| 因果一致(Causal) | 有因果关系的操作保序 | 版本向量(Version Vector) |
| 最终一致(Eventual) | 停止写入后最终收敛 | CRDT、反熵(Anti-entropy) |
先问业务需要哪一档,再选技术。订单金额、库存扣减需要线性一致,代价是分区时不可写;点赞数、在线状态、浏览计数只需要最终一致,可以换来分区期间两边都能写。
BEAM 里实现线性一致的代价不低::global 的名称注册在分区时会各自为政,Mnesia 需要显式多数派配置,真正的强一致要靠 :ra(RabbitMQ 的 quorum queue 与 Khepri 都基于它)。多数场景的正确做法是把强一致需求外移到数据库(Postgres 的事务、Redis 的原子命令),BEAM 集群只处理无状态或最终一致的部分。
CRDT:可收敛的共享状态
CRDT(Conflict-free Replicated Data Type,无冲突复制数据类型) 是最终一致的数学基础:它保证只要所有副本最终收到全部更新,无论顺序如何、无论是否重复,最终状态都相同。这消除了分布式系统中最麻烦的部分——冲突解决。
核心性质有两条:
- 交换律(commutativity):
merge(a, b) == merge(b, a) - 幂等性(idempotency):
merge(a, a) == a
幂等性尤其重要,它让消息重传变得安全——网络抖动导致重复投递不会破坏状态。
常用的 CRDT 类型:
| 类型 | 语义 | 典型用途 |
|---|---|---|
| G-Counter | 只增计数器 | 页面浏览量 |
| PN-Counter | 可增可减计数器 | 库存余量(需谨慎) |
| G-Set | 只增集合 | 已处理任务 ID |
| OR-Set | 增删集合,删优先于并发增 | 在线用户列表 |
| LWW-Register | 按时间戳取最后写入 | 用户昵称、配置项 |
| LWW-Element-Set | 带时间戳的元素集合 | 标签集合 |
Erlang 生态的 CRDT 库是 :delta_crdt(delta-state CRDT,只传输增量而非全量状态,网络开销从 O(状态大小) 降到 O(变更大小)):
# mix.exs
{:delta_crdt, "~> 0.6"}
alias DeltaCrdt
{:ok, crdt} = DeltaCrdt.start_link(DeltaCrdt.AWLWWMap)
# 本地更新
DeltaCrdt.put(crdt, :user_1_status, :online, :node_a)
# 读取本地副本(无需网络往返)
DeltaCrdt.read(crdt)
# 与另一个节点上的副本互联
DeltaCrdt.set_neighbours(crdt, [:"app@10.0.1.6"])
DeltaCrdt 的传播模型是邻居间同步:每个副本与若干邻居交换增量,增量像流言一样扩散到全网。set_neighbours/2 定义的邻居关系决定了收敛速度——全互联收敛最快但连接最多,环形拓扑连接少但延迟高。
DeltaCrdt.AWLWWMap 是「Add-Wins Last-Write-Wins Map」,语义是:并发的删除与新增,以新增为准;同一个 key 的并发更新,以时间戳晚的为准。时间戳依赖物理时钟,节点间时钟漂移会直接影响语义正确性,必须配 NTP。
Phoenix 生态的 Phoenix.Tracker 就构建在 DeltaCrdt 之上,用于跨节点跟踪进程(如「某个用户连在哪个节点」)。Horde 则用 CRDT 做分布式 Registry 与 DynamicSupervisor 的状态复制。
CRDT 的适用边界要清楚:它适合「可交换的更新」,不适合「需要读-改-写原子性」的操作。用 PN-Counter 实现库存扣减就是经典反例——它只能保证计数收敛,无法阻止超卖,因为「检查余量」与「扣减」不是原子操作。这类场景需要线性一致,应交给数据库。
脑裂治理
脑裂(split-brain) 指网络分区导致集群分裂成多个互不可见的子集群,每个子集群都认为自己是完整的。危害在于:同一份数据被两边分别修改,分区恢复后无法自动合并;或者同一个单例服务在两个分区里各起一份,造成重复处理。
检测
第一层检测靠 monitor_nodes 的 :nodedown。但要注意,:nodedown 不一定意味着对方宕机——网络抖动、GC 停顿、TCP 重传超时都会触发。真正的分区特征是「多个节点同时从视图里消失」:
def handle_info({:nodedown, node, info}, state) do
remaining = Node.list()
expected = state.expected_nodes
if length(remaining) < div(expected, 2) + 1 do
Logger.error("lost quorum: #{length(remaining)}/#{expected} nodes remaining")
# 触发自保:停止对外服务或降级为只读
end
{:noreply, %{state | down: [node | state.down]}}
end
多数派(Quorum)
多数派是脑裂治理最有效的原则:只有拥有超过半数节点的分区才能继续提供写服务,少数派分区主动降级为只读或停止服务。这样任意时刻最多只有一个分区能写,从根上避免了分叉。
判定逻辑简单直接:alive_nodes > total_nodes / 2。前提是 total_nodes 必须是一个静态配置的期望值,不能从 Node.list() 动态推导——分区后每个分区看到的节点数都不同,动态推导必然两边都认为自己是多数派。
OTP 的内建防护
OTP 25 引入了 prevent_overlapping_partitions 选项(OTP 26 起默认为 true),它在节点重连时检测「重叠分区」并主动断开旧连接,避免集群长期处于两个不相交的视图:
%% vm.args
-kernel prevent_overlapping_partitions true
这个机制解决的是「分区恢复后视图不一致」,不解决分区期间的写冲突。两者是互补的:前者保证最终视图收敛,后者需要应用层的多数派或 CRDT 来保证数据收敛。
Mnesia 的分区行为
Mnesia 在分区期间的行为需要特别理解。默认配置下,分区的两边都认为自己持有完整数据,都能读写;恢复连接后 Mnesia 检测到不一致,会报 inconsistent_database 并拒绝启动表。
mnesia:system_info(db_nodes). %% 全部数据库节点
mnesia:system_info(running_db_nodes). %% 当前可见的节点
%% 分区后人工恢复:以某个节点为准
mnesia:set_master_nodes([node1, node2]).
Mnesia 的多数派可以通过在启动时设置 extra_db_nodes 与 master_nodes 来约束,但配置复杂且容易出错。如果业务需要跨节点强一致的表,优先考虑 :ra 或外部数据库,Mnesia 更适合单节点或对一致性要求不高的场景,细节见 Mnesia 分布式数据
。
分区恢复流程
分区恢复不是「连上就好」,需要一个明确流程:
- 确认分区已消除:
Node.list()稳定包含全部期望节点,且持续数秒无抖动。 - 检查数据分叉:对 CRDT 状态做一次全量比对(
DeltaCrdt.read/1对比哈希),对数据库检查冲突记录。 - 重建单例:分区期间可能有两份单例进程,恢复后必须保证只剩一份。用
:global的名称注册会在重连后自动选主,但需要确认旧进程已被终止。 - 恢复写入:多数派降级的节点重新加入可写集合。
- 记录事件:把分区起止时间、影响范围写入日志,作为后续调参依据。
心跳调优:分区检测的速度与误判
BEAM 靠**心跳(heartbeat)**判断对端是否存活。每个节点定期向已连接节点发送 tick,若在 net_ticktime 窗口内没有收到任何 tick,就判定对端不可达并触发 nodedown。
%% vm.args
-kernel net_ticktime 60
net_ticktime 的默认值是 60 秒,含义是「大约 60 秒无心跳即判定断开」,实际触发时间在 net_ticktime 的 50%~100% 之间。这个参数是检测速度与误判率的权衡:
| 取值 | 分区检测延迟 | 误判风险 | 适用 |
|---|---|---|---|
| 15 秒 | 快(7~15 秒) | 高,GC 停顿或网络抖动即触发 | 同机架、低延迟内网 |
| 60 秒(默认) | 中(30~60 秒) | 低 | 通用 |
| 120 秒 | 慢(60~120 秒) | 极低 | 跨地域、网络不稳 |
调小 net_ticktime 能让脑裂更快被发现,代价是长 GC 停顿或短暂网络抖动也会触发 nodedown。BEAM 的 GC 是分代的,大堆进程的一次 full sweep 可能停顿数百毫秒,正常情况下远小于 15 秒;但如果进程持有巨量数据(如大 ETS 表的复制),停顿会显著变长。
经验规则是:net_ticktime 至少是 P99 最大 GC 停顿的 3 倍。用 recon 或 :erlang.statistics(:garbage_collection) 观察 GC 行为,确认最坏停顿量级后再定这个值。
除了心跳超时,还有两个相关参数影响分区期间的资源占用:
%% 发送队列水位,超过则断连而非无限缓冲
-kernel dist_buf_busy_limit 16MB
%% 连接建立超时
-kernel net_setuptime 15
dist_buf_busy_limit 尤其重要:分区期间如果一端持续向不可达节点发消息,发送队列会无限增长直到耗尽内存。设置上限后,超过水位会直接断开连接并丢弃消息——对「宁可丢消息也不能 OOM」的系统,这是必要的保护。
反熵与数据修复
CRDT 的收敛保证是「所有更新最终都被送达」。但网络分区、节点重启、消息丢弃都可能让某个副本长期缺失更新。反熵(anti-entropy) 是修复这类偏差的机制:定期在副本间比对状态摘要,发现差异则补齐。
最简单的实现是周期性全量比对:
defmodule AntiEntropy do
use GenServer
@interval :timer.minutes(5)
def handle_info(:sync, state) do
Enum.each(Node.list(), fn node ->
remote_digest = :rpc.call(node, DeltaCrdt, :read, [state.crdt])
local_digest = DeltaCrdt.read(state.crdt)
if remote_digest != local_digest do
# 触发一次全量状态合并
DeltaCrdt.sync(state.crdt, node)
end
end)
schedule()
{:ok, state}
end
end
全量比对的开销与状态大小成正比。状态较大时改用摘要(digest):对状态做哈希后只传输哈希,只有不一致才拉全量。Merkle 树是把摘要做到极致的形式——它按 key 空间分桶,比对时能定位到具体哪个桶不同,只同步差异桶,代价从 O(状态) 降到 O(差异)。
反熵的频率是另一个权衡:频率高则网络与 CPU 开销大,频率低则偏差存在时间长。对最终一致的数据,5 分钟一次是合理起点;对「偏差可容忍度低」的场景,应改用更主动的推送式同步(每次更新立即广播增量)。
分片与数据归属
集群规模上去之后,让每个节点都持有全量数据不可行,需要分片(sharding)——把 key 空间划分给不同节点。分片方案决定了扩缩容时的数据迁移量:
| 方案 | 扩缩容影响 | 实现复杂度 | 典型实现 |
|---|---|---|---|
| 取模(mod N) | 几乎全部 key 重新映射 | 低 | 手写 |
| 一致性哈希 | 仅相邻节点受影响 | 中 | :hash_ring |
| 虚拟节点哈希 | 负载更均匀 | 中高 | 自建 |
| Rendezvous 哈希 | 仅受影响节点迁移,无环结构 | 中 | 自建 |
取模哈希的问题是 N 变化时几乎全部 key 都要迁移——从 4 节点扩到 5 节点,80% 的数据要搬家。一致性哈希把节点与 key 都映射到一个环上,增删节点只影响环上相邻的一段,迁移量约为 1/N。
Erlang 生态里 :hash_ring 提供了现成的环实现:
ring = :hash_ring.new([{node1, 128}, {node2, 128}, {node3, 128}])
:hash_ring.find_node(ring, "user:12345")
第二个参数是每个物理节点的权重,也就是它在环上放置的虚拟节点数。虚拟节点越多,负载越均匀,但环的元数据也越大。128 个副本是常见取值,能让节点间负载偏差控制在 10% 以内。
分片带来的新问题是数据归属变更:扩缩容后,原本属于节点 A 的 key 现在归节点 B,必须迁移。迁移期间要处理「旧位置与新位置同时可读」的过渡状态,通常做法是双写 + 迁移完成后切换路由表。
对于进程本身的分片(而不是数据),Horde.DynamicSupervisor 提供了更省事的选择:它基于 CRDT 在所有节点间同步进程归属,扩容时自动重新分配,无需手工迁移。代价是它只保证最终一致,分区期间可能出现同名进程的短暂重复。
一致性的工具选择
BEAM 提供的一致性原语各有明确的适用边界:
| 工具 | 一致性 | 适用 | 限制 |
|---|---|---|---|
:global | 最终(分区期不一致) | 单例进程名注册 | 节点数 > 50 时性能下降 |
:pg | 无(进程组,不保证唯一) | 广播、进程发现 | 不提供互斥 |
Horde.Registry | CRDT 最终一致 | 分布式注册表 | 冲突时可能短暂重复 |
Horde.DynamicSupervisor | CRDT 最终一致 | 分布式进程管理 | 同上 |
:ra | 线性一致(Raft) | 配置、元数据、队列 | 需要多数派,写延迟高 |
| 外部数据库 | 取决于实现 | 强一致业务数据 | 引入网络依赖 |
:global 是最容易被误用的。它能在集群内保证名字唯一,但分区期间两边会各自注册同名进程,恢复后 :global 会按内部规则杀掉其中一个——如果那个进程正在处理任务,任务就丢了。用 :global 做单例时必须让进程可安全重启,或改用 Horde.Registry 这类基于 CRDT 的方案。
分布式锁同理::global.set_lock/2 在分区期间不保证互斥。真正需要跨节点互斥的场景,应当使用具备共识语义的实现,例如 Redis 与 ZooKeeper/etcd 的分布式锁
,其原理与 共识算法
一脉相承。
实践建议
- 先明确一致性需求再选技术。多数业务只需要最终一致,强行上 Raft 会换来分区时不可用。
- 集群规模控制在 50 节点以内。全互联的 O(N²) 连接成本是硬约束,超大规模应拆成多个小集群。
- 多数派判定用静态期望节点数。动态推导在分区时必然失效。
- CRDT 只用于可交换的更新。计数、集合、状态标记适用;读-改-写的原子操作不适用。
:global不当互斥锁用。分区期间它不保证唯一,会静默产生重复单例。- 时间戳类 CRDT 必须配 NTP。物理时钟漂移会直接破坏 LWW 语义。
- 分区恢复要有流程,不能靠「连上就自动好了」。数据比对、单例重建、日志记录缺一不可。
- 用 libcluster 自动组建集群,但别指望它给一致性。它只负责连接,不负责收敛。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。