Klemens Lukaszczyk

A 1ms partition is enough to corrupt a node

· 5 min read

  • elixir
  • distributed-systems
  • pubsub

Messages are guaranteed to be delivered as long as the distribution link between the nodes is active, but it gets tricky when a partition occurs - that means both nodes online, but the link between them is down.

Phoenix.PubSub.PG2 delivers a broadcast to the nodes connected at that instant. No buffer, no ack. If a node is unreachable when you broadcast, the message is gone for that node - and nobody finds out.

Step through it:

Fire and forget

Node A - no buffer, no acks

Node B

B comes back and carries on as if nothing happened. No error, no gap detection - msg 2 and msg 3 simply never existed as far as B is concerned. If those messages were cache invalidations, B now serves stale data forever. If they fed a projection, B’s derived state is quietly wrong.

That’s the part that bothers me: not that messages can be lost - networks lose things - but that the loss is silent. You can’t retry what you don’t know you missed.

The fix: remember what you sent, and to whom

Picture the broadcasting node with a small notebook.

  • The notebook has a fixed number of pages (buffer_size, say 5), and one message goes on one page. Fill the last page and you loop back to the first, writing over whatever was there. So the notebook always holds the 5 most recent messages and nothing older - that’s all “page” means here.
  • Every message also gets a number that never resets: msg 1, msg 2, msg 3, and on forever. That counter is write_cursor. The pages cycle; the numbers don’t.
  • Next to each other node, A keeps a bookmark (read_cursors[node]): “this node has confirmed everything up to msg N.” It moves only when that node writes back “got it”. No reply, no move - so anything unconfirmed is simply sent again next time. That is the entire at-least-once guarantee.

Notebook (5 pages) - press the buttons

write_cursor:

Node B bookmark: behind:

So: message numbers count up forever, pages get reused. Subtract a node’s bookmark from the current message number and you get one number - how far behind it is. That number answers both questions at once: what does the node still need, and is it still in the notebook?

Three cases, and that’s the whole decision:

bookmark = read_cursors[node]             # how far this node has confirmed
lag      = write_cursor - bookmark        # how far behind it is

lag == 0           -> :caught_up      # it has everything, send nothing
lag <= buffer_size -> {:messages, …}  # still in the notebook, resend those pages
lag >  buffer_size -> {:expired, …}   # written over, can't be resent

The same partition, this time with a notebook

Step through the partition from the top of the post and watch the bookmark. The second tab walks the other case, where B stays away too long:

Node A - notebook (5 pages)

Node B

bookmark:

No gap: B gets 2 and 3, in order, before it ever sees 4.

And when B stays away so long that the messages it needed are written over, A does not quietly carry on and let B think it is up to date. B is told “I lost your place, reload from the database” - it knows its state is stale and fixes it. That’s the actual guarantee:

Either you receive every message in order, or you are told you fell behind.

Getting message 3 means you already have 1 and 2. Never a silent hole.

The cost is duplicates. If a batch is delivered and handled but the ack is lost, the cursor doesn’t advance and the batch is re-sent. Handlers have to tolerate that - send absolute state rather than deltas, or dedupe by message id.

Demo in action

ScoreBoard is three nodes on one live match. Take a node offline and its column freezes while the peers buffer for it (the orange box counts what’s held); bring it back and the backlog replays in order. Overflow the ring buffer - 5 here - and the box turns red instead: that peer gets {:cursor_expired, node} and reloads from the source of truth. Converged either way.

Three ScoreBoard nodes: one goes offline and freezes, its peers buffer, then the backlog replays and the boards converge

Performance: batching is the whole story

Measured end-to-end on a Fly.io cluster (fra, performance-4x), every node receiving every message, pool_size=1.

Every run below delivered every message, with no overflow:

batch_interval*payload3 nodes4 nodes
100 ms10 B79,445 msg/s75,468 msg/s
100 ms200 B41,784 msg/s43,479 msg/s
0 ms10 B1,319 msg/s787 msg/s

* batch_interval is how long a producer collects writes before flushing them - one acked round-trip per node carrying everything written in that window, instead of a round-trip per message.

Two numbers worth internalizing:

  • batch_interval=0 collapses to ~1k msg/s - 60 to 95 times slower. Every message becomes its own acked round-trip, so latency is the wall, not CPU. The buffer isn’t what makes this fast; amortizing the round-trip across a full batch is.
  • A 200 B whole-object payload costs ~42% of 4-node throughput vs 10 B. Send small deltas, not whole objects.

Delivery is network-bound, so with two or more peers the flush fans out to one task per node and costs the slowest single round-trip, not the sum.

When not to use it

The buffer is in memory and bounded. It closes the brief-partition gap; it is not a durable log. Want persistence across a full cluster restart, replay from disk, cross-language consumers, or huge backlogs? Run a real broker. Want reliable cross-node messaging on a BEAM cluster you already operate, without standing up a message broker to get it? That’s the gap this fills.

It’s also not a drop-in replacement - at-least-once costs more than fire-and-forget. Keep the default PubSub for ordinary broadcasts and run EchoPubSub alongside it, for the streams where a missing message quietly corrupts state.

EchoPubSub on GitHub · docs

← Back to blog