ZAB and Total Order Broadcast in Practice

LESSON

Consensus and Coordination

008 30 min intermediate

ZAB and Total Order Broadcast in Practice

By the end of this lesson, you will be able to...

  • Trace how ZAB turns a client update into an ordered, committed ZooKeeper proposal.

  • Explain why a new leader must activate and synchronize a history before accepting new proposals.

  • Distinguish total-order delivery from packet arrival order and from a guarantee that every read is fresh.

Idea in one sentence: ZAB lets ZooKeeper replicas deliver one ordered stream of committed updates, and it makes a new leader repair that stream before it is allowed to extend it.

Core Insight

A control plane uses a three-server ZooKeeper ensemble to coordinate a handoff. It stores a marker saying that worker group blue is draining, then a marker saying that worker group green may accept work.

The tempting model is: “The data is replicated, so it is fine if every server eventually receives both writes.” That works only when the relative order of the writes has no meaning. Here it does. A replica that applies green may accept work before blue is draining describes a different coordination history.

The harder moment comes when the leader fails halfway through a proposal. Some replicas may have written it to durable storage; others may not have seen it. A replacement leader cannot safely start assigning fresh sequence numbers merely because an election named it leader.

ZAB, ZooKeeper's atomic broadcast protocol, addresses both pressures. It makes committed messages deliver in one total order, and it separates a leader's activation from normal broadcast. The activation phase establishes a history that is safe to continue. Only then may the new leader accept client updates.

The Small Situation: One Update Stream, Three Partial Views

Call the current leader L and the followers F1 and F2. A majority is two servers. The following is a teaching trace; the zxids are illustrative.

committed history on all three servers

17:44  worker/blue = ready
17:45  worker/green = standby

The control plane sends two ordered updates:

u46: worker/blue  = draining
u47: worker/green = accepting-work

The intended service story is u46 then u47. That does not mean that every packet must travel through the network in a globally identical physical order. It means that if one correct ZooKeeper server delivers u46 before u47, every correct server that delivers both must deliver them in that same order.

This distinction matters because network delivery is local. A packet can be delayed on the path to F2 while reaching F1 quickly. ZAB's job is to turn those different observations into one logical delivery sequence, not to make the network behave as one cable.

The Initial Model: Elect a Leader, Then Continue Writing

It is reasonable to think that a new election solves the failure. The elected server has a new epoch, so perhaps it can immediately accept u47 and call it the next update.

That approach works when all replicas already have the same durable prefix. It fails when the old leader left partial work behind. Consider this intermediate state after L begins proposing u46:

Server Has proposal 17:46? Has a commit decision?
L yes maybe
F1 yes maybe
F2 no no

An acknowledgement is not merely “I received a packet.” In ZooKeeper's protocol, it means the server has recorded the proposal in persistent storage. A quorum of acknowledgements is the evidence that lets the leader commit a proposal and tell replicas to deliver its message.

Now let L crash. F1 knows more than F2, but neither follower can infer the full global state from that fact alone. If F2 became active and immediately proposed a fresh suffix, it could try to build on a history that ignores a proposal that a quorum had already made durable or committed.

The missing idea is not another ordinary client write. It is a phase that establishes which history the new leader is allowed to extend.

The Better Model: Activation Before Active Broadcast

Plain meaning: The replacement leader first gets a quorum onto a compatible history. It proves that it can safely continue the old stream before it starts a new one.

In this scenario: The new leader learns the relevant highest history from a quorum, synchronizes followers to its chosen history, and commits a leader-activation proposal. It does not accept u47 until that activation completes.

Technical name: ZooKeeper calls the protocol ZAB (ZooKeeper Atomic Broadcast). Its two broad phases are leader activation and active messaging. ZAB identifies proposals with a zxid, which contains an epoch and a counter. A new leader uses a new epoch, so a zxid also makes leadership transitions visible in the ordered history.

During active messaging, the leader assigns proposals in order, sends them in that order to followers, waits for quorum acknowledgements, and issues commits in order. Followers process commits in order and deliver a proposal's message only once it is committed.

client update
  -> leader assigns zxid and sends proposal in order
  -> replicas persist and acknowledge it
  -> quorum evidence lets the leader issue COMMIT
  -> replicas deliver the message in zxid order

The protocol's central service property is stronger than “every write reaches several machines.” It is a stream property:

if any server delivers a before b,
every server delivers a before b.

Reliable delivery adds the other half: a message delivered by one server is eventually delivered by all servers. The mechanism uses quorum intersection so that a leader selected from a quorum cannot lose the evidence of an already committed proposal.

A Worked Trace: A Leader Dies Between Proposal and Next Write

Start from the common committed prefix ending at 17:45.

Step 1: L proposes the draining marker

L receives u46 and assigns zxid 17:46. It sends the proposal to F1 and F2 in that order relative to earlier proposals.

L  -> F1: PROPOSAL 17:46, blue = draining
L  -> F2: PROPOSAL 17:46, blue = draining

F1 persists it and acknowledges. F2 is slow and has not received it yet. With L and F1, the ensemble has a majority. The leader can commit 17:46 and send a COMMIT message. Suppose L crashes immediately after F1 receives that commit, before F2 receives either message.

L:  17:46 persisted and committed, then crashes
F1: 17:46 persisted and delivered
F2: last known state is still 17:45

The important observation is that F2 being behind does not erase the committed update. L and F1 formed a quorum, and any future quorum intersects it. That intersection is why the recovery phase has a route back to the evidence.

Step 2: election names a candidate; activation earns authority

Assume F1 becomes the candidate chosen to lead. Election alone is not permission to broadcast. The candidate must see the highest relevant zxid from the followers involved in the quorum and establish the new epoch.

F1 therefore synchronizes the history. In this trace, it sends the missing 17:46 proposal and commit information to F2, bringing F2 to the same delivered prefix. A follower that instead contained an uncommitted higher suffix could be told to discard that suffix when the quorum evidence shows it was not committed. This is a recovery rule, not a guess based on which server looks newest.

Once a quorum has synchronized with the new leader, F1 proposes and commits NEW_LEADER for epoch 18. The details vary with the implementation's recovery messages, but the safety boundary is stable: the new leader does not accept ordinary proposals until this activation proposal is committed.

before activation completes: no ordinary new proposal
after activation completes: 18:0 NEW_LEADER is committed
then:                       18:1 may be assigned to u47

Step 3: active messaging resumes in one order

Now F1 can accept u47, give it zxid 18:1, collect persistent acknowledgements from a quorum, issue COMMIT, and let every replica deliver it after 17:46.

delivered stream

17:45  worker/green = standby
17:46  worker/blue  = draining
18:0   NEW_LEADER (protocol transition, not a client message)
18:1   worker/green = accepting-work

F2 may have learned the individual packets late, but it cannot legally deliver 18:1 before 17:46. The ordered stream survives the leader change because activation repaired the prefix before normal broadcast resumed.

So far: ZAB uses a leader for fast, ordered broadcast, but it does not trust the bare fact of leadership after a failure. Quorum-backed synchronization and activation turn a candidate into a leader that has earned a safe suffix to extend.

What This Changes for a ZooKeeper Client

Before this model, it is easy to say, “The leader puts writes in order.” After it, the more useful statement is: “The system has a committed delivery stream, and a replacement leader must recover that stream before it creates more of it.”

This helps explain why zxids are meaningful diagnostics. They expose the order of proposals and encode an epoch change. A gap between replicas' zxids can point to a follower that needs synchronization; repeated epoch changes can point to leadership instability. These are observations to investigate, not proof that one particular client operation failed.

There is also a boundary that deserves careful wording. Total-order delivery of committed messages does not mean every read from every ZooKeeper server is fresh. ZooKeeper's ordinary reads can be served locally and can be stale. Its writes are linearizable, but a client that needs a read to reflect a particular synchronization point must use the service's documented synchronization or quorum-aware pattern for its version and use case. Ordered broadcast solves the write-history problem; it is not a universal freshness guarantee.

Costs, Failure Boundaries, and Signals

ZAB improves a coordination service's ability to expose one recoverable sequence of changes. The trade-off is a leader-centered write path, persistent acknowledgements, quorum availability, and an activation pause after leader changes.

It helps when the ensemble can form a quorum and replicas can synchronize. It can stop making progress when the leader loses quorum, storage is too slow for acknowledgements, or leader changes repeat before activation completes. Safety can remain intact while client writes wait.

Useful signals include a growing gap in zxids between peers, followers repeatedly resynchronizing, proposals waiting for quorum acknowledgements, long leader-activation periods, and frequent epoch changes. These signals say that the ordered path is struggling. They do not justify bypassing ordering or treating an uncommitted proposal as delivered.

ZAB also does not make external side effects transactional. If a watcher receives an ordered state change and then triggers an email, a network call, or a workload migration, that external action still needs its own retry and idempotency design. ZAB orders ZooKeeper's state changes; it does not automatically order every consequence outside ZooKeeper.

Common Confusions

Confusion: “Total order means every packet arrives everywhere in the same physical order.”

Why it is tempting: The word “broadcast” sounds like one simultaneous network event.

Better model: Network paths can differ. The guarantee is the logical order in which committed messages are delivered by the replicated service.

Confusion: “Electing a new leader makes it safe to accept a new write immediately.”

Why it is tempting: Election decides who will coordinate next.

Better model: Election chooses a candidate. Leader activation first synchronizes a safe history and commits the transition that permits normal proposals.

Confusion: “A higher zxid on one server always means that server's suffix is committed.”

Why it is tempting: Higher numbers feel more authoritative.

Better model: A proposal must have quorum-backed commit evidence. Recovery uses quorum information to retain committed history and discard an uncommitted suffix when appropriate.

Confusion: “A totally ordered write stream makes all reads linearizable.”

Why it is tempting: One ordered history sounds like every observer must see the same newest state.

Better model: Ordered committed writes and read freshness are separate guarantees. Local reads can be stale.

Check Your Understanding

Check: F2 did not receive proposal 17:46 before the old leader crashed, but L and F1 had committed it. Can a newly active leader deliver a new client update before recovering 17:46 to F2?

Think first, then reveal.

Answer: No. The leader must first activate a safe, synchronized history with a quorum. Extending the stream before that step could make a later update appear before a committed predecessor at a replica that was behind.

Check: Does a high zxid by itself prove that its proposal was delivered?

Think first, then reveal.

Answer: No. It identifies a proposal's position, including its epoch and counter, but delivery requires the proposal's quorum-backed commit decision.

Practice: Review a Failover Claim

An operator says: “Our leader failed after proposal 41:9. Replica R1 has it on disk, R2 does not, and R3 has just won the election. We should immediately accept the next client write as 42:1; R2 will catch up later.”

Write a short review of this claim. State what R3 must do first, why the disk record on R1 is not alone enough to classify 41:9 as committed, and one operational signal that would make you pause the change.

Model answer: R3 must complete leader activation: obtain the relevant quorum history, synchronize followers to the safe prefix, and commit the activation transition before accepting ordinary proposals in its epoch. R1's record proves that one replica persisted 41:9, not that a quorum acknowledged and committed it. The review should pause on a missing quorum, a growing zxid gap, repeated leader activation, or storage latency that prevents acknowledgements. The safe outcome may retain or discard 41:9 depending on the quorum-backed recovery evidence; it must not be decided by a single replica's suffix.

Connections

The previous lesson used Raft joint consensus to preserve an authority bridge while membership changes. ZAB addresses a different transition: a leader change. In both cases, a new authority cannot skip the evidence that connects it to the old history.

The next lesson compares primary-backup, multi-leader, and leaderless replication. ZAB is a concrete primary-backup-style coordination design whose ordered write path is worth its cost when the service's state decides who may act.

Resources

Key Takeaways

PREVIOUS Raft Membership Changes and Joint Consensus NEXT Replication Models: Primary-Backup, Multi-Leader, and Leaderless