Elixir: Distributed Elixir

Distribution is a BEAM feature, not a framework: two nodes that share a secret (the cookie) form a cluster where sending a message to a remote process is the same send/2 as a local one. This lesson covers the mechanics, the conveniences, and — honestly — the failure modes.

Nodes and Cookies

# Start two nodes on one machine and connect them.
iex --sname alpha@localhost --cookie secret
iex --sname beta@localhost  --cookie secret

# In alpha:
Node.connect(:"beta@localhost")     #=> true — now they're a cluster
Node.list()                         #=> [:"beta@localhost"]
Two BEAM nodes connected with a shared cookie, messaging processes across the cluster via erpc
Fig. 1 — Fully-connected mesh: every node maintains a direct TCP link to the others it has connected to.

Calling Across Nodes

# :erpc runs a function on a remote node — the simplest remote call.
:erpc.call(:"beta@localhost", MyApp.Health, :check, [], 5_000)

# Broadcast a computation to every node in the cluster:
:erpc.multicall(Node.list(), MyApp.Health, :check, [])

# You can even spawn a process on a remote node (rarely the right tool,
# but important to understand):
Node.spawn(:"beta@localhost", fn ->
  IO.puts("running on #{node()}")
end)

Style note: prefer services with APIs over Node.spawn for real systems — remote-spawn code couples deployment topology into function bodies. Use distribution for registration, PubSub, and cluster-aware coordination.

Clustering in Production

Manual Node.connect works for experiments; real clusters form automatically. libcluster discovers nodes via DNS, Kubernetes API, or multicast — you configure a strategy once and the mesh maintains itself.

# mix.exs: {:libcluster, "~> 3.3"}

# application.ex — start the cluster topology with the tree.
children = [
  {Cluster.Supervisor, [topologies()]},
  MyApp.Repo,
  # ...
]

defp topologies do
  [
    k8s: [
      strategy: Cluster.Strategy.Kubernetes,
      config: [kubernetes_node_suffix: System.get_env("NAMESPACE"),
               kubernetes_selector: "app=my-app"]
    ]
  ]
end

Global Names and PubSub

Three registration levels, three scopes. :global keeps exactly one process per name in the whole cluster (a leader election of sorts). Registry with keys: :duplicate plus Phoenix.PubSub broadcast to every interested process across nodes — the machinery behind LiveView presence and channels.

# One-per-cluster process (e.g., a scheduler that must run exactly once):
case :global.register_name(:my_app_scheduler, pid) do
  :yes -> :ok
  :no  -> :already_taken      # another node owns it — stand down
end

# Cluster-wide publish/subscribe (started in the tree as
# {Phoenix.PubSub, name: MyApp.PubSub}):
Phoenix.PubSub.subscribe(MyApp.PubSub, "prices:BTC")
Phoenix.PubSub.broadcast(MyApp.PubSub, "prices:BTC", {:tick, 67_432})

Net-Splits and the CAP Reality

The BEAM default is availability with risk: after a partition heals, disconnected nodes reconnect and both sides keep running in the meantime — which can mean two processes believing they own :global resources. Production answers, in order of commonality:

  • Keep cluster state advisory. Authoritative state lives in Postgres/Redis; the cluster coordinates, it does not own.
  • Consensus for critical coordination: Raft (used by RabbitMQ's Khepri) or libraries like horde + delta_crdt for eventually-consistent registration.
  • Restart partitioned nodes: some teams run the VM under a heart/supervisor that reboots a node when it detects isolation; most teams prefer keeping state advisory instead.

The pragmatic rule: use distribution for mesh convenience (pubsub, presence, monitoring), reach for consensus only when double-execution is catastrophic.

Practice

  1. Start two nodes with the same cookie, connect them, and call :erpc.call/4 both ways.
  2. Register a process under :global on both nodes and observe which wins; kill it and watch the other take over.
  3. Set up PubSub on both nodes and verify a broadcast from one arrives on the other.

Next: Data & Ecto