Architecture & System Design › Distributed Systems
Total Order Broadcast
Delivering the same messages in the same order to every node.
Also known as: total order broadcast, atomic broadcast, ordered broadcast
Total order broadcast delivers the same messages in the same order to all recipients: every node sees the identical sequence, enabling replicated state machines (apply the same ordered commands, reach the same state). It’s consensus’s communication face — ordering agreed before delivery, execution deterministic after.
broadcast m1, m2, m3 → ALL nodes deliver [m1, m2, m3] (same order everywhere)
Implementations sequence via a leader (Raft log order), sequencer, or consensus per message/batch; recipients buffer gaps and deliver in order. The guarantee powers replicated logs, deterministic simulation, and any “everyone processes identically” design.
The classic mistakes:
- Assuming FIFO equals total order. Per-sender FIFO doesn’t order across senders; total order needs agreement, not just transport. Different guarantee, different machinery.
- Gaps mishandled. A missing sequence number blocks delivery indefinitely unless retransmission/failure paths fill or skip explicitly. Buffer with timeouts and recovery.
- Leader/sequencer bottleneck. All ordering through one node caps throughput and creates the failure domain; batching and partitioning (per-shard order) scale it.
- Deterministic execution assumed. Ordered delivery plus nondeterministic handlers (random, wall-clock, local state) still diverges replicas. Determinism must cover handlers, not just delivery.
- Slow-consumer blocking. One lagging recipient stalling global delivery couples everyone to the slowest; bounded buffers with catch-up (snapshots + replay) isolate laggards.
- Confusing with causal order. Total order is stronger (one global sequence) and costlier; causal suffices where only cause-effect matters. Pay for totality only where replicas must match exactly.
- Membership changes mid-stream. Joining/leaving recipients need state transfer aligned to the order (snapshots at sequence points); ad-hoc joins see inconsistent prefixes.
When to need it: replicated state machines and deterministic replicas — same order, same logic, same state. Weaker orders serve everywhere exact replication isn’t required.