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#
| Behavior | Senior (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?" |
| Locks | Acquire lock, do work, release | Lease + fencing token checked by the resource; locks for efficiency vs locks for correctness | Standard library for leader election with fencing baked in; bans hand-rolled locks |
| Data | Store config and state in it | Metadata only: KBs per key, < 1–2GB total, watches not polling | Tracks 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 lost | Consensus cluster is a tier-0 dependency: cell-per-region, blast radius mapping, game days |
| Ownership | Infra team runs it | Platform owns ensemble; clients own session handling and correct recipes | Decides 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#
| Intent | Constraint | Strategy | Failure Mode | Correctness Bar |
|---|---|---|---|---|
| Leader election / mutual exclusion | Exactly one active actor per role; fast failover | Lease/ephemeral node + CAS on a leader key + fencing tokens | Split brain from paused ex-leader | Never two effective writers |
| Configuration & metadata store (shard maps, feature config, schema versions) | Linearizable, versioned, watched by many clients | Keys with revisions, watches, CAS updates, small values | Watch storms, oversized values, stale reads from followers | Every client converges on the same version |
| Membership & service discovery | Liveness in seconds, thousands of clients | Ephemeral nodes / leased keys + watches, or a dedicated discovery system | Session expiry storm on network blip deregisters the fleet | Eventually 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#
| Position | Rationale |
|---|---|
| Metadata only; kilobytes per key, < ~1–2GB total | Every 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 token | Pauses longer than the TTL are inevitable at scale. |
| 3 voters by default, 5 for tier-0; never even numbers | 4 tolerates the same 1 failure as 3 with slower writes. |
| Keep it off the per-request path | Consensus write latency is ~2–10ms and the cluster is shared; the data plane must survive its outage. |
| Watches, not polling | 10K clients polling every second is 10K reads/s of pure waste. |
| Dedicated low-latency disks | fsync latency is leader-election latency; a noisy neighbor on the disk causes elections. |
| Clients must tolerate unavailability | Cache 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.
Quorum math:
| Voters (N) | Majority | Failures tolerated | Write cost | Use |
|---|---|---|---|---|
| 1 | 1 | 0 | Lowest | Dev only |
| 3 | 2 | 1 | 1 follower RTT + fsync | Default |
| 4 | 3 | 1 | Worse than 3, same tolerance | Never |
| 5 | 3 | 2 | Slightly slower than 3 | Tier-0, or rolling upgrades while tolerating a failure |
| 7 | 4 | 3 | Noticeably slower | Rarely 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:
| Aspect | ZooKeeper (ZAB) | etcd (Raft) |
|---|---|---|
| Data model | Hierarchical znodes (/app/leader), ≤ 1MB each (jute.maxbuffer) | Flat key space with prefix ranges, MVCC with global revision; ~1.5MB max request |
| Default read | Served locally by any server — may be stale; sync() before read for linearizability | Linearizable by default (ReadIndex through leader); serializable option for local stale reads |
| Liveness primitive | Session + ephemeral znode (deleted when session expires) | Lease with TTL; keys attached to a lease are deleted when it expires |
| Ordering primitive | Sequential znodes (/lock/n-0000000042) | Global revision; create_revision per key |
| Watches | One-shot historically (re-register after firing); persistent recursive watches since 3.6 | Streaming watch from any revision still retained; no missed events if revision not compacted |
| Transactions | multi() — atomic batch of ops | Txn(If compare).Then(ops).Else(ops) — CAS mini-transactions |
| Read scaling | Observers: non-voting members serving reads | Learners: non-voting members (mainly for safe member addition) |
| Ecosystem | Hadoop, HBase, older Kafka, Solr, Curator recipes | Kubernetes, 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:
| TTL | Failover | False failovers (GC pause, network blip) | Who pays |
|---|---|---|---|
| 3s | Fast | Frequent — any 3s pause triggers election | On-call chasing flapping leaders; duplicated work without fencing |
| 10s | ~10–15s | Rare | Users see 10–15s of no-leader |
| 30s+ | Slow | Very rare | Users 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#
| Limit | ZooKeeper | etcd |
|---|---|---|
| Max value / node size | 1MB default (jute.maxbuffer) | ~1.5MB request limit (--max-request-bytes) |
| Total data | Entire tree in heap; practical ceiling low GBs | Default quota 2GB (--quota-backend-bytes); documented suggested max 8GB |
| History | Snapshots + transaction logs; purge with autopurge | MVCC keeps old revisions until compaction; then defrag reclaims disk |
| Throughput (writes) | Order of 10K–tens of thousands/s | Order 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"
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 purpose | Consequence of two holders | Tool |
|---|---|---|
| Efficiency (avoid duplicate cache rebuild, avoid two crawlers on one URL) | Wasted work | Redis SET NX PX is fine; so is etcd — pick the cheaper dependency |
| Correctness (one writer to a ledger, one primary per shard) | Corrupted data | Consensus-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
| Choice | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Short election timeout | Leader failure detected in ~1s | Disk stalls / GC → spurious elections, write unavailability during each | Every client of the cluster |
| Long election timeout | Stable leadership | Real leader failure = seconds of write unavailability | Writers during failover |
| Linearizable reads | Correct decisions | Leader RTT per read; reads fail without quorum | Read latency; availability in partitions |
| Serializable/stale reads | Fast, available in minority partition | Decisions on stale data | Correctness — acceptable for config display, not election |
| 5 voters vs 3 | Survives 2 failures / upgrade + failure | Slightly higher write latency, more hardware | Platform budget |
| Voters across regions | Survives region loss | Every 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#
| Dimension | ZooKeeper | etcd | Consul | Embedded Raft (KRaft, ClickHouse Keeper, CockroachDB) | DB conditional write (Postgres / DynamoDB) |
|---|---|---|---|---|---|
| Protocol | ZAB | Raft | Raft (servers) + gossip (agents) | Raft inside the product | The DB's own replication |
| Data model | Hierarchical znodes | Flat KV, MVCC revisions | KV + service catalog + health checks | Product-specific metadata log | Rows / items |
| Default reads | Local, sequentially consistent | Linearizable | Default mode (leader, mostly consistent); consistent and stale options | Internal | Depends (DynamoDB strongly consistent reads opt-in) |
| Liveness | Sessions + ephemeral nodes | Leases | Sessions + health checks | Internal heartbeats | Timestamp columns / TTL you manage |
| Runtime | JVM | Go single binary | Go | In-process | Existing DB |
| Ecosystem | Hadoop, HBase, Solr, legacy Kafka | Kubernetes, cloud-native | Service mesh, multi-DC discovery | — | Everywhere |
| Ops burden | Medium–high (JVM, session tuning) | Medium (disk, compaction, defrag) | Medium | Low (part of the product) | Zero new |
| Pick when | Existing Hadoop/Curator ecosystem | New systems, Kubernetes-adjacent, need watches + leases | Service discovery + health across DCs | You're building a distributed product | You 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)#
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#
| Resource | Rough figure | Note |
|---|---|---|
| Write commit latency (3 voters, same region) | ~2–10ms | Dominated by fsync + 1 RTT |
| Write throughput | ~10K–50K writes/s | Disk-bound; batching helps; not horizontally scalable |
| Linearizable read latency | ~1–5ms | Leader round trip |
| Serializable read throughput | 100K+/s across members | Scales with members/observers |
| Recommended data size (etcd) | 2GB default quota; ≤ 8GB suggested max | Whole keyspace effectively memory-resident |
| Max value | ~1–1.5MB | Keep values in KB |
| Voters | 3 or 5 | More voters = slower writes |
| Watchers | Tens of thousands with proxies/caches | Fan-out CPU is the limit |
| Leader failover (etcd defaults) | ~1–3s to elect | Plus app lease TTL for your own election |
How to Scale Anyway#
- Reduce writes: batch updates, lengthen keepalive intervals where safe, move chatty state elsewhere.
- 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).
- Split by keyspace: separate clusters per concern — Kubernetes can put high-churn Events in a separate etcd cluster from core objects.
- 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#
| Topology | Write latency | Survives | Use |
|---|---|---|---|
| One cluster, one region, 3 AZs | ~2–10ms | AZ loss | Default |
| One cluster, voters in 3 regions | ~50–150ms per write | Region loss | Tier-0 global metadata that changes rarely |
| Cluster per region + async mirror / manual failover | Local | Region loss with manual promotion | Most multi-region systems |
| Cluster per region, no shared state | Local | Region 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; ZooKeeperzk_server_statewithout a leader; client error rateUnavailable. - 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_secondsp99 > 10ms;etcd_disk_backend_commit_duration_secondsp99 > 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_countdrop / lease revoke rate; clientSessionExpiredevents; 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 showsNOSPACE. - 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_activegauge 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.
Operational Reality Matrix#
| Failure | Detection Signal | Blast Radius | Mitigation | Owner |
|---|---|---|---|---|
| Quorum loss | etcd_server_has_leader == 0 | All control-plane changes; data plane if wrongly coupled | Restore members; last-known-good in clients | Platform on-call |
| Election storm | leader changes > 3/10m; fsync p99 | Write latency for every client | Dedicated disks, tune timeouts | Platform |
| Expiry storm | Membership drop > 10%/min | Every service using membership | Dampening, freeze rebalances | Platform + consuming service |
| NOSPACE | DB size > 80% quota | All writes (Kubernetes scheduling halts) | Compact, defrag, disarm | Platform; offending client team |
| App split-brain | Duplicate run IDs, leader gauge > 1 | That application's data | Fencing | Application team |
| Watch overload | Slow watcher count, server CPU | Cluster latency for all | Proxies, narrower watches | Platform; client team |
When to Use vs. Alternatives#
| Need | Pick | Why |
|---|---|---|
| Leader election / shard ownership inside your platform, watches needed | etcd | Leases, CAS txns, reliable watch streams, simple ops |
| Existing Hadoop/HBase/Solr ecosystem, Curator recipes | ZooKeeper | Already there; mature recipes |
| One or two singleton jobs, no watches | Postgres/DynamoDB conditional write | No new tier-0 dependency |
| Service discovery with health checks across DCs | Consul or platform-native (Kubernetes endpoints, cloud service registry) | Built for availability and health, not just consistency |
| Efficiency locks (dedupe work) | Redis SET NX PX | Cheap; occasional double work is acceptable |
| You're building a distributed database/queue | Embedded Raft | Don't make your product depend on another ensemble |
| High-volume state, queues, sessions | Not a consensus store — DB, Kafka, Redis | Write 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.
Operational Concerns#
What the On-Call Actually Does#
- Watches three numbers: has-leader, leader changes per 10 minutes, and fsync p99. Almost every etcd incident shows up in one of them first.
- Does rolling operations one member at a time, verifying
endpoint healthand that the member caught up before touching the next. Membership changes are add learner → promote, never add a voter directly into a struggling cluster. - Runs compaction and defrag on schedule; defrag one member at a time outside peak.
- Takes snapshots (
etcdctl snapshot save) at least hourly to object storage and tests restore quarterly. For Kubernetes, the etcd snapshot is the cluster backup. - Finds the noisy client — per-client request rates, the biggest keys, the widest watches — because most overload is one misbehaving controller.
- For ZooKeeper: tunes JVM heap and GC (pauses > session timeout expire every session), enables
autopurgefor snapshots/logs, keeps the transaction log on its own disk.
Key Metrics & Alerts#
| Metric | Healthy | Alert |
|---|---|---|
etcd_server_has_leader | 1 | 0 for > 10s (page) |
etcd_server_leader_changes_seen_total | Rare | > 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 rate | 0 | Sustained 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 | ~0 | Spike > 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 Study | How It Is Used | Key Pattern |
|---|---|---|
| Distributed Consensus | The primitive itself: Raft/ZAB, quorum, leases | Fencing tokens, quorum placement |
| Distributed Job Scheduler | Singleton scheduler leader, per-partition ownership | Lease + fencing, idempotent job execution |
| Service Discovery | Registry with ephemeral registrations and watches | Availability vs consistency choice |
| Database Sharding | Shard map and primary-per-shard ownership | Control plane / data plane split |
| Replicated Data Store | Membership and leader per replica group | Per-shard leases |
| Distributed Coordination | The cross-cutting pattern | When 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 Say | What Interviewers Hear | What 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#
| Scenario | L5 Answer | L6 / Staff Answer | L7 / Principal Answer |
|---|---|---|---|
| "Ensure only one scheduler runs jobs" | ZooKeeper leader election | etcd lease 10s + CAS election + revision as fencing token checked by the job store + idempotent job execution | Could 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 nodes | fsync p99, leader changes, DB size, noisy client; dedicated SSDs; compaction | Which teams are writing, and should they be here? Per-client quotas and an intake review for new consumers |
| "Multi-region config" | One global etcd | Per-region etcd, config pushed by a rollout controller; clients cache last-known-good | Region independence as policy: no global consensus in any request or deploy path |
| "Migrate off ZooKeeper" | Replace with etcd | Inventory recipes, dual-run, fencing parity, cut over per consumer | Decide per system: embed Raft, move to etcd platform, or use DB CAS — and retire the ensemble on a date |
The Staff ZooKeeper & etcd Checklist#
- Scope it: "etcd holds leader keys and the shard map — about 5,000 keys, under 50MB."
- Size the ensemble: "3 voters across 3 AZs, dedicated SSDs, fsync p99 under 10ms."
- Lease and TTL: "10-second lease, keepalive every 3 seconds, failover in ~10–15 seconds."
- Fence: "The lease revision is the fencing token; storage rejects older tokens."
- Decouple the data plane: "If etcd is down, routers keep the last shard map and keep serving; only failover pauses."
- 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#
| Dimension | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| Correctness | Locks and elections | Leases, fencing, linearizable vs stale reads | Org-standard election library with fencing; bans hand-rolled recipes |
| Sizing | Node count | Quorum math, write latency, data/value limits | Number of clusters org-wide, consolidation plan |
| Failure | "Replicated" | Quorum loss, election storms, expiry storms, NOSPACE | Tier-0 dependency map; cell isolation; game days |
| Architecture | Coordination in the path | Control plane / data plane split | Which systems may depend on shared consensus at all |
| Choice | ZooKeeper or etcd | etcd vs ZK vs DB CAS vs embedded Raft | Build vs buy vs retire across the fleet |
Strong hire signals
| Signal | What 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
| Signal | Why It Misses the Bar |
|---|---|
| Lock without lease semantics | Believes in a guarantee that doesn't exist |
| Consensus per request | Throughput and availability ceiling |
| Even-sized or two-zone ensembles | Misunderstands quorum |
| Stores bulk data in it | Will 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.
| Option | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Shared etcd platform | One expert team, standard libraries, good ops | Correlated failure: one noisy client or bad upgrade hits every dependent system | Every team during a platform incident |
| Per-system ensembles | Isolation | Dozens of clusters run by non-experts; inconsistent versions and backups | Each team's on-call; security |
| Embedded consensus (KRaft-style) | Product owns its fate; no external dependency | Every product team must be Raft-competent | Product engineering |
| DB conditional writes | Zero new infrastructure | No watches; DB load; limited to simple leases | DB 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.
| Scale | Clusters | Infra/month | Headcount | On-call load |
|---|---|---|---|---|
| Startup | 0–1 (managed K8s hides etcd) | ~$0–1.5K | ~0 (use DB CAS / managed) | Near zero |
| Growth | 3–6 (per region + per concern) | ~$5–10K | ~1 FTE share of platform | Monthly incidents, mostly disk/quota |
| Enterprise | 20–50 (incl. legacy ZooKeeper) | ~$30–75K | 3–5 FTE with real consensus expertise | Weekly; 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#
One-Way Doors vs Two-Way Doors#
| Decision | Reversibility | Cost to Reverse |
|---|---|---|
| Putting consensus in the request path | One-way-ish — performance and availability assumptions spread through code | Re-architect to control/data plane split |
| Global (cross-region) quorum for a system | One-way-ish | Split into per-region state; data model changes |
| Recipe semantics clients depend on (ephemeral nodes, sequential ordering) | One-way once many clients use them | Every client rewritten during migration |
| ZooKeeper vs etcd for a new system | Two-way early, one-way later | Cheap before launch, expensive after |
| Ensemble size (3 → 5) | Two-way | Learner add + promote |
| Lease TTLs, timeouts | Two-way | Config |
🧭 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#
| Signal | What 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