System Design Problem

Design a Distributed Coordination Service (ZooKeeper)

Commonly Asked By:GoogleYahooTwitterMeta

Interview Setup

Interview Prompt

Design a distributed coordination service (ZooKeeper style) for 1M znodes, 100K reads/sec, 10K writes/sec, and 50K long lived client sessions, acting as the CP backbone for legacy Kafka deployments, Hadoop, and service discovery.

Clarifying Questions (ask before designing)

QuestionWhy it matters
Is this optimized for reads or writes?Coordination services are typically read heavy (10:1 or 100:1 read to write ratio).
What are the core primitives we must support?Leader election, configuration management, and distributed locks dictate the need for Ephemeral nodes and Watches.
What should happen to writes during a network partition?The minority side must not make quorum backed writes progress if the system is preserving one authoritative history.
What consistency should clients observe on follower reads?Follower reads can be older than the leader state, so the design should distinguish sequentially consistent reads from a read that must be synchronized first.

Scope

In scope

  • Hierarchical data model (ZNodes)
  • Watches and event notification
  • Consensus and Leader Election (ZAB protocol)
  • Client session management and Ephemeral node lifecycle
  • ACLs, versioned writes, and atomic multi operations

Out of scope (state explicitly)

  • Storing large files (blob storage)
  • High-throughput data streaming (Kafka use case)

Functional Requirements

Start by asking your interviewer which coordination primitives are in scope, such as leader election, distributed locks, or configuration management, and whether the entire znode tree must fit in RAM on every ensemble member. Confirm ensemble size before discussing quorum math because fault tolerance and quorum size depend on it.

  • Hierarchical Namespace: Data is organized in a file system like tree of nodes (called znodes).
  • Znode Creation Modes: Combine Persistent or Ephemeral lifetime with optional Sequential naming. Ephemeral znodes are removed when the owning client session expires.
  • Watch Mechanism: Support classic one time watches plus persistent and persistent recursive watches in modern ZooKeeper releases. Watch events are asynchronous notifications, so clients read current state after receiving an event.
  • Core API: create, delete, exists, getData, setData, getChildren, sync, and atomic multi operation support.
  • Sessions and ACLs: Maintain heartbeats and session timeouts for Ephemeral node ownership, and support ACLs for node access control.

Non-Functional Requirements

Your interviewer will stress test write scalability through a single Leader and whether reads from Followers can be slightly stale. They will also probe split brain prevention during a network partition, where a minority partition becomes unavailable rather than serving divergent data.

  • High Availability: The service remains writable while a quorum is available. A minority partition cannot make quorum backed writes progress.
  • Strict Ordering: Committed updates share one authoritative zxid order, and a client does not observe a previously seen state being rolled back.
  • High Read Throughput: Typical workloads are 10:1 read to write ratio. Reads should scale across Followers and Observers without increasing write quorum size.
  • Strongly Ordered Writes: Updates are totally ordered and committed only after a quorum accepts them.
  • Session Liveness: Heartbeats and session timeouts detect failed clients reliably without expiring healthy sessions during normal JVM pauses and transient network delay.
  • Bounded Payloads: Znode payloads stay small so the full in memory tree and replication traffic remain predictable.

Capacity Estimations

The entire znode tree must fit in RAM on every ensemble member, so size the in memory dataset before assuming arbitrary payloads can be stored in znodes.

MetricCalculationValue
Ensemble sizeGiven (5 node quorum)5 nodes
Quorum for writes(5/2)+13 nodes
Znodes in memoryGiven (coordination metadata)~1M znodes
Avg znode payloadGiven~256 bytes
In memory dataset1M x 256 B + tree overhead~500 MB
Read requests/secGiven (10:1 read:write)100K reads/s
Write requests/secGiven10K writes/s
Active client sessionsGiven50K
Watch notifications/sec (peak)Given200K

The entire znode tree must fit in RAM on every ensemble member (~500 MB for 1M small znodes, before JVM, session, watch, ACL, and replication buffers). Provision substantial headroom above the raw tree estimate. The 10K writes/s figure is a workload target rather than a hard ZooKeeper limit. Actual write throughput is constrained by single Leader serialization, quorum network latency, transaction log I/O, and transaction size. Higher churn should be handled by splitting independent coordination domains across ensembles or by choosing a system with the required write profile. Observer nodes scale reads without enlarging the write quorum. Add them when replica capacity or read locality becomes the bottleneck rather than using a fixed 100:1 threshold.

Architecture Diagram

ZooKeeper is the CP coordination backbone for systems that need small, strongly consistent metadata such as leader election, distributed locks, and service discovery, rather than bulk data storage. Every ensemble member holds the full znode tree in RAM. Writes funnel through a single Leader via ZAB quorum while reads scale horizontally across Followers and Observer nodes, with follower reads providing sequential consistency.

At 100K reads/sec and 10K writes/sec, the architecture trades write scalability for strict ordering. Adding read replicas or Observers increases read capacity, while increasing the voting ensemble can increase quorum communication overhead. Watches push change notifications over long lived TCP sessions so clients do not need to poll for lock release or configuration updates.

Loading...

In the room

Frame ZooKeeper as CP coordination for small in memory metadata, not a database or a message queue. Ask what happens if a client stores a 10 MB blob in a znode, because the entire tree must fit in RAM on every node.

Component Deep Dives

1. The Data Model (Znode Tree)

ZooKeeper primitives, including persistent, ephemeral, and sequential znodes alongside classic one time watches and persistent watches, map directly to distributed systems patterns. The sections below connect the in memory tree model to ZAB atomic broadcast, read and write consistency guarantees, watch driven notifications, and the transaction log recovery path that survives cluster wide power loss.

The hierarchical znode tree forms ZooKeeper's data model, where persistent and ephemeral lifetimes plus sequential naming map directly to configuration storage, failure detection, and fair lock ordering.

Unlike a standard key value store, ZooKeeper organizes keys in a tree structure similar to a file system. The entire tree is kept in memory on every server, which makes reads fast but also means memory usage grows with the full replicated tree.

Loading...
  • Persistent Znodes: Remain in the tree until explicitly deleted. Used for static configuration.
  • Ephemeral Znodes: Bound to the client's session. If the client crashes and its session expires after missed heartbeats, ZooKeeper automatically deletes the znode. This is critical for failure detection and service discovery.
  • Sequential Naming: ZooKeeper automatically appends a monotonically increasing 10 digit counter within the parent namespace (for example, node-00000001). Combined with Ephemeral nodes and predecessor watches, this provides the ordering primitive used for fair locks and coordination queues.

2. ZAB Protocol (ZooKeeper Atomic Broadcast) ⭐

ZAB is ZooKeeper's leader based atomic broadcast protocol for replicating an ordered transaction history across the ensemble. It combines leader establishment, recovery, and broadcast. See our guide on Replication, Failover, and Leader Election.


--- ZAB (ZooKeeper Atomic Broadcast) Write Protocol ---

1. Client sends WRITE request to any ZooKeeper server.
2. If the request reaches a Follower, that server forwards it to the Leader.
3. Leader assigns the next zxid and creates a PROPOSAL.
4. Leader broadcasts the PROPOSAL to the voting Followers.
5. Voting Followers append the transaction to their local transaction log.
6. Followers return ACKs after durably accepting the proposal.
7. Leader waits for a quorum of 3 out of 5 voting servers.
8. Leader COMMITs the transaction and applies it to its in-memory znode tree.
9. Leader returns SUCCESS to the client.
10. Followers apply the committed transaction to their in-memory trees.

Notes:
• Observers receive replicated state but do not participate in the write quorum.
• The transaction log plus fuzzy snapshots support crash recovery.
• During leader recovery, the new Leader reconciles follower histories before new proposals resume.

During leader recovery, the new leader establishes the current epoch and reconciles follower histories before new proposals are accepted. Normal transaction broadcasting resumes only after recovery completes, which prevents an old leader from advancing a conflicting history.

Because the Leader coordinates all writes, write throughput does not scale linearly with more voting members. Increasing the voting ensemble can increase quorum communication and acknowledgement overhead, while adding Observer nodes expands read capacity without changing the write quorum.

3. Read & Write Paths

Write and read paths provide distinct consistency guarantees. Updates are totally ordered and committed through the Leader, while reads from Followers may be slightly stale unless the client synchronizes the target server first. Explore the trade-offs further in CAP Theorem and Consistency Models.

Write Path:
1. Client sends write to any node.
2. If node is a Follower, it forwards the write to the Leader.
3. Leader executes ZAB protocol (Propose -> Ack -> Commit).
4. Once committed, Leader applies it to memory and replies to the client.

Read Path:
1. Client sends read to the node it is connected to (Leader, Follower, or Observer).
2. The node reads directly from its local in memory Znode tree.
3. Reads are fast, but a follower may return slightly stale data until it processes the latest commit.
   (If the client needs the latest committed state before a read, it can issue 'sync()' to synchronize that server first).

Wait, aren't reads stale? Yes. ZooKeeper provides Sequential Consistency rather than simultaneous identical views for all clients. A follower or Observer can return older state than the Leader, but a client does not observe an update it has already seen being rolled back. When the client needs the latest committed state before a read, it can call sync() and then read from that server.

4. Watches & Push Notifications

Classic watches eliminate polling through push notifications and are one time triggers. Clients normally re register after an event fires. Modern ZooKeeper also supports persistent and persistent recursive watches. Lock implementations should watch only their immediate predecessor to avoid thundering herd storms.

Polling for changes (for example, polling for a lock release every 100ms) would overwhelm the cluster. ZooKeeper provides Watches for event driven updates.

  • When a client calls getData("/config", true), the server records the client's session as interested in /config.
  • If another client updates /config, the server pushes a watch event to the client over its persistent TCP connection. The event is a notification, so the client reads the current znode state rather than relying on the event payload as the source of truth.
  • Herd Effect Prevention: Watches are one time triggers. After firing, the client must set a new watch. Furthermore, in distributed locks, clients should only watch the specific znode immediately preceding their own in the sequence, rather than having all clients watch the lock holder.

5. Persistence & Snapshots

Because the tree lives in RAM, the transaction log combined with fuzzy snapshots enables ZooKeeper to recover committed state after cluster wide power loss. The durability model is also covered in Write-Ahead Logging and Data Durability and Idempotency and Exactly-Once Effects.

Since the tree is in RAM, a cluster wide power loss would erase the in memory state. ZooKeeper recovers it with two disk backed mechanisms:

  • Transaction Log: ZooKeeper writes the transaction representing an update to its transaction log before treating the successful update as durable.
  • Fuzzy Snapshots: Periodically, the in memory tree is serialized and dumped to disk. It is "fuzzy" because writes continue while the dump happens, meaning the snapshot is not point in time perfect.
  • Recovery: On boot, a server loads the latest fuzzy snapshot and replays the transaction log from the snapshot point to reconstruct the exact state. This lets the server recover the in memory tree without requiring a perfect point in time snapshot.

6. Sessions, Ephemeral Nodes, and Distributed Locks

ZooKeeper's most important coordination pattern combines client sessions with ephemeral sequential znodes. The lock holder must keep its session alive for the lifetime of the lock. A lock contender creates an ephemeral sequential node under a lock path, checks the lowest sequence number, and acquires the lock if it owns that node. Otherwise, it watches only the immediately preceding contender. When that predecessor disappears because its session expires or the holder releases the lock, the next contender wakes and retries. This avoids a thundering herd because each waiter watches one predecessor instead of the lock holder.

The same session mechanism supports service discovery. A healthy service instance owns an Ephemeral znode. If its process fails and its session expires, the registration disappears automatically and consumers can observe membership changes without polling.

7. etcd, Kubernetes, and the Displacement Context

Greenfield systems often choose etcd and Kubernetes over ZooKeeper. Acknowledge this displacement while explaining why legacy ZooKeeper dependent stacks continue to use ZAB semantics.

Modern greenfield systems often use etcd for strongly consistent cluster metadata and Kubernetes control plane state. Kubernetes exposes EndpointSlice objects for service endpoint discovery, so applications typically do not need a separate ZooKeeper coordination cluster. ZooKeeper remains relevant for legacy Kafka broker metadata deployments that still use ZooKeeper mode, Hadoop YARN, and HBase. Kafka 4.0 and later use KRaft and no longer support ZooKeeper mode. Applications built around ZooKeeper can depend on ZAB ordering, ephemeral sequential znodes, and classic watch semantics. High-availability database stacks such as PostgreSQL clusters managed by Patroni can similarly use DCS backends such as etcd or ZooKeeper for leader election.

Interview framing: ZooKeeper provides CP coordination for small in memory metadata rather than acting as a database or a message queue.

API Design

ZooKeeper provides a simple, file system like API. These primitives are combined to build distributed mechanisms such as locks, leader election, service discovery, and configuration management.

// 1. Authenticate, then create an Ephemeral ZNode
// CREATOR_ALL_ACL grants permissions to the authenticated creator
zk.addAuthInfo("digest", "app_user:app_password".getBytes());
String path = zk.create("/services/payment/node_1", "10.0.0.5".getBytes(), Ids.CREATOR_ALL_ACL, CreateMode.EPHEMERAL);

// 2. Read Data & Set Watch
// The 'true' flag registers a one time data watch for this client
byte[] data = zk.getData("/config/db_url", true, stat);

// 3. Write Data
// Uses the observed version for compare and set semantics
zk.setData("/config/db_url", "jdbc:postgresql://new-db".getBytes(), stat.getVersion());

// 4. Get Children & Set a classic one time Watch
// Useful for Service Discovery or Distributed Locks
List<String> children = zk.getChildren("/locks/resource_1", true);

// 5. Persistent recursive Watch (ZooKeeper 3.6+)
// Remains installed across events and includes descendant changes
zk.addWatch("/config", watcher, AddWatchMode.PERSISTENT_RECURSIVE);

// 6. Atomic multi operation
// Multiple mutations succeed or fail as one transaction
zk.multi(Op.setData("/config/a", "v2".getBytes(), -1), Op.delete("/config/b", -1));
  • Watches: Classic watches are one time triggers. Persistent watches can remain installed across multiple changes.
  • Versions: A client can pass the last known version for compare and set semantics. If another client modified the znode concurrently, the conditional write fails.
  • Synchronization: sync() establishes a synchronization barrier so a subsequent read on that server reflects the committed operations that precede the sync request.
  • ACLs: Each znode can define permissions for identities that need to read, write, create children, or administer the node.

Data Model

The entire data tree is stored in RAM for extreme read performance. Each node in the tree is called a ZNode.

// A ZNode is not just a byte array, as it also holds critical metadata (Stat structure)
struct Stat {
    int64 czxid;      // Zxid of the transaction that created this znode
    int64 mzxid;      // Zxid of the transaction that last modified this znode
    int64 pzxid;      // Zxid of the transaction that last modified children
    int32 version;    // Incremented on every data change (used for CAS)
    int32 cversion;   // Incremented on every child change
    int32 aversion;   // Incremented on every ACL change
    int64 ephemeralOwner; // Session ID if ephemeral, 0 if persistent
    int32 dataLength; // Length of the byte array
    int32 numChildren;
};

// Internal representation in RAM
class DataNode {
    byte[] data;
    Long acl;
    StatPersisted stat;
    Set<String> children;
}
  • Ephemeral Nodes: Tied to the client session. If the client crashes and its session expires after missed heartbeats, ZooKeeper automatically deletes the node. Crucial for Service Discovery.
  • Sequential Naming: ZooKeeper appends a monotonically increasing 10 digit sequence suffix within the same parent namespace. Distributed lock contenders use that ordering and watch only their immediate predecessor.

Fault Tolerance

ZooKeeper requires a strict majority quorum to elect a leader and commit writes. For an ensemble of N nodes, the quorum size is floor(N/2) + 1.

Ensemble SizeQuorum SizeFault Tolerance (Nodes can fail)
321
532
743

Split-Brain Prevention: In a network partition of a 5 node cluster splitting into groups of 2 and 3, only the group of 3 can form or maintain the quorum required to elect a leader and commit writes. The group of 2 cannot make quorum backed writes progress, so it cannot create a divergent authoritative history. A server on the minority side may still serve local reads depending on its mode, but those reads can be stale.

Additional Considerations

Interview Walkthrough

  • 25-minute cut

    Skip deep architectural variants unless targeting staff level.

    • In memory znode tree data model and RAM footprint (8 min)
    • Write path routing through the leader with ZAB quorum consensus (9 min)
    • Read path execution from followers and sequential consistency semantics (8 min)
  • Frame the system as CP coordination for small metadata (such as configuration, locks, and leader election) rather than a general database or queue.
  • Explain znode types: persistent for static configuration, ephemeral for session bound failure detection, and sequential for fair lock ordering. For detailed locking implementations, consult Distributed Lock Manager.
  • Write path: all writes route to the Leader and require quorum backed commit. Reads execute against Followers or Observers with sequential consistency, and sync() can be used before a read that needs fresher state.
  • Classic watches are one time triggers, so the client normally re registers after notification. Persistent watches are available when a longer lived subscription is needed. Lock implementations should watch only the immediate predecessor node to avoid herd storms.
  • A 5 node ensemble tolerates 2 node failures, and any minority partition halts writes to prevent split brain inconsistencies.
  • Acknowledge etcd and Kubernetes displacement for modern greenfield stacks while noting that ZooKeeper still backs legacy Kafka and Hadoop clusters.
  • Common pitfall: storing large payloads in znodes when the entire tree must fit in RAM on every ensemble member.

Engineering Trade-offs

In Memory Data vs. Data Size

Because the entire Znode tree must fit in RAM, ZooKeeper is not a general-purpose database. It is designed to store megabytes, not gigabytes, of data (for example, configuration strings and host IP addresses). A single znode should generally stay around 1 MB or less as an operational guideline, with substantially smaller payloads preferred for coordination metadata. Attempting to store large binaries in ZooKeeper will cause severe garbage collection pauses and network saturation during snapshot synchronization.

Read Scalability vs. Write Bottlenecks

ZooKeeper scales reads effectively by adding Follower nodes or Observers (nodes that serve reads but do not participate in voting quorums). However, writes encounter a bottleneck because every write must be serialized through the single Leader and acknowledged by a quorum. More voting members can increase coordination overhead, while Observers do not change the write quorum.

ZAB vs. Raft vs. Paxos

Paxos: A family of consensus protocols that establishes the core quorum based safety model, but is often harder to implement and explain than Raft.
Raft (etcd): Designed for understandability. Joint consensus for membership changes. Strongly consistent reads by default.
ZAB (ZooKeeper): A leader based atomic broadcast protocol tailored to ZooKeeper's primary backup state machine replication model. It orders transactions through the leader and uses a recovery phase to reconcile follower histories before new proposals resume after leadership changes.

CP vs AP (CAP Theorem)

ZooKeeper is commonly classified as a CP coordination system because quorum backed writes preserve one authoritative history during partitions. If a server loses contact with the quorum, that side cannot become the authoritative leader or commit quorum backed writes. Read behavior depends on the server mode and can expose older local state. The key trade off is loss of quorum backed write availability on the minority side in exchange for preventing split brain writes.

💬Review

Help Us Improve

How helpful was this walkthrough?

Click a star to rate. We actively use this feedback to refine and update our system design content.

Placeholder
Optional but highly appreciated!

Discussion

Share your thoughts, ask questions, or help others.

Loading comments...