分布式一致性与网络分区:CRDT、libcluster 与脑裂治理

系统梳理 BEAM 集群的一致性问题:分布式 Erlang 的全互联拓扑与连接成本、libcluster 的 Gossip 与 Kubernetes 策略、CAP 取舍与一致性谱系、CRDT 类型与收敛语义、脑裂的检测与多数派治理、:global 与 Ra 的适用边界及分区恢复流程。

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²):

节点数连接数每节点连接数
5104
2019019
50122549
100495099

每个连接背后是一个独立的 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.GossipUDP 组播互探同网段、开发跨网段/云上组播常被禁
Cluster.Strategy.Epmd本地 epmd 注册表单机多节点仅限同主机
Cluster.Strategy.ErlangHosts.hosts.erlang 文件静态集群节点列表写死
Cluster.Strategy.Kubernetes.DNSheadless Service DNSK8s依赖 DNS 轮询
Cluster.Strategy.KubernetesK8s API 查 PodK8s需要 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 分布式数据 。

分区恢复流程

分区恢复不是「连上就好」,需要一个明确流程:

  1. 确认分区已消除:Node.list() 稳定包含全部期望节点,且持续数秒无抖动。
  2. 检查数据分叉:对 CRDT 状态做一次全量比对(DeltaCrdt.read/1 对比哈希),对数据库检查冲突记录。
  3. 重建单例:分区期间可能有两份单例进程,恢复后必须保证只剩一份。用 :global 的名称注册会在重连后自动选主,但需要确认旧进程已被终止。
  4. 恢复写入:多数派降级的节点重新加入可写集合。
  5. 记录事件:把分区起止时间、影响范围写入日志,作为后续调参依据。

心跳调优:分区检测的速度与误判

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.RegistryCRDT 最终一致分布式注册表冲突时可能短暂重复
Horde.DynamicSupervisorCRDT 最终一致分布式进程管理同上
:ra线性一致(Raft)配置、元数据、队列需要多数派,写延迟高
外部数据库取决于实现强一致业务数据引入网络依赖

:global 是最容易被误用的。它能在集群内保证名字唯一,但分区期间两边会各自注册同名进程,恢复后 :global 会按内部规则杀掉其中一个——如果那个进程正在处理任务,任务就丢了。用 :global 做单例时必须让进程可安全重启,或改用 Horde.Registry 这类基于 CRDT 的方案。

分布式锁同理::global.set_lock/2 在分区期间不保证互斥。真正需要跨节点互斥的场景,应当使用具备共识语义的实现,例如 Redis 与 ZooKeeper/etcd 的分布式锁 ,其原理与 共识算法 一脉相承。

实践建议

  1. 先明确一致性需求再选技术。多数业务只需要最终一致,强行上 Raft 会换来分区时不可用。
  2. 集群规模控制在 50 节点以内。全互联的 O(N²) 连接成本是硬约束,超大规模应拆成多个小集群。
  3. 多数派判定用静态期望节点数。动态推导在分区时必然失效。
  4. CRDT 只用于可交换的更新。计数、集合、状态标记适用;读-改-写的原子操作不适用。
  5. :global 不当互斥锁用。分区期间它不保证唯一,会静默产生重复单例。
  6. 时间戳类 CRDT 必须配 NTP。物理时钟漂移会直接破坏 LWW 语义。
  7. 分区恢复要有流程,不能靠「连上就自动好了」。数据比对、单例重建、日志记录缺一不可。
  8. 用 libcluster 自动组建集群,但别指望它给一致性。它只负责连接,不负责收敛。

继续阅读

探索更多技术文章

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

全部文章 返回首页

「erlang」更多文章

  1. BEAM 内存剖析与泄漏排查:recon、observer 与堆分析
  2. 缓存、限流与熔断:Cachex、Hammer 与降级策略
  3. Elixir 元编程与宏:AST、quote/unquote 与 DSL 设计