NOTE
ZooKeeper Atomic Broadcast (ZAB)
What ZAB is, its roles and node states, leader election, broadcast/two-phase commit, request processing, and crash recovery.
This is a historical learning note and may contain outdated or incomplete understanding.
1. What Is ZAB
- A distributed consensus algorithm implemented by ZooKeeper itself.
- ZAB (ZooKeeper Atomic Broadcast Protocol): an atomic message-broadcast protocol that supports crash recovery.
- What is crash recovery: when the cluster has just started or the Leader crashes, ZooKeeper enters recovery mode and needs to elect a new Leader. After the election is completed, it synchronizes data with the other machines. When most servers have finished synchronizing, recovery mode ends.
- What is broadcast: once the Leader and most Followers have synchronized their state, the system enters broadcast mode.
- Sequential consistency
- Sequential consistency is a form of eventual consistency.
- After the Leader returns a successful write to the client, some Followers may not yet have written the data.
- Sequential consistency is implemented through ZXID.
- epoch: every election must produce a Leader. This epoch identifies the era of that Leader.
- incrementing value: the value increases by 1 for each write.
- Sequential consistency is a form of eventual consistency.
2. ZAB Algorithm Flow
2.1. Roles
- Leader: handles write requests.
- Follower: handles read requests and participates in voting.
- Observer: handles read requests and does not participate in voting.
2.2. Node States
- Looking: the state of looking for a Leader. When a server is in this state, it considers that there is currently no Leader in the cluster, so it needs to enter the Leader-election state.
- Following: follower state. Indicates that the current server’s role is Follower.
- Leading: leader state. Indicates that the current server’s role is Leader.
- Observing: observer state. Indicates that the current server’s role is Observer.
2.3. Three Stages
2.3.1. Leader Election
- Rule: the node with the larger zxid is elected Leader. If the zxid is the same, the node with the larger myid is elected Leader.
- Each vote consists of two parts: (myid, zxid).
- zxid

- epoch: whenever a new Leader is produced, the epoch of all ZooKeeper Servers is updated.
- xid: represents the underlying transaction number.
- myid
- the id of each server.
- All Nodes are in Follower state.
- If a Follower does not receive the Leader’s heartbeat for a certain period of time or the Leader does not receive heartbeats from more than half of the Followers, it switches itself to Candidate state and starts an election.
- If more than half of the Nodes return true, it is elected Leader.
- If no more than half of the Nodes return true, wait for a while and start another election.
- If a request from a Leader is received during the election and its term > the node’s own term, give up the election and become a Follower.
- After the election ends, the node becomes Leader or Follower.
2.3.2. Broadcast
2.3.2.1. Two-Phase Commit
- The Leader receives a client write request.
- The Leader uses two-phase commit.
- After receiving the write request, the Leader broadcasts a proposal to all Followers. Followers that can write return ACK.
- After the Leader receives ACKs from more than half, it returns that the write request was processed successfully to the client, and then broadcasts COMMIT to all Followers to make the proposal effective.
- The Leader returns write success or failure to the client.
2.3.2.2. Request Processing
2.3.2.2.1. Read Requests
- Can be handled by any node.
- The more machines in the cluster, the higher the read-request throughput.
2.3.2.2.2. Write Requests
- Must be handled by the Leader, which sends the request to all nodes.
- Pre-commit stage: send the data to all nodes. Only when more than half of the nodes accept does it proceed to the second stage.
- Commit: notify all nodes to commit the data.
- The more machines in the cluster, the lower the write-request throughput.
- What if the request reaches a node in the failed half? It may read old data. That is, ZooKeeper guarantees eventual consistency (other consistency types include strong consistency and weak consistency).
2.3.3. Crash Recovery
- The old Leader crashes, a new Leader takes over, and the other Followers switch to the new Leader and begin synchronizing data.
- Leader election is the same as Stage 1.
3. References
- Analysis of the ZooKeeper ZAB Protocol - Jianshu
- Paxos Algorithm - Wikipedia, the Free Encyclopedia
- Detailed Explanation of the Paxos Algorithm (1) - 割肉机 - CNBlogs
- ZooKeeper Atomic Broadcast Protocol (ZAB) and implementation of ZooKeeper - CloudKarafka
- A Brief Analysis of ZooKeeper’s Consistency Principle - Zhihu
- Interview Question: What Is the Zab Protocol? - Juejin
- Talking About ZooKeeper’s Sequential Consistency
- Study Notes: The ZooKeeper-Based Zab Protocol - SegmentFault
- ZooKeeper Interview Questions - Personal Article - SegmentFault
- Explaining ZooKeeper’s Election Mechanism in Plain Language - DockOne.io
Discussion
Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub