Hiring BarSupport

Design with ZooKeeper & etcd — Staff-Level Technology Guide

Technology guide37 min read6 diagrams

Why This Matters#

ZooKeeper and etcd are not databases, and they are not lock services either. They are small, strongly consistent sources of truth about who is in charge — which process is leader, which config version is live, which shard lives where, which member of a group is still alive. Their entire value is linearizability for a few megabytes of critical metadata. Put 8GB of application data in them, or put them in the per-request path, and you have built the most fragile component in your architecture.

They appear in Staff interviews whenever the design needs exactly one of something: one scheduler firing a cron job, one primary per shard, one writer to a partition, one active config. "Design a distributed job scheduler", "Design a key-value store", "How do you handle leader election?", "How does your service discovery work?" — all end at a consensus system. The L5 candidate says "use ZooKeeper for leader election." The L6 candidate says "the leader holds an etcd lease with a 10-second TTL and every write to storage carries the lease's revision as a fencing token, because a GC-paused ex-leader will wake up still believing it's leader." The L7 candidate asks how many consensus clusters the org is running, who is on call for them, and whether this system should depend on one at all — or on a conditional write in the database it already has.

The L5 → L6 gap is not knowing Raft vs ZAB. It is knowing that a lock is a lease, a lease can expire while you're still holding it, and only fencing makes that safe.

The L5 → L6 → L7 Contrast#

BehaviorSenior (L5)Staff (L6)Principal (L7)
First move"ZooKeeper for coordination""What must be exclusive, and what's the cost of two holders for 10 seconds? That decides lease TTL and fencing.""Do we need a consensus cluster at all, or can the database we already run give us a conditional write?"
LocksAcquire lock, do work, releaseLease + fencing token checked by the resource; locks for efficiency vs locks for correctnessStandard library for leader election with fencing baked in; bans hand-rolled locks
DataStore config and state in itMetadata only: KBs per key, < 1–2GB total, watches not pollingTracks every consumer of the shared cluster, quotas per client, capacity reviews
Failure"It's replicated, 3 nodes"Quorum math, disk fsync latency, session expiry storms, what degrades when quorum is lostConsensus cluster is a tier-0 dependency: cell-per-region, blast radius mapping, game days
OwnershipInfra team runs itPlatform owns ensemble; clients own session handling and correct recipesDecides centralized vs embedded consensus per system; consolidates or retires clusters
Scale"Add nodes""Adding voters makes writes slower. Add observers/learners for reads, or split the keyspace into separate clusters."Multi-cluster topology with clear ownership boundaries
Why "Locks" separates levels

A distributed lock in ZooKeeper or etcd is held by a session or lease that survives only as long as heartbeats arrive. A process can pause — a 15-second GC, a VM migration, a kernel stall — longer than the lease. The coordination service correctly expires the lease and grants the lock to someone else. The paused process wakes up, still inside its critical section, and writes. Two writers. The only fix is at the resource: every write carries a monotonically increasing token (ZooKeeper zxid/sequence number, etcd revision), and the resource rejects tokens older than the newest it has seen. The Senior answer is correct in the happy path; the Staff answer is correct during the pause.

Why "First move" separates levels at L7

Running a consensus cluster is a real operational commitment: dedicated low-latency disks, careful upgrades, a team that understands quorum loss. Many "we need ZooKeeper" designs only need a conditional write — UPDATE leases SET owner=?, epoch=epoch+1 WHERE name=? AND expires_at < now() in Postgres, or a DynamoDB conditional put. The Principal move is to ask whether the new dependency pays for itself.

The 60-Second Pitch#

"For leader election and cluster membership I'd use etcd — a Raft-replicated, linearizable key-value store. Three or five members in separate zones tolerate one or two failures. Leaders hold a lease with a ~10-second TTL, kept alive every ~3 seconds; if the leader dies, the lease expires and a standby wins a compare-and-swap on the leader key within about 10–15 seconds. Every action the leader takes on shared storage carries the lease's revision as a fencing token, so a stalled ex-leader can't corrupt anything. Clients watch keys instead of polling, so config changes propagate in milliseconds. It holds metadata only — a few thousand keys, well under a gigabyte. I would not use it for high-volume data, per-request locks, or anything that a conditional write in our main database already handles."

The Three Intents#

IntentConstraintStrategyFailure ModeCorrectness Bar
Leader election / mutual exclusionExactly one active actor per role; fast failoverLease/ephemeral node + CAS on a leader key + fencing tokensSplit brain from paused ex-leaderNever two effective writers
Configuration & metadata store (shard maps, feature config, schema versions)Linearizable, versioned, watched by many clientsKeys with revisions, watches, CAS updates, small valuesWatch storms, oversized values, stale reads from followersEvery client converges on the same version
Membership & service discoveryLiveness in seconds, thousands of clientsEphemeral nodes / leased keys + watches, or a dedicated discovery systemSession expiry storm on network blip deregisters the fleetEventually accurate within seconds; availability over precision

🎯 Staff Move: "I'll scope etcd to leader election and the shard map — both small, both need linearizability. Service discovery for 20,000 pods I'd put on a system built for it, because discovery wants availability during partitions and consensus systems choose consistency."

The Staff Positions#

PositionRationale
Metadata only; kilobytes per key, < ~1–2GB totalEvery write goes through one leader and is fsynced on a quorum; the whole dataset lives in memory.
Every lock is a lease; every lease needs a fencing tokenPauses longer than the TTL are inevitable at scale.
3 voters by default, 5 for tier-0; never even numbers4 tolerates the same 1 failure as 3 with slower writes.
Keep it off the per-request pathConsensus write latency is ~2–10ms and the cluster is shared; the data plane must survive its outage.
Watches, not polling10K clients polling every second is 10K reads/s of pure waste.
Dedicated low-latency disksfsync latency is leader-election latency; a noisy neighbor on the disk causes elections.
Clients must tolerate unavailabilityCache the last known config and keep serving; losing quorum should freeze changes, not stop traffic.

Architecture & Internals#

Five internals change design decisions: the consensus protocol and quorum, the read path, sessions and leases, watches, and the storage engine limits.

Consensus: ZAB (ZooKeeper) and Raft (etcd)#

Both are leader-based replicated logs. A write goes to the leader, is appended to the leader's log, replicated to followers, and committed once a majority has fsynced it. Then it's applied to the in-memory state machine.

Diagram: Consensus: ZAB (ZooKeeper) and Raft (etcd)

Quorum math:

Voters (N)MajorityFailures toleratedWrite costUse
110LowestDev only
3211 follower RTT + fsyncDefault
431Worse than 3, same toleranceNever
532Slightly slower than 3Tier-0, or rolling upgrades while tolerating a failure
743Noticeably slowerRarely justified

Place voters in 3 separate availability zones. With 3 voters in 2 zones, losing the zone with 2 voters loses quorum.

Differences that matter for design:

AspectZooKeeper (ZAB)etcd (Raft)
Data modelHierarchical znodes (/app/leader), ≤ 1MB each (jute.maxbuffer)Flat key space with prefix ranges, MVCC with global revision; ~1.5MB max request
Default readServed locally by any server — may be stale; sync() before read for linearizabilityLinearizable by default (ReadIndex through leader); serializable option for local stale reads
Liveness primitiveSession + ephemeral znode (deleted when session expires)Lease with TTL; keys attached to a lease are deleted when it expires
Ordering primitiveSequential znodes (/lock/n-0000000042)Global revision; create_revision per key
WatchesOne-shot historically (re-register after firing); persistent recursive watches since 3.6Streaming watch from any revision still retained; no missed events if revision not compacted
Transactionsmulti() — atomic batch of opsTxn(If compare).Then(ops).Else(ops) — CAS mini-transactions
Read scalingObservers: non-voting members serving readsLearners: non-voting members (mainly for safe member addition)
EcosystemHadoop, HBase, older Kafka, Solr, Curator recipesKubernetes, CoreDNS, many cloud-native systems

The Read Path — Linearizable vs Stale#

This is the most-missed detail in interviews.

  • etcd linearizable read: follower asks the leader for the current commit index (ReadIndex), waits to apply up to it, then reads locally. Costs roughly one RTT to the leader. Guarantees you see every write that completed before your read started.
  • etcd serializable read: read local state immediately. Fast, may be stale by up to the replication lag (usually ms, but unbounded during partitions).
  • ZooKeeper read: local and possibly stale; ZooKeeper guarantees each client sees its own writes and a monotonic view (sequential consistency per session), not linearizability across clients. Call sync() first for up-to-date reads.

🎯 Staff Insight: "If a standby checks 'is there a leader?' with a stale local read, it can see an old leader key that has already been deleted. For election decisions I'd use etcd's default linearizable reads or ZooKeeper sync() + read — and still fence, because being right at read time doesn't make you right at write time."

Sessions and Leases#

ZooKeeper session:
  client heartbeats every ~timeout/3; session timeout negotiated (typ. 6-40s,
  bounded by 2x-20x tickTime; tickTime default 2000ms)
  on expiry: all ephemeral znodes of that session are deleted, watches fire

etcd lease:
  Grant(TTL=10s) -> lease_id; KeepAlive stream refreshes every ~TTL/3
  Put(key, value, lease=lease_id)
  on expiry: attached keys deleted in one revision, watchers notified

Failover time ~= TTL (detection) + watch delivery (ms) + CAS election (1 RTT)
               ~= 10-15s with TTL=10s

Choosing the TTL is a Staff decision with two victims:

TTLFailoverFalse failovers (GC pause, network blip)Who pays
3sFastFrequent — any 3s pause triggers electionOn-call chasing flapping leaders; duplicated work without fencing
10s~10–15sRareUsers see 10–15s of no-leader
30s+SlowVery rareUsers see 30s+ outage on real failure

Default: 10 seconds, heartbeat every ~3s, and fencing so a false failover is harmless.

Watches#

A watch notifies a client when a key or prefix changes. It's how 5,000 service instances learn about a new shard map in milliseconds without 5,000 polling loops.

  • ZooKeeper (classic): one-shot — after firing you must re-read and re-register, and changes between fire and re-register are collapsed. Design as "watch says something changed; re-read current state."
  • etcd: a long-lived gRPC stream from a revision. If you disconnect and reconnect with your last revision, you get every event since — unless that revision has been compacted, in which case you must re-list and restart the watch.

The watch storm: 10,000 clients watching one key; the key changes; 10,000 notifications, 10,000 immediate re-reads. With ZooKeeper's classic one-shot watches plus a naive lock recipe, this is the herd effect. The fix in lock recipes: each waiter watches only its immediate predecessor's sequential node.

Storage Limits#

LimitZooKeeperetcd
Max value / node size1MB default (jute.maxbuffer)~1.5MB request limit (--max-request-bytes)
Total dataEntire tree in heap; practical ceiling low GBsDefault quota 2GB (--quota-backend-bytes); documented suggested max 8GB
HistorySnapshots + transaction logs; purge with autopurgeMVCC keeps old revisions until compaction; then defrag reclaims disk
Throughput (writes)Order of 10K–tens of thousands/sOrder of 10K–tens of thousands/s, disk-bound
Throughput (reads)Scales with servers/observers~100K+/s serializable; linearizable costs a leader RTT

When etcd hits its quota it raises a NOSPACE alarm and becomes read-only for writes. In Kubernetes this means no pod can be created, deleted, or rescheduled.


Core Usage — "The Entire Game": Recipes Done Right#

With ZooKeeper and etcd, nobody fails on the data model. They fail on recipes — the few patterns built on top of the primitives — implemented subtly wrong. Use battle-tested libraries (Apache Curator for ZooKeeper, go.etcd.io/etcd/client/v3/concurrency for etcd) and understand what they guarantee.

Recipe 1: Leader Election with Fencing#

etcd (concurrency.Election semantics):
  lease = Grant(ttl=10s); KeepAlive(lease)
  loop:
    txn = If(CreateRevision("/election/scheduler") == 0)
          Then(Put("/election/scheduler", my_id, lease))
          Else(Get("/election/scheduler"))
    if txn.succeeded:
        fencing_token = txn.header.revision        # monotonic across leaders
        become_leader(fencing_token)
        break
    else:
        watch("/election/scheduler") until DELETE  # then retry
                                                   # (library queues waiters by
                                                   #  create_revision: no herd)

become_leader(token):
    every action on shared resources carries token
    storage: UPDATE jobs SET ... WHERE id=? AND ? >= last_token   (reject stale)
    if KeepAlive fails or lease lost: STOP immediately, do not "finish up"
Diagram: Recipe 1: Leader Election with Fencing

The pause scenario, resolved:

t=0      A is leader, token 1042
t=+1s    A enters a 15s GC pause
t=+10s   lease expires; key deleted; B wins, token 1107
t=+11s   B writes job state with token 1107 -> store last_token = 1107
t=+16s   A wakes, still believes it's leader, writes with token 1042
         -> store rejects (1042 < 1107). No corruption.
Without fencing: A's write lands. Two schedulers fired the same job.

🎯 Staff Move: "Leader election gives me a best-effort 'usually one leader.' Fencing turns it into 'never two effective writers.' If the downstream resource can't check a token — say it's an email provider — then I need idempotency keys there instead, because election alone won't stop a duplicate."

Recipe 2: Distributed Lock — Efficiency vs Correctness#

Lock purposeConsequence of two holdersTool
Efficiency (avoid duplicate cache rebuild, avoid two crawlers on one URL)Wasted workRedis SET NX PX is fine; so is etcd — pick the cheaper dependency
Correctness (one writer to a ledger, one primary per shard)Corrupted dataConsensus-backed lease plus fencing at the resource; or avoid the lock with a conditional write

ZooKeeper lock recipe (Curator InterProcessMutex): create an ephemeral sequential node /locks/job-42/lock-0000000017; if yours is the lowest, you hold the lock; otherwise watch only the next-lower node (avoids herd effect). Session expiry deletes your node, releasing the lock automatically.

Recipe 3: Configuration and Shard Maps#

key:   /config/shard-map          value: {"version": 57, "shards": {...}}  (< 100KB)
write: Txn(If(ModRevision("/config/shard-map") == 9912)
           Then(Put("/config/shard-map", new_map)))          # CAS: no lost updates
read:  clients Get once, then Watch from returned revision
       on event: apply new map atomically; keep last good map on disconnect
  • Keep values small. A 1MB shard map watched by 5,000 clients is 5GB of fan-out per change. Split into per-shard keys, or store a pointer + version and put the blob in object storage.
  • Version everything. Clients report the version they're running; the rollout is done when all report the new version.
  • Last-known-good caching means a coordination outage freezes config — it doesn't take down serving.

Recipe 4: Group Membership#

Each member writes /members/<id> attached to its lease (etcd) or as an ephemeral node (ZooKeeper). The coordinator watches the prefix. Membership changes drive rebalancing — which is exactly how pre-KRaft Kafka and many sharded systems assigned partitions.

Danger: a network blip between the fleet and the coordination cluster expires thousands of sessions at once → mass deregistration → mass rebalancing → the "cure" is worse than the blip. Mitigate with longer session timeouts for membership than for leadership, rebalance dampening (wait N seconds before acting), and caps on how many members can be removed per minute.

Recipe 5: Barriers, Counters, Queues — Mostly Don't#

ZooKeeper's docs describe queue and barrier recipes. At scale, don't build queues on a consensus store — every enqueue is a quorum fsync. Use Kafka or SQS. Sequence/ID generation via consensus works but caps at the cluster's write rate; allocate ranges (lease 10,000 IDs per call) if you must.


The Tunable Tradeoff — Consistency, Availability, and Timeouts#

Consensus systems are CP: during a partition, the minority side cannot write (and with linearizable reads, cannot read). The knobs you tune are about how fast you detect failure versus how often you falsely detect it.

etcd timing (defaults):
  --heartbeat-interval = 100ms     (leader -> followers)
  --election-timeout   = 1000ms    (follower starts election if no heartbeat)
  rule: election-timeout >= 10x heartbeat and >> round-trip time
  cross-region RTT 70ms -> heartbeat ~150-300ms, election ~1.5-3s

ZooKeeper timing:
  tickTime = 2000ms; initLimit / syncLimit in ticks
  session timeout bounded to [2 x tickTime, 20 x tickTime] = [4s, 40s] by default

App-level lease TTL (your choice): ~10s
ChoiceWhat WorksWhat BreaksWho Pays
Short election timeoutLeader failure detected in ~1sDisk stalls / GC → spurious elections, write unavailability during eachEvery client of the cluster
Long election timeoutStable leadershipReal leader failure = seconds of write unavailabilityWriters during failover
Linearizable readsCorrect decisionsLeader RTT per read; reads fail without quorumRead latency; availability in partitions
Serializable/stale readsFast, available in minority partitionDecisions on stale dataCorrectness — acceptable for config display, not election
5 voters vs 3Survives 2 failures / upgrade + failureSlightly higher write latency, more hardwarePlatform budget
Voters across regionsSurvives region lossEvery write pays cross-region RTT (~50–150ms)Every writer, always

🎯 Staff Move: "I'd keep the voting members inside one region across three zones — 1–2ms RTT, ~5ms commits — and handle region failure at a higher layer, with a per-region cluster and a documented manual failover. Stretching one quorum across regions makes every write pay 70ms to protect against an event that happens once a year."


Anti-Patterns — What Kills ZooKeeper & etcd Deployments#

1. Using It as a Database#

Storing user sessions, job payloads, or per-request state. Data grows, snapshots get large, followers take minutes to catch up after restart, and in etcd the 2GB quota trips NOSPACE — the cluster stops accepting writes. Kubernetes clusters have hit this with large numbers of Events or oversized custom resources.

2. Consensus in the Per-Request Path#

Acquiring an etcd lock per API request. Throughput caps at the cluster's write rate (~10K/s), p99 inherits fsync latency, and a 30-second election means a 30-second outage for the whole product. Fix: coordinate at the control plane (who owns shard 17?), then serve locally.

3. Locks Without Fencing#

Covered above; the most common correctness bug. Symptom in production: rare duplicate side effects that "can't happen," always correlated with GC pauses or node migrations.

4. Even-Numbered or Two-Zone Ensembles#

4 nodes tolerate 1 failure (same as 3) with slower writes. 3 nodes in 2 AZs: lose the AZ with 2, lose quorum.

5. Shared Disks and Noisy Neighbors#

etcd/ZooKeeper on the same disk as application logs or a database. A log burst raises fsync latency to 500ms, heartbeats slip, elections start, writes pause. etcd's guidance: wal_fsync_duration_seconds p99 under ~10ms; use dedicated SSDs.

6. Hand-Rolled Recipes#

"We implemented our own lock with a TTL key and a polling loop." It will have a herd effect, a race between check and set, or no fencing. Use Curator or the etcd concurrency package.

7. Watching Huge Prefixes with Many Clients#

Every client watching / or /registry/ and receiving every change. Watch fan-out CPU on the server, bandwidth to clients, and in etcd, slow watchers force the server to buffer events. Fix: watch narrow keys, use a caching proxy (etcd gRPC proxy, Kubernetes API server's watch cache), or push through a fan-out tier.

8. Forgetting Compaction and Defrag (etcd)#

Without periodic compaction, MVCC history grows until the quota. Compaction frees logical space; defragmentation returns disk space and briefly blocks the member — run it one member at a time, never on the leader first during peak.


The Technology Landscape — Head-to-Head Comparison#

DimensionZooKeeperetcdConsulEmbedded Raft (KRaft, ClickHouse Keeper, CockroachDB)DB conditional write (Postgres / DynamoDB)
ProtocolZABRaftRaft (servers) + gossip (agents)Raft inside the productThe DB's own replication
Data modelHierarchical znodesFlat KV, MVCC revisionsKV + service catalog + health checksProduct-specific metadata logRows / items
Default readsLocal, sequentially consistentLinearizableDefault mode (leader, mostly consistent); consistent and stale optionsInternalDepends (DynamoDB strongly consistent reads opt-in)
LivenessSessions + ephemeral nodesLeasesSessions + health checksInternal heartbeatsTimestamp columns / TTL you manage
RuntimeJVMGo single binaryGoIn-processExisting DB
EcosystemHadoop, HBase, Solr, legacy KafkaKubernetes, cloud-nativeService mesh, multi-DC discovery—Everywhere
Ops burdenMedium–high (JVM, session tuning)Medium (disk, compaction, defrag)MediumLow (part of the product)Zero new
Pick whenExisting Hadoop/Curator ecosystemNew systems, Kubernetes-adjacent, need watches + leasesService discovery + health across DCsYou're building a distributed productYou need one lease or leader and already run the DB

Direction of travel: systems are embedding consensus rather than depending on an external ensemble. Kafka replaced ZooKeeper with KRaft; ClickHouse built ClickHouse Keeper as a ZooKeeper-compatible replacement. The lesson: for a product, an external coordination cluster is an operational tax; for an application, a shared managed one (or your database) is usually enough.


Patterns#

Pattern 1: Control Plane / Data Plane Split (the Staff default)#

Diagram: Pattern 1: Control Plane / Data Plane Split (the Staff default)

Requests never touch etcd. If etcd loses quorum, routers keep the last shard map and storage nodes keep serving; only changes (failover, rebalancing) pause. This is the architecture behind Kubernetes (API server + etcd vs kubelet-run workloads) and most sharded databases.

Pattern 2: Singleton Worker (Cron / Scheduler / Compactor)#

N replicas campaign; one wins and runs the singleton loop, fencing its writes. Standbys are hot. Failover ≈ lease TTL. Use for job schedulers, compaction coordinators, billing-run triggers. See Distributed Job Scheduler.

Pattern 3: Per-Shard Leadership#

Instead of one leader, each of 1,000 shards has a leader lease. Spreads load, limits blast radius of a single lease expiry. Lease count × keepalive rate must fit etcd's budget: 1,000 leases refreshed every 3s ≈ 333 writes/s — fine. 1M leases is not.

Pattern 4: Config Distribution with Staged Rollout#

Config in etcd under versioned keys; a rollout controller moves a current_version pointer per cohort (canary 1% → 10% → 100%). Clients watch their cohort's pointer. Rollback = move the pointer back. The consensus store gives atomicity; the controller gives safety.

Pattern 5: Database as Coordinator (Often the Right Answer)#

-- One row per lease; fencing via epoch
UPDATE leases
SET    owner = 'worker-7', epoch = epoch + 1, expires_at = now() + interval '10 seconds'
WHERE  name = 'billing-run'
  AND  (expires_at < now() OR owner = 'worker-7')
RETURNING epoch;   -- epoch is the fencing token

One Postgres primary already provides linearizable single-row updates. For a handful of singleton jobs, this avoids a new tier-0 dependency entirely. DynamoDB conditional writes (attribute_not_exists / version checks) serve the same role on AWS. Use when you need a few leases and no watches.


Scaling#

The Numbers#

ResourceRough figureNote
Write commit latency (3 voters, same region)~2–10msDominated by fsync + 1 RTT
Write throughput~10K–50K writes/sDisk-bound; batching helps; not horizontally scalable
Linearizable read latency~1–5msLeader round trip
Serializable read throughput100K+/s across membersScales with members/observers
Recommended data size (etcd)2GB default quota; ≤ 8GB suggested maxWhole keyspace effectively memory-resident
Max value~1–1.5MBKeep values in KB
Voters3 or 5More voters = slower writes
WatchersTens of thousands with proxies/cachesFan-out CPU is the limit
Leader failover (etcd defaults)~1–3s to electPlus app lease TTL for your own election

How to Scale Anyway#

  1. Reduce writes: batch updates, lengthen keepalive intervals where safe, move chatty state elsewhere.
  2. Scale reads: serializable reads where staleness is OK; ZooKeeper observers; etcd gRPC proxy or an application cache (the Kubernetes API server's watch cache is exactly this).
  3. Split by keyspace: separate clusters per concern — Kubernetes can put high-churn Events in a separate etcd cluster from core objects.
  4. Split by region/cell: one cluster per region; nothing global unless it must be.

You cannot scale writes by adding voters. Adding voters makes writes slower.

Multi-Region#

TopologyWrite latencySurvivesUse
One cluster, one region, 3 AZs~2–10msAZ lossDefault
One cluster, voters in 3 regions~50–150ms per writeRegion lossTier-0 global metadata that changes rarely
Cluster per region + async mirror / manual failoverLocalRegion loss with manual promotionMost multi-region systems
Cluster per region, no shared stateLocalRegion loss (independent)Cell architectures

Failure Modes & Recovery#

1. Quorum Loss#

  • Symptom: All writes fail; leader elections never succeed; linearizable reads fail.
  • Root cause: Majority of voters down or partitioned — an AZ outage with bad placement, a botched rolling upgrade (restarting the second node before the first rejoined), disk full on two members.
  • Detection: etcd_server_has_leader == 0; ZooKeeper zk_server_state without a leader; client error rate Unavailable.
  • Fix: Restore members. As an absolute last resort, force a new cluster from one member's snapshot (--force-new-cluster) — accepting possible loss of the last committed writes.
  • Prevention: 3 AZs, rolling operations gated on "all members healthy," disk alerts at 60%.

2. Leader Election Storm (Flapping)#

  • Symptom: Leader changes every few seconds; write latency p99 spikes into seconds.
  • Root cause: Slow fsync (shared disk, burstable cloud volume out of credits), CPU starvation, network jitter vs a tight election timeout.
  • Detection: rate(etcd_server_leader_changes_seen_total[10m]) > 3; etcd_disk_wal_fsync_duration_seconds p99 > 10ms; etcd_disk_backend_commit_duration_seconds p99 > 25ms.
  • Fix: Move to dedicated provisioned-IOPS SSDs; raise election timeout; isolate CPU.
  • Prevention: Dedicated nodes/disks; fsync latency SLO.

3. Session Expiry Storm / Mass Deregistration#

  • Symptom: Thousands of ephemeral nodes/leases expire together; every member rebalances; downstream thrash.
  • Root cause: Network partition between fleet and ensemble, ensemble GC pause (ZooKeeper JVM), or leader election longer than session timeouts.
  • Detection: Spike in zk_ephemerals_count drop / lease revoke rate; client SessionExpired events; membership size drops > 10% in a minute.
  • Fix: Freeze rebalancing; let sessions re-establish.
  • Prevention: Rebalance dampening, max-removals-per-minute, session timeouts ≥ election time, clients that re-register gracefully.

4. Database Space Exceeded (etcd NOSPACE)#

  • Symptom: mvcc: database space exceeded; all writes rejected. In Kubernetes: nothing can be scheduled or updated.
  • Root cause: No compaction, oversized objects, a runaway controller writing on every loop.
  • Detection: etcd_mvcc_db_total_size_in_bytes / quota > 80%; alarm list shows NOSPACE.
  • Fix: Compact to current revision, defrag each member, etcdctl alarm disarm; find and stop the writer.
  • Prevention: Auto-compaction (e.g. retention 1h/5m periodic), scheduled defrag, per-client write rate monitoring.

5. Split-Brain at the Application Layer#

  • Symptom: Two instances both acting as leader; duplicate jobs, conflicting writes. Consensus cluster looks perfectly healthy.
  • Root cause: Lease expired during a pause/partition; app kept acting; resource didn't check a fencing token.
  • Detection: Duplicate side effects keyed by job/run ID; a leader_active gauge summing to > 1 across replicas.
  • Fix: Reconcile duplicates; add fencing.
  • Prevention: Fencing tokens at every shared resource; stop work immediately on keepalive failure; idempotency keys for external effects.
Diagram: 5. Split-Brain at the Application Layer

Operational Reality Matrix#

FailureDetection SignalBlast RadiusMitigationOwner
Quorum lossetcd_server_has_leader == 0All control-plane changes; data plane if wrongly coupledRestore members; last-known-good in clientsPlatform on-call
Election stormleader changes > 3/10m; fsync p99Write latency for every clientDedicated disks, tune timeoutsPlatform
Expiry stormMembership drop > 10%/minEvery service using membershipDampening, freeze rebalancesPlatform + consuming service
NOSPACEDB size > 80% quotaAll writes (Kubernetes scheduling halts)Compact, defrag, disarmPlatform; offending client team
App split-brainDuplicate run IDs, leader gauge > 1That application's dataFencingApplication team
Watch overloadSlow watcher count, server CPUCluster latency for allProxies, narrower watchesPlatform; client team

When to Use vs. Alternatives#

NeedPickWhy
Leader election / shard ownership inside your platform, watches neededetcdLeases, CAS txns, reliable watch streams, simple ops
Existing Hadoop/HBase/Solr ecosystem, Curator recipesZooKeeperAlready there; mature recipes
One or two singleton jobs, no watchesPostgres/DynamoDB conditional writeNo new tier-0 dependency
Service discovery with health checks across DCsConsul or platform-native (Kubernetes endpoints, cloud service registry)Built for availability and health, not just consistency
Efficiency locks (dedupe work)Redis SET NX PXCheap; occasional double work is acceptable
You're building a distributed database/queueEmbedded RaftDon't make your product depend on another ensemble
High-volume state, queues, sessionsNot a consensus store — DB, Kafka, RedisWrite throughput and size limits

When NOT to Use ZooKeeper or etcd#

  • Anything per request. Their write path is a quorum fsync on a shared cluster.
  • More than a few GB of data or values larger than ~100KB.
  • Queues and work distribution at volume — use Kafka or SQS.
  • When you only need one lease and already run a strongly consistent database.
  • Discovery for tens of thousands of short-lived endpoints where availability during partition beats precision.
Diagram: When NOT to Use ZooKeeper or etcd

Operational Concerns#

What the On-Call Actually Does#

  1. Watches three numbers: has-leader, leader changes per 10 minutes, and fsync p99. Almost every etcd incident shows up in one of them first.
  2. Does rolling operations one member at a time, verifying endpoint health and that the member caught up before touching the next. Membership changes are add learner → promote, never add a voter directly into a struggling cluster.
  3. Runs compaction and defrag on schedule; defrag one member at a time outside peak.
  4. Takes snapshots (etcdctl snapshot save) at least hourly to object storage and tests restore quarterly. For Kubernetes, the etcd snapshot is the cluster backup.
  5. Finds the noisy client — per-client request rates, the biggest keys, the widest watches — because most overload is one misbehaving controller.
  6. For ZooKeeper: tunes JVM heap and GC (pauses > session timeout expire every session), enables autopurge for snapshots/logs, keeps the transaction log on its own disk.

Key Metrics & Alerts#

MetricHealthyAlert
etcd_server_has_leader10 for > 10s (page)
etcd_server_leader_changes_seen_totalRare> 3 in 10m
etcd_disk_wal_fsync_duration_seconds p99< 10ms> 25ms
etcd_disk_backend_commit_duration_seconds p99< 25ms> 100ms
etcd_mvcc_db_total_size_in_bytes / quota< 50%> 80%
etcd_server_proposals_failed_total rate0Sustained increase
etcd_network_peer_round_trip_time_seconds p99< 10ms same region> 50ms
ZooKeeper zk_avg_latency / outstanding requests< 10ms / < 10> 100ms / > 100
Client session expirations / min~0Spike > 1% of sessions

Security and Change Control#

Mutual TLS between members and clients, RBAC per client prefix (etcd roles on key ranges, ZooKeeper ACLs per znode). A consensus store with write access for everyone is a single place where one bad script can delete the shard map for the whole company.


Interview Application — Staff-Level Plays#

Which Case Studies Use ZooKeeper & etcd#

Case StudyHow It Is UsedKey Pattern
Distributed ConsensusThe primitive itself: Raft/ZAB, quorum, leasesFencing tokens, quorum placement
Distributed Job SchedulerSingleton scheduler leader, per-partition ownershipLease + fencing, idempotent job execution
Service DiscoveryRegistry with ephemeral registrations and watchesAvailability vs consistency choice
Database ShardingShard map and primary-per-shard ownershipControl plane / data plane split
Replicated Data StoreMembership and leader per replica groupPer-shard leases
Distributed CoordinationThe cross-cutting patternWhen coordination is worth its cost

Every System Design Question Has a Coordination Moment#

  • Rate limiter: "Limit policies live in etcd and push to gateways via watch; the counters never touch etcd."
  • Chat: "Which gateway owns which user connection is soft state in Redis; which node owns which partition of the message log is a lease in etcd."
  • Payment processing: "The nightly settlement run is a singleton — a lease with fencing, and the settlement table rejects stale epochs."
  • Web crawler: "URL-frontier shard ownership per crawler is a lease; losing it means another crawler picks the shard up in ~10s, and politeness state is rebuilt from storage."

What Interviewers Probe#

After You Say...They Will Ask...What They're Evaluating
"Leader election with ZooKeeper""Leader pauses 20 seconds. Then what?"Fencing tokens
"5-node etcd for safety""Why not 4? Why not 7?"Quorum math and write cost
"Store config in etcd""How big? How many watchers?"Size limits, watch fan-out
"Distributed lock""What happens if the lock expires mid-operation?"Lease semantics
"Replicate across regions""What's your write latency now?"Cross-region quorum cost
"etcd for coordination""What happens to requests if etcd is down?"Control/data plane separation

Common Interview Mistakes#

What Candidates SayWhat Interviewers HearWhat Staff Engineers Say
"The lock guarantees only one writer"Doesn't know leases expire"The lock is a lease; the resource enforces a fencing token."
"4 nodes for extra safety"No quorum math"3 tolerates 1, 5 tolerates 2; even numbers add cost, not safety."
"Read the leader key from any node"Stale-read election bugs"Linearizable read, and fence anyway."
"Store sessions in ZooKeeper"Using consensus as a database"Sessions in Redis; ZooKeeper holds metadata only."
"Every request checks the lock"Consensus on the hot path"Coordinate ownership; serve locally."
"etcd across 3 regions"Hasn't priced write latency"Per-region clusters; global quorum only for rare-change metadata."

L5 vs L6 vs L7 Responses#

ScenarioL5 AnswerL6 / Staff AnswerL7 / Principal Answer
"Ensure only one scheduler runs jobs"ZooKeeper leader electionetcd lease 10s + CAS election + revision as fencing token checked by the job store + idempotent job executionCould the job table's conditional update do this without a new dependency? If etcd exists as a platform, use the standard election library
"etcd is slow"Add nodesfsync p99, leader changes, DB size, noisy client; dedicated SSDs; compactionWhich teams are writing, and should they be here? Per-client quotas and an intake review for new consumers
"Multi-region config"One global etcdPer-region etcd, config pushed by a rollout controller; clients cache last-known-goodRegion independence as policy: no global consensus in any request or deploy path
"Migrate off ZooKeeper"Replace with etcdInventory recipes, dual-run, fencing parity, cut over per consumerDecide per system: embed Raft, move to etcd platform, or use DB CAS — and retire the ensemble on a date

The Staff ZooKeeper & etcd Checklist#

  1. Scope it: "etcd holds leader keys and the shard map — about 5,000 keys, under 50MB."
  2. Size the ensemble: "3 voters across 3 AZs, dedicated SSDs, fsync p99 under 10ms."
  3. Lease and TTL: "10-second lease, keepalive every 3 seconds, failover in ~10–15 seconds."
  4. Fence: "The lease revision is the fencing token; storage rejects older tokens."
  5. Decouple the data plane: "If etcd is down, routers keep the last shard map and keep serving; only failover pauses."
  6. Operate it: "Auto-compaction, weekly defrag, hourly snapshots with tested restore."

🎯 Staff Insight: Don't use a consensus store for data, for per-request decisions, or for locks without fencing. The strongest signal is volunteering the GC-pause scenario before the interviewer asks.

Evaluation Rubric#

DimensionSenior (L5)Staff (L6)Principal (L7)
CorrectnessLocks and electionsLeases, fencing, linearizable vs stale readsOrg-standard election library with fencing; bans hand-rolled recipes
SizingNode countQuorum math, write latency, data/value limitsNumber of clusters org-wide, consolidation plan
Failure"Replicated"Quorum loss, election storms, expiry storms, NOSPACETier-0 dependency map; cell isolation; game days
ArchitectureCoordination in the pathControl plane / data plane splitWhich systems may depend on shared consensus at all
ChoiceZooKeeper or etcdetcd vs ZK vs DB CAS vs embedded RaftBuild vs buy vs retire across the fleet

Strong hire signals

SignalWhat It Sounds Like
Fencing unprompted"A paused leader will wake up and write — the token stops it."
Quorum math"5 voters tolerate 2 failures, 4 tolerate 1."
Scoped usage"Metadata only, kilobytes per key."
Data plane independence"Losing quorum freezes changes, not traffic."
Questions the dependency"For one lease, Postgres is enough."

Lean no-hire signals

SignalWhy It Misses the Bar
Lock without lease semanticsBelieves in a guarantee that doesn't exist
Consensus per requestThroughput and availability ceiling
Even-sized or two-zone ensemblesMisunderstands quorum
Stores bulk data in itWill hit size limits and slow recovery

Common false positives

  • Reciting Raft's election rules ≠ designing with it. Can they fence?
  • "We used ZooKeeper for Kafka" ≠ coordination judgment — that was the product's choice, not theirs.
  • Knowing Paxos ≠ knowing when you don't need consensus.

The Principal Lens#

Why L7 Sees This Problem Differently#

At Staff level a consensus store is a correctly configured dependency. At Principal level it is a tier-0 blast-radius amplifier: every system that depends on it inherits its failure modes, and a shared ensemble couples teams that otherwise share nothing. Organizations accumulate ZooKeeper ensembles — one for Kafka, one for HBase, one for the scheduler, one someone set up in 2017 — each with its own half-understood runbook. The L7 question is "how few consensus clusters can we run, who is qualified to be on call for them, and which systems shouldn't depend on one at all?"

The Org-Level Fault Line#

One shared coordination platform vs. per-system ensembles vs. no external consensus at all.

OptionWhat WorksWhat BreaksWho Pays
Shared etcd platformOne expert team, standard libraries, good opsCorrelated failure: one noisy client or bad upgrade hits every dependent systemEvery team during a platform incident
Per-system ensemblesIsolationDozens of clusters run by non-experts; inconsistent versions and backupsEach team's on-call; security
Embedded consensus (KRaft-style)Product owns its fate; no external dependencyEvery product team must be Raft-competentProduct engineering
DB conditional writesZero new infrastructureNo watches; DB load; limited to simple leasesDB owners (lightly)

The Principal default: a small number of cell-scoped coordination clusters run by the platform team (one per region or failure cell, never a single global one), a standard election/lock library with fencing built in, DB-CAS for simple singletons, and embedded consensus only for teams building storage/streaming products.

Cost Model#

Assumptions: dedicated 3-voter clusters on 4 vCPU / 16GB nodes with provisioned SSDs ($300–500/node/month), snapshots to object storage, loaded engineer ~$250K/yr. The dominant cost is people, not machines.

ScaleClustersInfra/monthHeadcountOn-call load
Startup0–1 (managed K8s hides etcd)~$0–1.5K~0 (use DB CAS / managed)Near zero
Growth3–6 (per region + per concern)~$5–10K~1 FTE share of platformMonthly incidents, mostly disk/quota
Enterprise20–50 (incl. legacy ZooKeeper)~$30–75K3–5 FTE with real consensus expertiseWeekly; high severity when it happens

The expensive line is invisible: an etcd quorum loss that stops a Kubernetes control plane can freeze deploys for every team for an hour. At 500 engineers, that's ~500 engineer-hours per incident.

The 3-Year Evolution Path#

Diagram: The 3-Year Evolution Path

One-Way Doors vs Two-Way Doors#

DecisionReversibilityCost to Reverse
Putting consensus in the request pathOne-way-ish — performance and availability assumptions spread through codeRe-architect to control/data plane split
Global (cross-region) quorum for a systemOne-way-ishSplit into per-region state; data model changes
Recipe semantics clients depend on (ephemeral nodes, sequential ordering)One-way once many clients use themEvery client rewritten during migration
ZooKeeper vs etcd for a new systemTwo-way early, one-way laterCheap before launch, expensive after
Ensemble size (3 → 5)Two-wayLearner add + promote
Lease TTLs, timeoutsTwo-wayConfig

🧭 Principal Move: "The door I'd guard is 'no consensus in the request path and no global quorum.' Whether a team picks etcd or ZooKeeper matters much less than whether their users can be served while the coordination cluster is down."

The Standard I'd Write#

RFC-COORD-001: Use of Coordination Services

Scope: All services using ZooKeeper, etcd, Consul, or DB-based leases for
leadership, locking, membership, or configuration.

MUST
  1. Not call the coordination service on the per-request serving path.
  2. Continue serving with last-known-good state when the coordination service
     is unavailable; only control-plane changes may pause.
  3. Use the platform election/lock library; locks guarding correctness MUST
     pass a fencing token that the protected resource verifies.
  4. Keep total data per client namespace < 100MB and values < 100KB.
  5. Use cell/region-scoped clusters; cross-region quorums require architecture
     review approval.
SHOULD
  6. Prefer a DB conditional write for <= 3 singleton leases with no watches.
  7. Use narrow watches; wide-prefix watchers go through a caching proxy.

Exceptions: Architecture review with platform lead sign-off; time-boxed;
recorded in the dependency registry.

Operations (platform): 3 or 5 voters across 3 AZs, dedicated SSDs, fsync p99
< 10ms, hourly snapshots, quarterly restore drill, twice-yearly quorum-loss
game day.

Success metrics: zero split-brain incidents; zero request-path outages caused
by coordination unavailability; count of unowned ensembles -> 0 by year 2.

What I'd Tell the VP#

"A handful of small clusters decide who is in charge across our systems. When one fails, whole categories of work stop — last quarter a full disk on one froze deploys for every team for 50 minutes. I'm proposing we run fewer of them, run them well, and make sure customer traffic keeps flowing even when they're down, so a failure delays changes instead of causing an outage. It costs about one engineer's time on the platform team and retires four clusters nobody really owns. The payoff is removing a category of company-wide incident."

Principal Interview Signals#

SignalWhat It Sounds Like
Dependency skepticism"Do we need consensus here, or a conditional write?"
Blast-radius mapping"Which systems stop if this cluster loses quorum?"
Standardizes the dangerous part"The election library does fencing, so no team can forget it."
Retires infrastructure"Kafka on KRaft lets us delete three ZooKeeper ensembles."
Cell thinking"Coordination is regional; nothing global in the request or deploy path."

Staff answers that L7 interviewers find insufficient:

  • "We'll run a 5-node etcd" — without asking how many other clusters exist and who's on call for them.
  • "Use fencing tokens" — as a per-system fix, not a library the whole org uses.
  • "Migrate from ZooKeeper to etcd" — without per-consumer recipe inventory or a retirement date.

In the Wild#

Google — Chubby#

Google's Chubby (described in a 2006 paper) is a Paxos-based lock service with a small file-system-like interface, used for coarse-grained locking, leader election, and storing small metadata for systems like GFS and Bigtable. The paper explicitly emphasized coarse-grained locks held for hours, not fine-grained per-operation locking, and introduced "sequencers" — essentially fencing tokens — for exactly the paused-holder problem.

Staff insight: The original production lock service already said: coarse-grained only, and fence your writes. Citing that makes the fencing argument land.

Kubernetes — etcd as the Cluster's Source of Truth#

Every Kubernetes object — pods, deployments, secrets, config maps — lives in etcd, accessed only through the API server, which maintains a watch cache so thousands of controllers and kubelets don't watch etcd directly. The recommended practices — dedicated etcd with fast disks, separate etcd for high-churn Events in large clusters, regular snapshots — are the same practices described above. Running workloads keep serving if etcd is down; only changes stop.

Staff insight: Kubernetes is the canonical control-plane/data-plane split. When your design needs coordination, draw it the way Kubernetes does.

Apache Kafka — Removing ZooKeeper (KRaft)#

Kafka depended on ZooKeeper for metadata and controller election for over a decade. KIP-500 replaced it with an internal Raft-based metadata quorum (KRaft), production-ready in 3.3, with ZooKeeper mode removed in 4.0. Motivations included operating one system instead of two, faster controller failover, and supporting far more partitions.

Staff insight: A product that depends on an external consensus ensemble pays for it in ops and scale limits. If you're building the product, embed consensus; if you're using products, prefer ones that already have.


Practice Drill#

Prompt: "Design leader-based ownership for a sharded key-value store with 4,096 shards across 200 nodes. Each shard has one primary; failover should take under 30 seconds; the system must never accept conflicting writes for the same shard."

Staff Answer

Control plane vs data plane first: a 3-voter etcd cluster across 3 AZs in the region holds (a) a controller leader key and (b) the shard assignment map — 4,096 small keys, a few MB total. A shard controller runs as 3 replicas, elects a leader via a 10s lease, and assigns shard primaries by CAS on /shards/<id> (value: node, epoch). Each storage node holds one lease (10s TTL, keepalive every 3s) and its shard keys reference its node ID; the node heartbeat is ~67 lease refreshes/s cluster-wide — trivial. On node failure the lease expires (~10s), the controller sees the watch event, picks a caught-up replica for each affected shard (~20 shards), and CASes new assignments with epoch + 1. Promotion completes in ~12–20s, under the 30s bound. Conflicting writes are prevented at the storage layer, not by etcd: every replicated write carries the shard epoch, replicas reject writes from an older epoch, and a primary that loses its lease stops accepting writes immediately (and its lease TTL is shorter than the time replicas wait before accepting a new epoch). Routers cache the shard map and watch for updates; if etcd loses quorum, routing continues on the cached map and only failover pauses.

Why this is L6:

  • Separates control plane from data plane; etcd never sees a client request.
  • Epoch-based fencing at the replicas makes "never conflicting writes" true even with paused ex-primaries.
  • Quantifies lease load, data size, and failover time against the requirement.
  • Defines degraded behavior under quorum loss.

What L7 adds:

  • Asks whether this product should embed Raft per shard group (like modern distributed databases) instead of depending on an external etcd — and prices both.
  • Scopes the etcd cluster to one region/cell and forbids a cross-region quorum in the failover path.
  • Makes epoch fencing part of the storage engine's contract so future features can't bypass it.

Quick Reference Card#

Role:              linearizable metadata: leaders, shard maps, config, membership
Protocols:         ZooKeeper = ZAB, etcd = Raft; write = majority fsync
Voters:            3 (tolerate 1) or 5 (tolerate 2); never even; 3 AZs
Write latency:     ~2-10ms same region; ~50-150ms if quorum spans regions
Write throughput:  ~10K-50K/s; adding voters makes it slower
Reads:             etcd linearizable by default; ZooKeeper local/stale unless sync()
Data limits:       values KBs (hard ~1-1.5MB); etcd quota 2GB default, <= 8GB
Liveness:          ZK sessions + ephemeral nodes; etcd leases (TTL ~10s, keepalive TTL/3)
etcd timing:       heartbeat 100ms, election timeout 1000ms (defaults)
Disk:              dedicated SSD; wal fsync p99 < 10ms
Maintenance:       auto-compaction, defrag one member at a time, hourly snapshots
Locks:             a lock is a lease; fence with zxid / revision / epoch
Architecture:      coordinate in control plane; data plane serves on last-known-good

RED FLAGS
  - Locks without fencing tokens
  - Consensus call per request
  - 4 voters, or 3 voters in 2 zones
  - Sessions, queues, or blobs stored in ZooKeeper/etcd
  - One global quorum across regions for everything
  - Hand-rolled lock recipes
  1. Loading the index…