Skip to content

Day 5: Talking across nodes

Days 1 through 4 got you a formed cluster: nodes discovered, handshaken, meshed, and monitored. Today is processes on different nodes finding each other and exchanging messages. This is the 02_processes section of the companion repo and Section 2 of the talk.

The thesis does its heaviest lifting today, so name it up front and keep checking against it. Everything in this briefing is a registry or a router, a way of answering “who should receive this message?” across the cluster. Each of these tools answers that question with some flavor of local opinion, replicated opinion, or deterministic computation. None of them answers it with agreement. Each registry is a different way to be wrong during churn, and the job on stage is to show that picking a registry means picking which wrongness you can tolerate. That framing collapses the whole zoo (Registry, :pg, :global, :syn, Horde, HRW) into one coherent slide instead of five library advertisements.

1. The baseline: sends are location-transparent, discovery is not

Section titled “1. The baseline: sends are location-transparent, discovery is not”

Worth thirty seconds of stage time to re-anchor. Once nodes are connected, these all work and all compile down to the SEND / REG_SEND control messages from Day 3:

send(pid, msg) # pid can live anywhere, SEND opcode
send({:my_server, :"b@host"}, msg) # registered name + explicit node, REG_SEND
GenServer.call({MyServer, :"b@host"}, :ping) # same thing with call semantics
:erpc.call(:"b@host", M, :f, [a]) # spawn-and-wait on the remote side

The transparency layer makes delivery trivial. What it deliberately does not provide is discovery: nothing built into sends tells you which node holds the process you want. The second form above hardcodes the node name, which is exactly what breaks under rolling deploys (Day 4: node names change every deploy in Kubernetes). Today is about the gap between “I can message any pid in the cluster” and “I know which pid to message”, and every tool below sits in that gap.

Also re-anchor the failure semantics, because they shape everything. Remote sends are fire-and-forget. send/2 to a dead pid, a nonexistent name, or across a connection that is about to drop returns :ok identically. Delivery is at-most-once, ordering is only guaranteed pairwise between two processes, and the only backpressure is the distribution buffer filling up (the busy_dist_port condition, numbers on Day 10’s cheat sheet). Monitors and links buy failure visibility back, at the cost of more dist traffic.

:pg is the right first tool to teach because it is easy to describe accurately: a cluster-wide multimap from group name to member pids, with strong eventual consistency and nothing more. OTP 23 introduced it (contributed from WhatsApp’s production stack) and OTP 24 deleted its predecessor pg2 entirely, so any audience member on a supported OTP has :pg and does not have pg2.

The design difference from pg2 explains everything about the module. pg2 serialized every join and leave through :global cluster-wide locks, which made membership changes O(cluster) and made a hot-joining group a throughput disaster. :pg dropped the locks. Joins and leaves take effect locally and immediately, and propagate to peers asynchronously. Reads never leave the node. You trade away “all nodes see the same membership at the same time” and get back fast, partition-tolerant membership operations. On the CAP menu, this is the AP corner.

Mechanics to have cold:

Scopes. A scope is an independent :pg overlay with its own process and its own ETS table. The default scope is named pg and only starts if the kernel parameter start_pg is set, so in practice Elixir apps start their own:

application.ex
children = [
%{id: :pg, start: {:pg, :start_link, [NetzLive.PG]}},
...
]
# anywhere
:ok = :pg.join(NetzLive.PG, {:feeder, substation_id}, self())
pids = :pg.get_members(NetzLive.PG, {:feeder, substation_id})
local = :pg.get_local_members(NetzLive.PG, {:feeder, substation_id})

Scopes exist to cut membership gossip. Each scope only syncs with the same scope name on other nodes, so unrelated subsystems do not pay for each other’s churn. A scope name that exists on only one node is effectively a local-only group server, which is occasionally useful.

Join semantics. Only local pids can be joined (:pg.join/3 with a remote pid fails). That restriction is a feature: each node is authoritative for its own members, and that is precisely what makes lock-free replication sound. A process can join the same group multiple times and must leave the same number of times. Groups are created implicitly on first join and vanish when empty. Members that die are removed automatically, and all members on a node that disconnects are removed from every other node’s view.

Reads. get_members/2 is a local ETS lookup of the replicated view, get_local_members/2 returns only pids on this node, both in no particular order, both optimized for speed. The local/all distinction is the backbone of the PubSub trick in the next section.

Monitoring, OTP 25.1 and later. :pg.monitor/2 subscribes you to one group and returns {ref, current_members}; :pg.monitor_scope/1 subscribes to everything in a scope. Updates arrive as {ref, :join, group, pids} and {ref, :leave, group, pids}. This is the monitor_nodes pattern from Day 4 lifted one layer up: the runtime hands you a per-observer-consistent event stream about group membership, with the same caveat that two observers’ streams need not agree in the moment.

Two caveats for the slide. First, eventual means eventual: after a nodeup, there is a window where node A sees B’s members and B does not yet see A’s. Code that assumes get_members is symmetric across nodes during churn is wrong. Second, membership is not transitive. :pg syncs over direct dist connections, so if A and C are both connected to B but not to each other (a connect_all false topology, or a partial mesh mid-formation), A and C do not see each other’s members even though B sees both. In the default full mesh this never bites. Mention it anyway, because it kills the “pg is a cluster-wide truth” intuition, and because hidden nodes are invisible to :pg for exactly this reason.

3. Phoenix.PubSub: :pg plus one message per node

Section titled “3. Phoenix.PubSub: :pg plus one message per node”

Phoenix.PubSub is a good dissection candidate because the whole room uses it daily and almost nobody knows it is about forty lines of :pg glue. Current release is v2.4.0 (hexdocs title verified; the adapter is still called Phoenix.PubSub.PG2 purely for historical continuity with the pg2 era).

It has two layers, and the division of labor between them is the useful part.

The local layer is Registry. Subscriptions are per-node. Phoenix.PubSub.subscribe(MyApp.PubSub, "topic") writes into a partitioned Registry running in duplicate mode on the local node, and local_broadcast/3 is a Registry.dispatch/3 over local subscribers. No distribution involved at all. Registry is ETS, partitioned by default by scheduler count, so local fan-out to tens of thousands of subscribers is parallel and cheap.

The distributed layer is :pg with one worker per node. The PG2 adapter starts a :pg scope named Phoenix.PubSub plus a small pool of PG2Worker processes, and each worker joins a :pg group. The group members are therefore one worker pid per node per pool partition, not subscribers. A cross-cluster broadcast does this:

  1. Pick a broadcast group by hashing the calling pid: elem(groups, :erlang.phash2(self(), tuple_size(groups))). Each sending process consistently uses one partition, spreading load across the pool.
  2. :pg.get_members for that group, then send exactly one message, {:forward_to_local, topic, message, dispatcher}, to each member pid whose node differs from the local node, and do a local_broadcast for the local subscribers.
  3. Each remote worker receives that one message and runs local_broadcast on its own node, fanning out via its local Registry.

So a broadcast to a 6-node cluster with 50,000 subscribers costs five messages on the distribution links, then 50,000 purely local deliveries done in parallel on each node. The dist mesh carries per-node traffic, not per-subscriber traffic. This is the most important scaling property of PubSub, and it falls directly out of the get_members vs get_local_members split in :pg. Good whiteboard material: draw the naive version (sender delivers to every remote subscriber pid over the dist link) next to the real version (forward once per node, fan out locally).

direct_broadcast/4 targets a single node by sending to the worker’s registered name as {group_name, node}, classic REG_SEND, no :pg lookup at all. The docs warn against using it for the local node because a send through your own worker serializes through that one process, whereas local_broadcast dispatches from the caller.

Guarantees, stated plainly for the talk: PubSub broadcast is fire-and-forget at every hop. send/2 to the remote worker is at-most-once; if the dist connection drops mid-broadcast, some nodes got it and some did not, and nobody is told. There is no ack, no retry, no ordering across senders. For presence-style and cache-invalidation-style traffic this is the right trade, and Phoenix.Presence layers a CRDT on top precisely because raw PubSub cannot carry state. The one-sentence summary for the talk: PubSub reaches every node that is currently reachable and says nothing about the rest.

Two version notes for accuracy on stage. Pool defaults have drifted: the PG2 adapter docs say pool_size defaults to 1 while child_spec/1 documents one partition per 4 cores. The code resolves this, but verify before stage which number you quote. And custom dispatchers are deprecated in 2.4 in favor of a :sender option on subscribe/3, so avoid presenting the old fastlane/dispatcher mechanism as current API. Verify before stage if you plan to mention fastlanes at all.

4. :global: the replicated name table that chooses C over A

Section titled “4. :global: the replicated name table that chooses C over A”

:global sits at the opposite corner of the menu from :pg, and teaching them back to back is the cleanest CAP illustration the BEAM offers.

The read path is fully local. Every node holds a complete replica of the name table, and :global.whereis_name(name) never leaves the node. Lookups scale perfectly.

The write path is fully global. :global.register_name(name, pid) is synchronous across the entire cluster. The docs commit to “the name is either registered on all nodes or none”, and registration and lock services are atomic with all involved nodes sharing the same view. Under the hood the global name server takes a cluster-wide lock, replicates the entry everywhere, and only then returns :yes. The consequences: registration latency grows with cluster size, concurrent registrations serialize against each other, and a slow or flaky node slows everyone’s writes. :global is for a small number of long-lived names, not for churny per-entity registration. Registering a process per device or per user via :global is the classic misuse; that workload belongs to :syn, Horde, or HRW routing.

Usage is pleasantly boring in Elixir:

GenServer.start_link(NetzLive.Coordinator, arg, name: {:global, :grid_coordinator})
GenServer.call({:global, :grid_coordinator}, :status)
# or the via form, interchangeable with Registry via tuples:
name = {:via, :global, :grid_coordinator}

Now the part that sets up Days 6 and 7. What happens when two partitions each registered the same name and then heal? On reconnect, :global detects the clash and calls a resolve function once per conflicted name. The three built-ins:

  • random_exit_name/3 (the default for register_name/2 and re_register_name/2): keeps one pid at random and kills the other.
  • random_notify_name/3: keeps one at random, sends the loser {:global_name_conflict, name} and lets it decide.
  • notify_all_name/3: unregisters both, tells both, lets the application sort it out.

Worth stating the default explicitly on stage: the standard library’s default conflict policy is to pick a survivor at random and unconditionally exit the other process. For a stateless worker that is fine. For a singleton coordinator holding in-flight state, random process exit on heal is the landmine, and the Swoosh incident (Day 7) is what happens when application code is linked to the losing side. If the resolver crashes or returns anything but one of the two pids, the name is simply unregistered on both sides, which is its own flavor of surprise.

One forward reference to plant today for tomorrow: since OTP 25 the kernel flag prevent_overlapping_partitions defaults to true, and :global now actively disconnects nodes to force clean, fully-connected partitions instead of overlapping ones. It changes how netsplits look in practice and it will change how your demo behaves. Tomorrow’s entire briefing is failure. Today, just register that :global is the module that cares most about partitions, because its correctness depends on the full mesh (the docs are explicit that reliable behavior requires both connect_all and prevent_overlapping_partitions enabled).

With :pg and :global understood as the two poles, the rest of the menu sorts itself. The question that matters for each row is what the tool does when two nodes disagree, not what API it has:

Tool Scope Unique names? On conflict / partition Write cost
Registry single node unique or duplicate n/a, local only local ETS
:global cluster unique resolver on heal, default kills one at random synchronous, all nodes, locked
:pg cluster no, groups of many views diverge, merge on heal, nothing killed local, async propagation
:syn cluster both (registry + groups) automatic resolution via callback, a loser is picked local, async propagation
Horde cluster unique (Registry API) CRDT merge, duplicate process terminated local, delta-CRDT gossip

Points worth making per row, beyond the table.

Registry is not distributed, and that is its virtue. It is the local building block everything else composes with (PubSub’s layer 1). The honest pattern for “distributed registry” in many real systems is local Registry plus a routing rule for which node owns the key, which is where HRW comes in below. Mention that Registry is partitioned ETS with :via support, then move on; the audience knows it.

:syn (current version 3.4.2) is “what if :pg also did unique names”. Same architectural family: scopes, node-local authority, async replication, strong eventual consistency, no global locks. The differences: it offers a registry (unique name to one pid, with metadata attached to registrations) alongside groups, and it resolves registration conflicts automatically after a netsplit heal rather than leaving both. By default it picks a survivor and terminates the loser, and the policy is customizable through its event handler behaviour (verify the exact default policy and callback name against the 3.4 docs before putting it on a slide). :syn is the pragmatic choice for high-churn per-entity registration at serious scale, the workload :global cannot take.

Horde (current version 0.10.0, last release Nov 2025) attacks the same problem with delta-CRDTs and adds the piece nobody else has: Horde.DynamicSupervisor, which distributes child processes across nodes and restarts them elsewhere when a node dies, plus Horde.Registry with a Registry-compatible API. The caveats to state honestly: the CRDT sync interval means there is a real window where a name resolves nowhere or a process runs twice; after a heal, Horde terminates the duplicate, so your processes must be built to die and hand off state; and the project has historically moved slowly, so check issue tracker health before recommending it from a stage. Horde is the maximal “make the cluster look like one supervisor” abstraction, and the thesis cuts against it: the more a library pretends the cluster is one machine, the more surprising its behavior during the minutes when the cluster visibly is not one machine.

The summary for the slide: every distributed registry either rejects writes until everyone agrees (:global), or accepts writes everywhere and later picks a loser (:pg merges, :syn and Horde kill the duplicate). There is no third option, because agreement under partition is exactly the thing distribution does not give you. That summary is also the transition into the Raft half of the talk.

Rendezvous hashing (highest random weight) earns its slot in the talk because it removes the registry problem rather than solving it. The question “which process owns key K” becomes a pure function, computed independently and identically on every node:

defmodule NetzLive.HRW do
@doc "Deterministically pick the owner of `key` from `nodes`."
def owner(key, nodes) do
Enum.max_by(nodes, fn node -> :erlang.phash2({key, node}) end)
end
end
# routing a rate-limit check, jola.dev style:
def hit(ip) do
nodes = [Node.self() | Node.list()]
node = NetzLive.HRW.owner(ip, nodes)
if node == Node.self() do
RateLimiter.local_hit(ip)
else
try do
GenServer.call({RateLimiter, node}, {:hit, ip}, 1_000)
catch
:exit, _ ->
# owner unreachable: degrade to a local check rather than failing open loudly
RateLimiter.local_hit(ip)
end
end
end

Why this works with zero coordination: :erlang.phash2/1 is documented as portable, the same input hashes to the same value on every node, architecture and OTP version notwithstanding. Two nodes given the same member list always compute the same owner, without ever exchanging a message about it. No lock, no replication, no resolver. HRW also has the minimal-disruption property: when a node leaves, only the keys it owned move (each remaining node keeps winning every comparison it already won); when a node joins, it steals roughly 1/n of the keys, evenly from everyone.

Against consistent hashing (ExHashRing in the Elixir world): HRW is stateless, needs no ring to build or rebalance, and costs O(n) hash evaluations per lookup instead of O(log n) against a prebuilt ring. Jola’s write-up concludes, reasonably, that for clusters under about ten nodes HRW is simpler and distributes at least as well, and NETZlive-sized clusters live comfortably in that range.

HRW is also the most instructive failure case of the day, so run the thesis through it. The function is deterministic, but its input is [Node.self() | Node.list()], each node’s local opinion of membership from Day 4. During a deploy or a partition, node A’s list and node B’s list differ for seconds to a minute, and during that window the same IP has two owners, each counting hits in its own ETS table. The rate limit is enforced at up to 2x for those keys until views converge, and nobody is notified. Jola’s post is candid about this: the limiter is “mostly” consistent, accurate “as long as the cluster is healthy”, and the right failure posture is to stay available and over-admit rather than block traffic on registry health. HRW did not remove the agreement problem. It moved the problem into the membership list and then chose availability. For rate limiting that trade is right. For “exactly one process may write this ledger entry” it is wrong, and that difference is what separates Section 2 of your talk from Section 3.

The same routing pattern generalizes past rate limiting: node-local caches with HRW routing, per-device connection ownership, sharded in-memory aggregation. It fits any workload where keys want a home node and a brief double-ownership window costs money rather than correctness.

7. Production notes to carry into the talk

Section titled “7. Production notes to carry into the talk”
  • Match the tool to churn and cardinality. Low-churn, low-cardinality, must-be-unique: :global. High-churn, high-cardinality: :syn, Horde, or HRW over local state. Many-recipients fan-out: :pg directly or PubSub on top. Putting per-user processes in :global and cluster singletons in eventually-consistent registries are the two symmetric mistakes.
  • Fan-out pressure lands on the dist buffer. PubSub’s forward-once-per-node design keeps dist traffic low, but a hot :pg group with large payloads can still fill a connection’s buffer and trigger busy_dist_port, which suspends any process that sends to that node, including innocent bystanders. Payload discipline (send ids, not blobs) matters more across nodes than within one. Numbers belong on Day 10’s cheat sheet.
  • Every registry read is a snapshot. get_members, whereis_name, and Node.list are all answers-as-of-now. The pid you got can be dead before your send executes; design with monitors and timeouts, not with lookup-then-trust.
  • Scopes are cheap isolation. Separate :pg scopes (and :syn scopes) per subsystem keep one noisy subsystem’s membership churn from delaying another’s, and make Observer archaeology far easier at 3am.
  • For the demo repo: the Day 4 :peer harness extends directly. Start three peers, start a :pg scope on each via :erpc, and all of today’s behavior is observable in a script, including the asymmetric-view window if you join members in a tight loop right after connect.

Active block (45 to 60 min, repo section 02_processes)

Section titled “Active block (45 to 60 min, repo section 02_processes)”

Build 02_processes/registry_tour.exs on top of the Day 4 :peer harness: (1) start three peers and start a :pg scope on each via :erpc.call/4; (2) on each peer, spawn a subscriber process that joins group {:topic, :grid} and prints anything it receives; (3) from ctl, implement a ten-line PubSub: get_members, send one {:forward, msg} per remote node (pick one member per node) and have that member fan out via get_local_members, then print per-node delivery counts to prove the one-message-per-node property; (4) register a :global singleton from two different peers concurrently and show one register_name returning :no; (5) implement HRW.owner/2 and map 10,000 synthetic keys across the three nodes, print the distribution, then :peer.stop/1 one node and print how many keys moved (expect roughly one third, all of them from the dead node). Stretch: in step 5, take the membership list on each node via :erpc immediately after the stop and find a moment where two nodes compute different owners for the same key.

  1. A Phoenix.PubSub broadcast goes out on a 6-node cluster with 40,000 subscribers spread evenly. Exactly how many messages cross distribution links, which processes send and receive them, and what performs the final delivery on each node?
  2. :global.register_name/2 and :pg.join/3 sit at opposite ends of a trade. State precisely what each one does synchronously versus asynchronously, what each costs as the cluster grows, and what each does about two nodes claiming the same name or group after a partition heals.
  3. Your HRW-routed rate limiter briefly enforces limits at 2x during a rolling deploy. Walk the exact mechanism: which input diverges, on which nodes, for how long (tie it to Day 4’s nodeup/nodedown machinery), and why the same mechanism makes HRW unsafe for singleton ownership.

Say it out loud (60 seconds, to the Haarlem room)

Section titled “Say it out loud (60 seconds, to the Haarlem room)”

“Once nodes are connected, sending a message anywhere in the cluster is trivial; knowing where to send it is the whole game. Explain the menu: a replicated name table that blocks until every node accepts the write and kills a random survivor after a split (:global), process groups that accept writes instantly everywhere and let views disagree for a moment (:pg, and Phoenix.PubSub is just :pg forwarding once per node), and a pure hash function that computes the owner with no registry at all (HRW), which only agrees as long as everyone’s member list does. End on the point: all three are reach with different failure flavors, and none of them is agreement.”