Hiring BarSupport

Design a Consensus / Coordination Service — Staff-Level Case Study

Case study70 min read10 diagrams

Technologies referenced in this case study: ZooKeeper & etcd · Apache Kafka · DynamoDB · PostgreSQL · Redis

Related: Distributed Coordination · Service Discovery · Replicated Data Store · Distributed Job Scheduler · Consistency Models · Build vs Buy

How to Use This Case Study#

Organized for interview use first, reference second. The algorithm is the smallest part of this interview; the fault lines and the "should we run this at all" judgment are the largest.

ModeTimeWhat to Read
Quick Review15 minExecutive Summary → Interview Walkthrough → Fault Line 2 (Leases vs Fencing) → Drills 1–3
Targeted Study1–2 hrsExecutive Summary → Walkthrough → Section 3 (Fault Lines) → Section 4 (Failure Modes) → Deep Dives 1–2
Deep Dive3+ hrsEverything, including Section 11 (Principal Lens) and Appendix A (Raft mechanics)
What is a Consensus / Coordination Service? — Why interviewers pick this topic

A coordination service lets a group of machines agree on a small amount of critical state — who the leader is, who holds a lock, what the current configuration is, which nodes are alive — even when some of those machines crash or the network drops messages. Under the hood it's a replicated log kept consistent by a consensus protocol (Raft, Multi-Paxos, ZAB). On the outside it looks like a tiny, very reliable key-value store with watches, sessions and compare-and-swap. Chubby, ZooKeeper, etcd and Consul are the canonical examples.

Before vs After — the "two leaders" scenario:

Without consensus-backed election + fencing:
t=0:       Scheduler A is leader (it holds a Redis key with a 10s TTL).
t=+2s:     A enters a 15s stop-the-world GC pause.
t=+12s:    Key expires. Scheduler B acquires it and becomes leader.
t=+17s:    A resumes. Its local state says "I'm leader". It never re-checked.
t=+17s:    A and B both dispatch the nightly billing job.
t=+18s:    14,000 customers invoiced twice. Finance on the phone by 9am.

With consensus-backed lease + fencing token:
t=0:       A holds leadership with fencing token 41 (a monotonically increasing epoch).
t=+2s:     A pauses for 15s.
t=+12s:    Lease expires. B elected with token 42. B writes to the job store with token 42.
t=+17s:    A resumes and tries to dispatch with token 41.
t=+17s:    Job store rejects: "token 41 < highest seen 42". A steps down.
t=+17s:    Zero duplicate invoices. One log line and a metric increment.

Why interviewers reach for this question: It tests whether you understand that distributed agreement is possible but expensive and subtle — and that the subtle part is almost never the protocol. It's clocks, pauses, quorum placement, membership changes, and the fact that the coordination service becomes the most critical shared dependency in the company. The best answer often includes "we should not build this."

Mechanics Refresher: Consensus Protocols and Primitives
Protocol / PrimitiveHow It WorksProsCons
RaftElected leader appends entries to a log; committed once a majority persists them; terms prevent stale leadersUnderstandable; widely implemented (etcd, Consul, CockroachDB, TiKV, Kafka KRaft)Leader is a throughput bottleneck; every write costs a majority fsync + RTT
Multi-PaxosStable proposer skips phase 1 for successive slots; acceptors vote by majorityTheoretically flexible; foundation of Chubby, SpannerMany unspecified details; hard to implement correctly
ZABZooKeeper's atomic broadcast: leader-based, primary orderMature, 15+ years in productionTied to ZooKeeper's design and JVM operations
LeaseTime-bounded grant of a right (leadership, lock)Avoids a round trip per operationCorrectness depends on bounded clock drift and pause length
Fencing tokenMonotonic number issued with each lease; resources reject older tokensMakes stale holders harmlessEvery protected resource must check it
Session + ephemeral nodeClient heartbeats a session; its ephemeral keys vanish when the session expiresLiveness detection for membership and locksSession timeout tuning trades false positives vs detection time
WatchClient is notified when a key changesPush-based config and electionHerd effect when thousands watch one key

For most production systems: don't implement any of these — run etcd or ZooKeeper (or use a managed equivalent), use leases with fencing tokens, keep the data small (MBs, not GBs), and keep the coordination service off the request path. The protocol is not the interview. Failure semantics and ownership are.


Executive Summary

If you only read one section, read this. Everything in the case study flows from the contrast below.

What This Interview Actually Tests#

Consensus is not a Raft-recitation question. Anyone can draw leader → followers → majority ack.

It is a correctness-boundary and blast-radius question that tests:

  • Whether you know what a consensus service actually guarantees — and the gap between "the service knows who the leader is" and "the old leader has stopped acting"
  • Whether you size quorums against the failures you actually expect (a zone, a region, a bad deploy)
  • Whether you keep consensus off the hot path and keep its data small
  • Whether you recognize that this becomes the company's most shared dependency, and design its blast radius accordingly

The key insight: Consensus gives you agreement inside the cluster. It does not give you mutual exclusion outside it. Every lock and leader election leaks at the boundary where a paused or partitioned process still believes it's in charge — and only fencing at the resource closes that gap.

The L5 vs L6 Contrast — Start Here#

BehaviorSenior (L5)Staff (L6)Principal (L7)
First moveExplains Raft leader election and log replicationAsks "What are clients using this for — election, locks, config, or as a database?" and caps the data size and write rateAsks "Why does the org need another coordination cluster? Who runs the ones we already have?"
Locks"Acquire a lock with a TTL, release when done""A lease is advisory. The protected resource must reject stale fencing tokens, or the lock is decoration"Mandates fencing in the platform lock API; locks without a fenceable resource are flagged in design review
Quorum"Use 3 or 5 nodes""5 nodes across 3 zones as 2-2-1 survives a zone loss plus one node; cross-region quorum adds 60–80ms per write — only if we need region survival"Decides per tier: zone-survivable clusters per region by default; region-survivable only for the handful of global control planes, with the latency budget signed off
Reads"Read from any node""Follower reads are stale. Linearizable reads go through ReadIndex or leader lease; watches are for config fan-out"Publishes a consistency menu (linearizable, bounded-stale, watch) and makes clients choose explicitly per call
Failure"Raft handles failover automatically""Raft handles node failure. It does not handle slow disks, GC pauses, full databases or a watch storm — those are our incidents"Runs game days on the coordination tier; tracks it as a correlated-failure risk across all dependent teams
Build vs buyDesigns a custom Raft implementation"Run etcd or ZooKeeper; building a consensus implementation is a multi-year correctness project"Retires redundant clusters; picks one supported coordination product and a managed path; owns the deprecation plan for the others
Why "locks" separates levels

L5: Designs a correct lock inside the coordination service — sequential ephemeral nodes, watches on the predecessor, session expiry releases the lock. All true. But the lock holder is a separate process, and that process can pause for longer than its session timeout (GC, VM migration, a swapped page, a SIGSTOP) and then keep writing.

L6: "The lock tells the holder it had the lock at some point. It can't tell the holder it still has it at the moment of the write. So I attach a fencing token — the lock's revision or a monotonic epoch — to every write, and the storage layer rejects any token lower than the highest it's seen. Without that, I'd call this an efficiency lock, safe only for work that's idempotent anyway."

L7: "Half the teams using locks don't need mutual exclusion for correctness — they need it to avoid wasted work. I'd split the API: efficiency_lock with no guarantees and correctness_lease that returns a fencing token and requires the caller to register which resource enforces it."

Why "quorum" separates levels

L5: "Five nodes tolerate two failures." Correct arithmetic, no placement. Five nodes in two zones (3+2) lose quorum when the zone with 3 goes down — the cluster tolerates two random node failures but not the correlated failure that actually happens.

L6: Places nodes against failure domains: "5 nodes across 3 zones as 2-2-1. Any single zone loss leaves ≥ 3. Across regions, I'd need 3 regions minimum and every write pays the second-closest region's RTT — about 60–80ms coast-to-coast. I'll only do that for state that must survive a region loss."

L7: Frames placement as an org-wide policy with a cost: "Region-survivable coordination costs ~50× the write latency and 3× the infra. I'd allow it for three global control planes — identity, global config, DNS control — and require everyone else to be region-scoped with async replication of what they need."

Why "failure" separates levels

L5: "Raft elects a new leader in a few hundred milliseconds." True for a crashed leader on a healthy cluster.

L6: Knows the real incidents: "The common etcd outages aren't crashes — they're a leader on a slow disk whose fsync p99 goes to 500ms, causing heartbeat misses and leader flapping every few seconds; or the database hitting its size quota and going read-only; or 10,000 clients reconnecting at once after a network blip. Each has a metric and a runbook."

L7: Recognizes correlated risk: "When etcd degrades, the Kubernetes API degrades, which means deploys stop, which means the fix for whatever caused it can't ship. The coordination tier needs a break-glass path that doesn't depend on itself."

The Staff Positions#

PositionRationale
Buy, don't build, the consensus layeretcd, ZooKeeper and Consul encode a decade of fixed edge cases; a new implementation needs formal methods and years of soak
Leases always ship with fencing tokensA lease without a fence protects against crashed holders, not paused ones
Keep it small: MBs of data, hundreds–low thousands of writes/sEvery write is a majority fsync; the whole dataset lives in memory and in every snapshot
Off the hot pathRequest paths read cached config or leader identity; a coordination blip must not fail user requests
5 nodes across 3 zones by default; 3 for non-criticalSurvives a zone loss; 7+ adds write latency for little gain
Linearizable reads only where neededReadIndex/leader reads cost a round trip; config fan-out uses watches + cached revisions
One coordination platform, many tenants — with per-tenant quotasTen teams running ten ZooKeepers is ten on-call rotations and ten upgrade backlogs

The Three Intents#

IntentConstraintStrategyFailure ModeCorrectness Bar
Coordination primitives (leader election, locks, membership)Safety: at most one actor effective at a timeSessions/leases + fencing tokens; small keyspace; watchesTwo actors both effective (split brain at the resource)Linearizable writes; fence enforced at every protected resource
Configuration / metadata store (feature flags, cluster state, service registry)Read fan-out to 10K+ clients; consistent snapshotsWatch-based distribution, revisioned reads, client cachesWatch storms; stale config served; DB size growthMonotonic revisions; bounded staleness at clients
Replicated state machine for data (consensus per shard in a database)Throughput and horizontal scaleMulti-Raft: one group per range/partition, thousands of groupsHot ranges, rebalancing, cross-group transactionsLinearizable per range; transactions across ranges need extra protocol

🎯 Staff Move: "I'll design the coordination service — the Chubby/etcd shape — used for leader election, locks and config by other systems. It holds megabytes, not terabytes, and handles hundreds of writes per second, not millions. If you want consensus as the replication layer of a database, that's multi-Raft sharding, and I'd treat it as a different design."

The Five Fault Lines#

#Fault LineThe Tension
1Safety vs Liveness (Quorum Sizing & Placement)Bigger, wider quorums survive more failures but every write gets slower and elections get harder
2Leases vs FencingTrust bounded clocks and pauses for speed, or require every resource to check tokens?
3Linearizable Reads vs Read ScaleRoute reads through the leader (correct, bottlenecked) or serve from followers/caches (fast, stale)?
4Build vs Buy vs AvoidOwn a consensus implementation, run etcd/ZooKeeper, or use a database's conditional writes instead?
5Control Plane vs Data PlaneLet request paths depend on coordination (simple, fragile) or cache and degrade (robust, stale)?

In the Wild: Real Production Systems#

Why this section belongs here: These three systems define the vocabulary interviewers use. Cite them by mechanism, not by name only.

Google Chubby — Coarse-Grained Locks as a Service#

Google's Chubby paper (2006) describes a lock service built on Paxos, typically five replicas per cell, used by systems like GFS and Bigtable for master election and for storing small amounts of metadata. Chubby deliberately offers coarse-grained locks held for hours or days, not fine-grained per-request locks. Clients hold sessions maintained by KeepAlives, cache data with server-driven invalidation, and receive a sequencer with a lock that downstream servers can check — Chubby's form of a fencing token.

Staff insight: The paper's most-quoted lesson is that engineers used Chubby as a name service and small database far more than as a lock service, and that load was dominated by reads and KeepAlives. Design your coordination service for how it'll be abused, not how it's specified.

etcd in Kubernetes — The Cluster's Single Source of Truth#

Every Kubernetes object — pods, services, deployments, secrets — lives in etcd, a Raft-based key-value store. Controllers don't poll; they watch resources from a revision and reconcile. The API server fronts etcd so that thousands of clients don't hit it directly. etcd documents a default storage quota of 2 GiB (configurable, with ~8 GiB suggested as a practical maximum) and a request size limit of about 1.5 MiB; when the quota is exceeded, etcd raises a NOSPACE alarm and rejects writes until compaction and defragmentation free space.

Staff insight: Kubernetes shows the right architecture — a caching, watch-multiplexing layer (the API server) in front of consensus — and the right failure: when etcd is slow, the whole control plane is slow, but running pods keep serving. Control plane down ≠ data plane down.

Apache Kafka — Removing ZooKeeper (KRaft)#

Kafka historically stored cluster metadata (brokers, topics, partition leaders) in ZooKeeper. KIP-500 replaced it with KRaft, a Raft-based metadata quorum inside Kafka itself, and Kafka 4.0 removed ZooKeeper mode. Motivations publicly cited include operating one system instead of two, faster controller failover, and supporting far more partitions because metadata is a replicated log that brokers consume rather than state loaded from ZooKeeper on failover.

Staff insight: "Build vs buy" for consensus flipped only when coordination became core to the product's scaling limit, with a mature team and years of runway. That's the bar for building: consensus is your product, not your dependency.

What Interviewers Probe#

After You Say...They Will Ask...(What They're Evaluating)
"We'll use a distributed lock""The holder pauses for 30 seconds. What stops it from writing after its lock expired?"Fencing tokens at the resource
"5 nodes tolerate 2 failures""Where are the 5 nodes? What happens if a zone goes down?"Failure-domain placement
"Reads go to followers for scale""Can a client read a value older than one it already saw?"Linearizability vs session guarantees
"Raft elects a leader automatically""The leader's disk is slow but not dead. What happens?"Gray failure, leader flapping
"We'll build it on Raft""Why not etcd? Who maintains your Raft library in 3 years?"Build vs buy judgment
"Services look up the leader on each request""The coordination cluster is unavailable for 2 minutes. What breaks?"Control plane off the hot path

System Architecture Overview#

Diagram: System Architecture Overview

Reading the diagram: Three separations matter. The access layer coalesces watches and enforces quotas so 10,000 clients look like 50 to the core. The consensus core is small and zone-spread so any single zone loss leaves a majority. The protected resource enforces fencing tokens — that's where mutual exclusion is actually guaranteed. The metrics panel watches the things that cause real incidents: fsync latency, leader churn, DB size vs quota and watcher count.

Quick-Reference: The 30-Second Cheat Sheet#

TopicThe L5 AnswerThe L6 Answer — Say This
ProtocolExplains Raft in detail"Raft via etcd. The protocol is solved; the edges aren't."
Locks"TTL lock, release when done""Lease plus fencing token; the resource rejects stale tokens."
Quorum"5 nodes""5 nodes, 3 zones, 2-2-1. Region-survivable only if we accept +60–80ms per write."
Reads"Any node""Linearizable via ReadIndex where needed; watches + cached revision for config."
Scale"Add nodes""Adding voters slows writes. Scale reads with learners/proxies; scale writes by splitting clusters."
Failure"Automatic failover""Slow disk → leader flapping; full DB → read-only; reconnect storm → overload. Each has an alert."
Ownership"Infra runs it""One platform team, per-tenant quotas and key prefixes, and a break-glass path."

Key Numbers Worth Memorizing#

MetricValueWhy It Matters
Majority quorum⌊n/2⌋ + 13 nodes tolerate 1 failure, 5 tolerate 2, 7 tolerate 3
Raft paper election timeout150–300ms randomizedThe baseline for fast failover on a LAN
etcd defaultsheartbeat 100ms, election timeout 1,000msTune both upward for cross-zone/cross-region latency
ZooKeeper tickTime default2,000ms; session timeout 2–20 ticksSessions of 4–40s: detection speed vs false expiry
Chubby default session lease~12s (per the paper), extended by KeepAlivesCoarse-grained locking, not per-request
Disk fsync (NVMe/SSD)~0.1–2ms typical; alert when p99 > 10msEvery committed write waits on a majority fsync
Intra-zone / cross-zone RTT< 0.5ms / ~1–2msZone-spread quorum adds only a couple of ms
US cross-region RTT (east↔west)~60–80msRegion-spread quorum makes every write ≥ this
etcd storage quota2 GiB default, ~8 GiB suggested maxHit it and writes stop (NOSPACE alarm)
etcd max request size~1.5 MiBCoordination stores are not blob stores
Practical write throughput (3–5 node, SSD, small values)~10–30K writes/s with batchingA shared cluster's budget — allocate it per tenant
Kubernetes supported cluster sizeup to ~5,000 nodesetcd write and watch load is a primary limiter

Interview Walkthrough

The most common mistake: Candidates spend 20 minutes explaining Raft terms, votes and log matching — accurately — and never reach the question the interviewer is waiting for: "The old leader was paused. What stops it from doing damage?" Compress the protocol to 3 minutes and spend the rest on fencing, placement, reads and blast radius.


Phase 1: Requirements & Framing (2–3 minutes)#

State the scope in one breath:

"We need a service that lets other systems elect leaders, take locks, register membership and distribute configuration — with a guarantee that they agree even through node failures and partitions."

Then clarify the intent and cap it:

"Is this a coordination service for other systems, or the replication layer inside a database? I'll assume coordination: small data — under a few GB total — writes in the hundreds to low thousands per second, reads and watches from up to tens of thousands of clients. Correctness is non-negotiable for leader election and locks; config can be slightly stale at clients."

Then name the failure envelope:

"What must we survive — a node, a zone, or a region? I'll assume a zone, within one region, and treat region survival as a separate decision because it changes write latency by 30–50×."

🎯 Staff Move: Capping data size and write rate in the first two minutes prevents the most common drift in this interview — someone deciding to store user data in the coordination service by minute 25.


Phase 2: Core Entities & API (1–2 minutes)#

  • Key: path-like name, value (≤ 1 MiB), create_revision, mod_revision, version, optional lease_id
  • Lease / Session: lease_id, TTL (e.g., 10s), owner; keys attached to a lease disappear when it expires
  • Revision: a cluster-wide monotonic counter incremented on every write — the basis for fencing, watches and consistent snapshots
  • Watch: subscription from a revision on a key or prefix
Put(key, value, lease?, if_mod_revision?)   → { revision }          # CAS via if_mod_revision
Get(key | prefix, consistency=linearizable|serializable, at_revision?) → { kvs, revision }
Delete(key, if_mod_revision?)
Txn(compare[], then_ops[], else_ops[])       → { succeeded, revision }
LeaseGrant(ttl) → lease_id ;  LeaseKeepAlive(lease_id) (stream) ;  LeaseRevoke(lease_id)
Watch(key | prefix, from_revision)           → stream of events
# Recipes on top:
Campaign(election, lease) → { leader_key, fencing_token = create_revision }
Lock(name, lease)         → { fencing_token }

🎯 Staff Move: "The API is deliberately tiny — CAS, leases and watches. Elections and locks are client-side recipes on these primitives. The fencing token falls out for free: it's the revision at which the leader key was created, and revisions never go backwards."


Phase 3: High-Level Architecture (≤5 minutes)#

Diagram: Phase 3: High-Level Architecture (≤5 minutes)

Walk the write path in 90 seconds:

  1. Client sends Put to the proxy; proxy forwards to the current leader.
  2. Leader appends to its log, fsyncs, sends AppendEntries to 4 followers in parallel.
  3. Once 2 followers (3 of 5 including the leader) have fsynced, the entry is committed.
  4. Leader applies to its state machine, bumps the revision, responds, and notifies watchers.
  5. Latency ≈ one intra-region RTT + one fsync ≈ 2–5ms zone-spread.

Key points to hit:

  1. Five voters, three zones, 2-2-1 — survives a zone
  2. Proxy layer — clients never connect directly to consensus nodes at 10K scale
  3. Leases for liveness, fencing tokens for safety
  4. Off the hot path — clients cache; a coordination outage degrades control, not traffic
  5. Buy it — etcd or ZooKeeper; the design work is placement, quotas, clients and operations

🎯 Staff Move: After drawing: "This is the part that's well-understood. The incidents happen elsewhere — at the resource that didn't check the token, the disk that got slow, the 10,000 watchers reconnecting. That's where I'd like to go deep."


Phase 4: Transition to Depth (1 minute)#

"Three areas decide whether this is safe in production: fencing — what happens when a leader doesn't know it's been replaced; quorum placement and gray failures — slow disks and flapping leaders; and read semantics and blast radius — who can depend on this and how. Which is most interesting to you?"

Default if no preference: fencing. It's where most candidates are weakest and most real-world bugs live.


Phase 5: Deep Dives (25–30 minutes)#

Deep dive 1: Leases and fencing (7–8 min)

"A lease says 'you're leader until T' — by the server's clock. The holder checks its own clock, and it can be wrong in two ways: clock drift, which is bounded if we use monotonic clocks and a safety margin, and process pauses, which aren't bounded at all — GC, VM live migration, page faults. So I make leadership safe at the resource: each leader gets a fencing token equal to the revision of its leader key. Every write to the job store carries the token; the store keeps the max token it has seen per resource and rejects lower ones. A paused ex-leader wakes up and gets rejected on its first write."

Quantify: "Lease TTL 10s, renewal every 3s, and the holder stops acting at 10s minus a safety margin of 1–2s measured on its own monotonic clock. That handles drift. Fencing handles the pause."

Who pays: "Every team that owns a protected resource has to implement the token check. That's the cost — and I'd pay it with a shared library, not by skipping it."


Deep dive 2: Quorum placement and gray failure (6–7 min)

"2-2-1 across three zones. Losing any zone leaves ≥ 3 voters. Losing a zone plus one more node elsewhere still leaves a majority only if the lost zone was the 1 — so we survive 'one zone' reliably and 'zone plus node' partially. The gray failure I worry about more is a leader with a slow disk: fsync p99 goes from 2ms to 400ms, heartbeats stall, followers time out, elect a new leader, and the old one — no longer overloaded — rejoins. If the slow disk is on a node that keeps winning elections, we flap every few seconds. etcd mitigates with pre-vote and check-quorum; I'd add an alert on leader_changes_total > 3 in 15m and a runbook to remove the sick node."


Deep dive 3: Reads and watches (5–6 min)

"Three read modes. Linearizable: ReadIndex — the leader confirms it's still leader with a heartbeat round, then serves once applied up to that index; costs one RTT, no disk write. Serializable: any member serves from local state — may be stale, fine for dashboards. Watches: clients subscribe from a revision and receive every change in order. For config at 10K clients, watches via the proxy; for leader discovery, watches; for 'am I still leader before I act', I don't read — I fence."


Deep dive 4: Blast radius and ownership (4–5 min)

"This cluster becomes the thing everything depends on. So: per-tenant key prefixes with quotas on keys, bytes and write rate; clients that cache the last-known config and keep serving if the cluster is unavailable; separate clusters per criticality tier — the Kubernetes control plane doesn't share etcd with feature flags. The platform team owns the pager; tenants own their usage, and a tenant exceeding quota gets throttled, not the cluster."


Phase 6: Wrap-Up (2–3 minutes)#

"The protocol guarantees agreement inside the cluster. Safety outside it comes from fencing at the resource. Availability comes from placement against real failure domains, and from keeping the cluster off the request path. And the biggest risk isn't consensus at all — it's that everything ends up depending on one cluster."

The build-vs-buy closer:

"I'd run etcd, not write Raft. If our product were coordination — like Kafka's move to KRaft — we'd staff a team for years and use formal methods and deterministic simulation testing. For everyone else, the work is operations, quotas and client libraries."

🎯 Staff Move: End with "who else depends on this, and what happens to them when it's down." That's the question the interviewer wants to hear you ask.


Common Timing Mistakes#

MistakeL5 Does ThisL6 Does This Instead
Protocol lecture15 min on terms, log matching, commit index3 min on Raft, then fencing and placement
Paxos vs Raft debateCompares proofs"Raft via etcd; the choice doesn't change the design"
No failure domains"5 nodes""5 nodes, 3 zones, 2-2-1"
Locks without fencesEphemeral-node lock recipe, done"Recipe gives liveness. Fence gives safety."
Coordination on hot pathEvery request checks leaderCache + watch; degrade on outage
Unbounded dataStores job payloads in the cluster"Pointers only; payloads go in object storage or a DB"

1. The Staff Lens#

1.1 Why This Problem Exists in Staff Interviews#

Consensus is the interview where textbook knowledge is abundant and production judgment is rare. Most candidates can describe Raft; few have debugged a leader that flaps because one disk is slow, or a lock that "worked" in every test and double-processed during a 20-second GC pause. The Staff-level signal is recognizing that the protocol is a solved component and that the boundaries — clocks, pauses, placement, clients, quotas, and the organization's growing dependence on it — are the design.

It also tests restraint. The strongest answers often reduce the scope: "We don't need a lock here; a conditional write on the row is enough." Coordination is expensive — every write costs a majority fsync — and it becomes a shared single point of failure. Knowing when not to use it is part of the answer.

1.2 The L5 vs L6 Contrast — Visual#

Diagram: 1.2 The L5 vs L6 Contrast — Visual

1.3 The Staff Question That Cuts Through Everything#

"Your leader process freezes for 30 seconds — longer than any timeout you've set — and then resumes mid-operation. Walk me through exactly what prevents it from corrupting state."

A candidate who answers "the lease expired so it's not leader" has described the coordination service's view. A candidate who answers "its next write carries token 41, the store has seen 42, the write is rejected, and the process exits on the rejection" has described the system's view. Only the second is safe.


2. Problem Framing & Intent#

2.1 The Three Intents — Explained#

Coordination primitives → safety at the edges

  • Constraint: at most one effective actor per role at any time
  • Strategy: leases for liveness, fencing tokens for safety, coarse-grained locks held for seconds to hours
  • Failure mode: split-brain at the resource when fencing is absent
  • Who pays for imperfection: the team whose data got written twice

Configuration / metadata store → fan-out and freshness

  • Constraint: tens of thousands of readers; changes must propagate in seconds and in order
  • Strategy: watches from a revision via a proxy layer; clients cache and survive outages with last-known-good
  • Failure mode: watch storms on reconnect; bad config propagating to everyone in seconds
  • Who pays: every service that consumes the config — a bad flag is a global incident

Replicated state machine for data → throughput and sharding

  • Constraint: millions of keys, high write throughput, horizontal scale
  • Strategy: multi-Raft — thousands of small consensus groups (one per range), a placement driver, range splits and merges (CockroachDB, TiKV, and Spanner's Paxos groups follow this shape)
  • Failure mode: hot ranges, rebalancing storms, cross-range transactions
  • Who pays: the database team; this is a database design interview, not a coordination one

2.2 When NOT to Use a Consensus Service#

  • Idempotent work where duplicates are merely wasteful. Two workers both regenerating a cache is fine. Use a best-effort lock in Redis or none at all.
  • Mutual exclusion on a single row in one database. A conditional write (UPDATE … WHERE version = ?) or SELECT … FOR UPDATE in PostgreSQL, or a condition expression in DynamoDB, already gives you linearizable CAS and the database is the resource — no fence gap. Adding ZooKeeper here adds a failure domain for nothing.
  • Leader election for a single-writer queue consumer. Kafka's consumer group protocol already assigns partitions exclusively; don't layer another election on top.
  • High-volume data. Job payloads, user sessions, metrics — anything beyond small metadata. The coordination store keeps everything in memory and in every snapshot.
  • Rate limiting or counters on the request path. Every increment is a majority fsync; use local-first counters (Rate Limiting).
  • Cross-region config where seconds of staleness are fine. Replicate asynchronously from a regional source of truth rather than stretching a quorum across regions.

🎯 Staff Insight: "The best coordination is the coordination you don't need. If the resource itself can do a compare-and-swap, the lock and the fence collapse into one operation — and that's strictly safer than any external lock."

2.3 What the Interviewer Leaves Underspecified#

  • Which failure domain must be survived — node, zone, region
  • Who the clients are — 50 controllers or 50,000 sidecars changes the access layer entirely
  • Whether locks guard correctness or efficiency — decides whether fencing is mandatory
  • Data size and write rate — kilobytes of leader keys vs gigabytes of "metadata" someone plans to add
  • Read consistency needs — is stale config for 2 seconds acceptable?
  • Multi-tenancy — is this one team's cluster or the company's?

2.4 Precise Terminology#

TermWhat It MeansWhy It Matters
SafetyNothing bad happens: no two leaders commit in the same term; committed entries are never lostConsensus protocols guarantee safety always
LivenessSomething good eventually happens: a leader gets elected, writes completeGuaranteed only when a majority can communicate and timing is reasonable (FLP)
QuorumA majority of voters: ⌊n/2⌋ + 1Any two quorums overlap — the source of safety
Term / epochMonotonic leadership generationLets nodes reject messages from stale leaders
LeaseTime-bounded right, renewed by heartbeatsLiveness tool; unsafe alone under pauses
Fencing tokenMonotonic number checked by the protected resourceMakes stale holders harmless
LinearizableEvery operation appears to take effect at one instant between call and returnWhat "consistent" should mean for locks and elections
Serializable read (etcd sense)Served from a member's local state; possibly staleCheap; not for correctness decisions
Learner / observerNon-voting replica that receives the logScales reads without slowing writes
Joint consensusRaft's two-phase membership change using old+new majoritiesPrevents two disjoint majorities during reconfiguration
Split brainTwo nodes both acting as leaderPrevented inside consensus by quorum overlap; outside only by fencing

🎯 Staff Insight: When an interviewer says "strongly consistent", ask: "Linearizable reads, or linearizable writes with possibly stale reads? etcd and ZooKeeper differ by default — ZooKeeper reads can be stale unless you sync first."


3. The Five Fault Lines#

Each fault line has a protocol-level answer everyone knows and an operational/organizational answer that decides the level.

3.1 Fault Line 1: Safety vs Liveness — Quorum Sizing & Placement#

The tension: More voters in more failure domains survive more failures but make every write wait on a slower majority and make elections slower and more fragile. Fewer, closer voters are fast but die with a zone or region.

ChoiceWhat WorksWhat BreaksWho Pays
3 nodes, 3 zonesSurvives one zone; ~2–4ms writesZone loss leaves 2 — one more failure or a rolling upgrade = no quorumOn-call (no maintenance headroom during zone events)
5 nodes, 3 zones (2-2-1)Survives a zone; maintenance headroom~20% slower commits than 3 (waits for 3rd-fastest ack)Platform (2 more nodes to run)
7 nodesSurvives 3 random failuresWrite latency and leader load grow; rarely matches real failure domainsEveryone (latency) for little gain
5 nodes, 3 regions (2-2-1)Survives a regionEvery write ≥ second-closest-region RTT (~30–80ms); election timeouts must rise to secondsAll writers (latency), platform (tuning)
Diagram: 3.1 Fault Line 1: Safety vs Liveness — Quorum Sizing & Placement

Staff default: "Five voters, three zones, 2-2-1, one region. Region survival for coordination is usually the wrong goal — most consumers are region-scoped anyway. I'd rather have one cluster per region and replicate the few global keys asynchronously."

When to deviate:

  • Global control planes (global DNS control, identity, cross-region leader for a globally consistent database): region-spread quorum, accept the latency, keep write rates in the tens per second.
  • Tiny, non-critical clusters (a team's internal scheduler): 3 nodes; the extra two voters aren't worth the operations.
  • Read-heavy scale: add learners/observers (non-voting) rather than voters — reads scale, writes don't slow.

🧭 Principal Move: "Quorum placement is an org policy, not a per-cluster choice. I'd publish three tiers — zone-survivable (default), region-survivable (by exception, with a latency sign-off), and single-zone (dev only) — and have the platform provision only those shapes."

❌ Common L5 Trap: "Five nodes, so we survive two failures." Then five nodes placed 3+2 across two zones: lose the 3-node zone and quorum is gone. Node counts without failure-domain placement are arithmetic, not availability.


3.2 Fault Line 2: Leases vs Fencing#

The tension: Leases let a leader act without a round trip per operation, which is fast — but correctness depends on assumptions about clocks and pauses that production violates. Fencing tokens are safe regardless of timing, but every protected resource must implement a check.

ChoiceWhat WorksWhat BreaksWho Pays
Lease onlySimple; no resource changesGC pause, VM migration, clock jump → two effective leadersResource owner (corrupt/duplicate writes)
Lease + self-check before each actionNarrows the windowPause can occur between check and action — the window never closesSame, less often — the worst kind of bug
Lease + fencing token at resourceSafe under arbitrary pausesEvery resource must store max-seen token and compareResource teams (implementation), platform (library)
Resource-native CAS (no external lock)Lock and fence are one operationOnly works when the resource supports conditional writesNobody — prefer when available
Diagram: 3.2 Fault Line 2: Leases vs Fencing

Staff default: "Every correctness lease comes with a fencing token — in etcd, the create_revision of the lease-holder key; in ZooKeeper, the zxid or the sequential node number. The storage layer checks it. If the resource can't check tokens — say, a third-party API — the lock is best-effort, and I design the operation to be idempotent instead."

When to deviate:

  • Efficiency-only locks (avoid duplicate cache warming): lease only; duplicates are harmless.
  • Resource is a single database row: skip the external lock; use the row's version column as the fence.
  • Hardware-backed bounded clocks (e.g., TrueTime-style uncertainty intervals): lease reads become safe inside the database because uncertainty is bounded and waited out — this is how Spanner's leaders serve reads — but only within that system.

🧭 Principal Move: "I'd split the platform API into efficiency_lock() and correctness_lease(). The second one requires the caller to declare the fenced resource in the service catalog. Design review rejects correctness claims on locks with no declared fence."

❌ Common L5 Trap: "Before each write, the worker checks it still holds the lock." The check and the write aren't atomic; a pause between them reopens the gap. The interviewer knows this — Martin Kleppmann's widely-read 2016 critique of Redlock makes exactly this point.


3.3 Fault Line 3: Linearizable Reads vs Read Scale#

The tension: Linearizable reads must confirm the leader is still leader (a quorum round) or trust a leader lease (clock assumptions). Serving from followers or caches scales reads ~N× but returns stale data — and can even go backwards across servers.

ChoiceWhat WorksWhat BreaksWho Pays
Reads through the Raft logTrivially linearizableEvery read is a write-cost operation (fsync)Everyone (throughput collapse)
ReadIndex (leader confirms via heartbeat round)Linearizable; no disk writeOne RTT to a majority; leader CPU boundLeader (load concentration)
Leader lease readsNo extra round tripUnsafe if clocks drift beyond the lease marginCorrectness (rare stale reads)
Follower/serializable readsScale with replicasStale; non-monotonic across servers unless session-pinnedClients (must tolerate)
Watch + client cacheScales to 10K+ clients; push-basedStaleness = propagation delay (~10–500ms); reconnect stormsPlatform (proxy layer)
Diagram: 3.3 Fault Line 3: Linearizable Reads vs Read Scale

Staff default: "Correctness decisions use fencing or CAS, not reads. Humans and deploy tools use linearizable reads. Fleets use watches through a proxy with cached, revisioned state. Serializable reads for observability only."

When to deviate:

  • ZooKeeper clients that must see their own writes elsewhere: sync() before read, or pin the session to one server (ZooKeeper guarantees per-session ordering).
  • Very high read rates on a few keys: add learners in each zone; don't add voters.

🎯 Staff Insight: "The dangerous read is the one used to make a decision that then acts on a different system. By the time the decision lands, the read is history. That's why I fence rather than read."


3.4 Fault Line 4: Build vs Buy vs Avoid#

The tension: A custom consensus implementation can be tailored to your product's scaling limits. It is also one of the hardest correctness projects in software: edge cases in membership changes, snapshotting and log truncation surface years later.

ChoiceWhat WorksWhat BreaksWho Pays
Build your own Raft/PaxosTailored; no external dependencyYears to maturity; subtle bugs; needs formal methods (TLA+) and deterministic simulation testingA dedicated team of 4–8 for years
Embed a library (etcd/raft, hashicorp/raft)Proven core; your own state machineYou own snapshots, storage, transport, membership opsYour team (integration + ops)
Run etcd / ZooKeeper / ConsulMature; well-known failure modes; ecosystemOps burden (upgrades, defrag, monitoring); JVM tuning for ZooKeeperPlatform team (1–2 FTE per fleet)
Managed service / cloud primitiveNo ops; conditional writes in DynamoDB, managed etcd in Kubernetes offeringsLess control; vendor limits; lock-inBudget, and vendor risk
Avoid (resource-native CAS)Zero new infraOnly when a single store holds the stateNobody

Staff default: "Avoid if the resource has CAS. Otherwise run etcd — or use a managed offering if we're on one cloud. I'd only embed a Raft library if consensus is inside our product's data path, and I'd only build one if we're a database company."

When to deviate:

  • You're building a database or a broker where metadata scale is the product limit (Kafka's KRaft): embed or build, with a team and a multi-year plan.
  • Air-gapped or edge deployments where running etcd isn't possible: embed a library.

🧭 Principal Move: "Every coordination cluster in the company is a pager rotation, an upgrade backlog and a CVE surface. I'd count them. When I did this exercise the answer is usually 'more than the platform team knew about' — and the first project is consolidation, not new features."


3.5 Fault Line 5: Control Plane vs Data Plane#

The tension: It's simplest for services to consult the coordination service directly when they need a leader, a lock or a config value. That puts a majority-fsync system with 1–2s elections on the request path of every product.

ChoiceWhat WorksWhat BreaksWho Pays
Consult per requestAlways freshEvery election or brownout = user-facing errors; read load explodesUsers (errors), platform (load)
Cache + watch, fail staticCoordination outage degrades change, not trafficStale during outage; must define max stalenessProduct (sign off on staleness)
Snapshot to local file / sidecarSurvives long outages; boot without coordinationPropagation slower (seconds)Platform (distribution tooling)

Staff default: "Coordination is control plane. The data plane runs from cached state and keeps running — fail static — if coordination is unavailable. The only things that stop during an outage are changes: new leaders, new locks, new config. That's the right thing to lose."

When to deviate:

  • Operations that are themselves control-plane (a deploy, a failover): they should block on coordination; better to wait than act on stale state.
  • Security revocations (kill-switches, key revocation): may need a max-staleness bound (e.g., 30s) after which clients fail closed — a deliberate, signed-off exception.

🧭 Principal Move: "I'd write 'fail static' into the platform contract: every client library caches last-known-good, exposes its staleness as a metric, and keeps serving. Teams that want fail-closed must file an exception and own the resulting availability risk."


4. Failure Modes & Operational Reality#

4.1 Leader Flapping From a Slow Disk (Gray Failure)#

t=0:       Node 2 (leader) disk degrades: fsync p99 2ms → 450ms. Not dead.
t=+1s:     Heartbeats from leader delayed behind fsync. Followers hit election timeout (1s).
t=+1.2s:   Node 4 elected. Node 2 steps down. Writes fail for ~1.2s.
t=+20s:    Node 2, now a follower with less load, looks healthy; wins next election after node 4 hiccups.
t=+25s:    Node 2 leader again → flaps again.
t=+10min:  leader_changes_total = 38. Kubernetes API p99 12s. Controllers retrying. Deploys stall.
t=+12min:  Page: leader_changes > 3 in 15m. Runbook: identify node with worst wal_fsync p99.
t=+15min:  Node 2 removed from membership (member remove). Cluster stable on 4 voters.
t=+2h:     Replacement node added as learner, caught up, promoted.

Detection: etcd_server_leader_changes_seen_total, etcd_disk_wal_fsync_duration_seconds p99 per member, etcd_server_proposals_failed_total.

Mitigation: remove the sick member; don't just restart it.

Prevention: dedicated local SSD for WAL (not network storage shared with noisy neighbors); pre-vote and check-quorum enabled; election timeout ≥ 10× heartbeat and ≥ 5× p99 RTT; per-member disk latency alerts that fire before elections.

Owner: coordination platform on-call.

4.2 Split Brain at the Resource — The GC Pause#

t=0:       Billing scheduler A holds leadership (lease 10s, no fencing).
t=+3s:     A starts a full GC on a 64 GB heap: 24s pause.
t=+13s:    Lease expires. B elected. B starts nightly run.
t=+27s:    A resumes mid-loop, dispatches remaining 6,000 invoices.
t=+28s:    Both A and B dispatch. Duplicate invoices: 6,000.
t=+9h:     Customer complaints. Incident opened.

Detection: jobs.duplicate_dispatch_total (dedup at the downstream), JVM gc_pause_seconds_max, leader.stale_token_rejections_total (only exists once fencing is in).

Mitigation: credit notes; dedup on invoice idempotency key.

Prevention: fencing token on the job store; downstream idempotency; GC tuning to keep pauses well below lease TTL — but never rely on that alone.

Owner: the scheduler's team (resource owner) with the platform lock library team.

4.3 Quorum Loss During a Zone Outage + Maintenance#

t=0:       3-node cluster across zones A/B/C. Rolling OS patch: node C drained.
t=+5min:   Zone A has a power event. Node A down.
t=+5min:   1 of 3 voters alive. No quorum. All writes fail. Reads (serializable) still served.
t=+5min:   Leader elections impossible. Controllers can't update state. Deploys halt.
t=+25min:  Node C patch completes, rejoins. 2 of 3: quorum restored.

Detection: etcd_server_has_leader == 0, proposals_pending rising.

Mitigation: complete maintenance faster; never force a new cluster from one survivor unless data loss is accepted and signed off (it discards uncommitted entries and risks divergence).

Prevention: 5 voters for anything critical so maintenance of one node plus a zone loss is survivable; maintenance windows blocked during declared cloud-provider incidents.

Owner: platform on-call; change management owns the maintenance policy.

4.4 Watch / Reconnect Storm#

t=0:       Network blip 8s between the fleet and the coordination cluster.
t=+8s:     40,000 clients reconnect simultaneously; each re-establishes session + 5 watches.
t=+9s:     200,000 watch creations + full-state re-reads. Leader CPU 100%. Heartbeats late.
t=+11s:    Leader election triggered by overload. Clients see errors, reconnect again.
t=+30s:    Positive feedback loop — the storm keeps the cluster unhealthy.

Detection: watchers_count step changes, grpc_server_started_total rate, leader CPU.

Mitigation: proxy layer absorbs reconnects (1 upstream watch per key regardless of client count); clients reconnect with jittered backoff (0–30s) and resume watches from their last revision rather than re-reading everything.

Prevention: clients never connect directly at fleet scale; watch-resume from revision is mandatory in the client library.

Owner: platform (proxy and client library).

4.5 Database Size Exceeded — The Read-Only Cluster#

t=0:       A team starts storing 200 KB job specs as keys. 40,000 jobs/day.
t=+6 days: DB size 1.9 GiB of 2 GiB quota. No alert on size.
t=+7 days: NOSPACE alarm. Cluster rejects writes. Every tenant's elections and locks fail.

Detection: etcd_mvcc_db_total_size_in_bytes / quota > 70% warn, > 85% page.

Mitigation: compact old revisions, defragment member by member, disarm alarm; evict the offending tenant's keys.

Prevention: per-tenant byte quotas at the proxy; value size limit per tenant (e.g., 64 KB); scheduled compaction (etcd supports periodic auto-compaction).

Owner: platform; offending tenant owns the cleanup and migration of payloads to object storage.

4.6 Membership Change Gone Wrong#

Replacing a node by adding the new one before removing the old one briefly makes a 6-member cluster needing 4 for quorum; if the new node can't catch up (large snapshot, slow network), the cluster is now less available. Replacing two nodes at once can produce a moment where old and new configurations have disjoint majorities if the implementation doesn't use joint consensus or single-server changes.

Prevention: add as learner first, wait until caught up (lag < 1,000 entries), then promote; change one voter at a time; automate it — humans shouldn't type member add during an incident.

Owner: platform; runbook-automated.

4.7 Operational Reality Matrix#

FailureDetection SignalBlast RadiusMitigationOwner
Slow disk → flappingleader_changes > 3/15m, wal_fsync_p99 > 10msAll tenants' writes; K8s APIRemove sick memberPlatform on-call
Split brain at resourcestale_token_rejections, downstream duplicatesOne tenant's dataFencing; idempotencyResource team + lock library
Quorum losshas_leader == 0All writes, all tenantsRestore members; never force-new casuallyPlatform on-call
Reconnect stormwatcher/connection rate spikesCluster-wide overloadProxy, jittered backoff, resume from revisionPlatform
DB quotadb_size / quota > 85%Cluster-wide read-onlyCompact, defrag, evictPlatform + tenant
Bad config pushedError-rate spike correlated with revisionEvery consumer of the keyRevert revision; staged rollout of configConfig owner
Membership mistakeMember lag, quorum size changeCluster availabilityLearner-first, one at a timePlatform
Clock jump on holderLease-margin violationsLease holder correctnessMonotonic clocks; fencingLibrary team

🎯 Staff Insight: Most coordination outages are multi-tenant blast radius problems: one tenant's behavior (big values, many watches, write floods) takes down everyone's elections. Quotas at the access layer are not a nice-to-have.


5. Evaluation Rubric#

5.1 Level-Based Signals#

DimensionSenior (L5)Staff (L6)Principal (L7)
FramingStarts with RaftSeparates coordination vs config vs data replication; caps size and write rateAsks how many coordination clusters the org already runs and who owns them
SafetyQuorum overlap explained wellFencing tokens at the resource; efficiency vs correctness locksPlatform API that makes fencing the default and unfenced correctness claims visible in review
Availability"5 nodes tolerate 2"Placement across zones; gray failure; learner-first membership changesTiered placement policy with latency sign-off; game days; correlated-risk register
ReadsFollower reads for scaleReadIndex vs lease vs serializable vs watch — chosen per useConsistency menu in the client library; defaults chosen centrally
Operations"Monitor the cluster"fsync p99, leader changes, DB size vs quota, watcher counts — with runbooksError budgets for the coordination tier that gate tenant onboarding
Build vs buyDesigns custom consensusRuns etcd/ZooKeeper; avoids when CAS sufficesConsolidates clusters; plans deprecations; negotiates managed options

5.2 Strong Hire Signals#

SignalWhat It Sounds Like
Separates inside vs outside safety"Consensus prevents two leaders committing. Only fencing prevents two leaders acting."
Places quorums against real failures"2-2-1 across three zones; zone loss leaves three."
Knows the gray failures"The outage I expect is a slow disk, not a dead node."
Keeps it off the hot path"The data plane fails static on cached config."
Reduces scope"This doesn't need a lock — the row has a version column."
Names tenants and quotas"One tenant's 200 KB values shouldn't be able to fill everyone's quota."

5.3 Lean No-Hire Signals#

SignalWhy It Misses the Bar
"Redis SETNX with TTL is a safe distributed lock"No consensus, no fencing; fails under failover and pauses
"Even number of nodes for balance"4 nodes tolerate 1 failure, same as 3, with slower writes
"Read from followers, it's strongly consistent"Confuses replication with linearizability
"We'll build our own Paxos" with no testing strategyUnderestimates a multi-year correctness project
Stores application data in the coordination storeIgnores memory, snapshot and quota limits
No answer for "the leader paused"The central question of the interview

5.4 Common False Positives#

  • Reciting Raft's log-matching property ≠ designing a coordination service. Protocol fluency is table stakes; the level is decided at the boundaries.
  • Mentioning CAP ≠ understanding tradeoffs. "It's CP" is a label. What matters is what becomes unavailable (writes, elections) and what keeps working (cached reads, the data plane).
  • Mentioning Jepsen ≠ testing culture. The Staff answer names what they'd test: partitions during membership change, clock skew on lease holders, disk stalls on the leader.
  • Big cluster sizes ≠ reliability. "Nine nodes for safety" usually signals someone who hasn't paid the write latency.

6. Interview Flow & Pivots#

6.1 Typical 45-Minute Shape#

PhaseTimeGoal
Framing0–3 minCoordination vs data replication; size and write caps; failure domain
API3–5 minKV + CAS + leases + watches; recipes on top
Architecture5–10 min5 voters, 3 zones; proxy; write path latency
Fencing10–18 minLeases vs pauses; token at the resource
Placement + gray failure18–26 minQuorum math with domains; slow disk; membership changes
Reads + blast radius26–34 minRead modes; fail static; quotas
Pivot34–42 minMulti-region, build vs buy, 50K clients, lock service at scale
Wrap42–45 minInside vs outside safety; buy it; one platform

6.2 How Interviewers Pivot — And What They're Testing#

PivotWhat They're TestingStrong Response Shape
"Make it survive a region outage"Do you know the latency price?3 regions minimum; +30–80ms per write; election timeouts in seconds; or regional clusters
"Why not Redis for locks?"Fencing and failover semanticsAsync replication can lose the lock on failover; no monotonic token by default
"100K clients need config"Access-layer designProxy with watch coalescing; client cache; resume from revision
"Implement it yourself"Correctness engineering maturityTLA+ spec, deterministic simulation, fault injection, years
"The leader is slow, not dead"Gray failurePre-vote, check-quorum, disk alerts, remove member
"Add 2 nodes to the running cluster"Membership safetyLearner first, one at a time, joint consensus

6.3 What to Deliberately Skip#

  • Paxos vs Raft proofs — one sentence: "equivalent safety, Raft's structure is easier to implement correctly."
  • Byzantine fault tolerance — say it's out of scope unless nodes are mutually untrusted (blockchains, multi-party systems); BFT needs 3f+1 nodes.
  • Log compaction internals — mention snapshots exist; skip formats.
  • Transport details — gRPC vs custom; irrelevant to the design.

6.4 Follow-Up Questions to Expect#

  1. "What's your fencing token, concretely, and who checks it?"
  2. "Five nodes in two zones — what's your availability during a zone outage?"
  3. "A client reads config from a follower, then reads from another follower and sees an older value. Is that allowed?"
  4. "How do you replace a failed node without risking quorum?"
  5. "One tenant is writing 5,000 keys/s. What happens to everyone else?"
  6. "How long does failover take, and what's the tradeoff in tuning it faster?"
  7. "When would you not use this service?"

7. Active Drills#

Drill 1: The Opening#

Prompt: "Design a distributed coordination service like ZooKeeper."

Staff Answer

"Before I draw — what are clients using it for? Leader election and locks, config distribution, or as the replication layer of a database? I'll assume the first two: small data, under a few GB, writes in the hundreds to low thousands per second, tens of thousands of clients reading and watching. I'll design for surviving a zone, not a region, and I'd run etcd or ZooKeeper rather than implement consensus.

I'll cover: a minimal API — CAS, leases, watches; five voters across three zones; an access layer that coalesces watches; fencing tokens for correctness locks; read modes; and quotas so one tenant can't take everyone down."

Why this is L6:

  • Distinguishes three intents and caps data and write rate
  • States the failure domain and its cost up front
  • Buys the consensus core, putting design effort where incidents happen

What L7 adds:

  • Asks how many coordination clusters already exist and proposes consolidation
  • Frames the service as a tiered platform with onboarding criteria and error budgets
❌ Common L5 Trap

"I'll implement Raft. Nodes start as followers; if they don't hear a heartbeat within a randomized timeout, they become candidates…"

Why this misses: Accurate and unasked-for. Ten minutes later the candidate hasn't addressed who uses the service, how much data it holds, or what protects resources from stale leaders.


Drill 2: Fencing — Make It Concrete#

Prompt: "You said fencing tokens. Where does the number come from and who checks it?"

Staff Answer

"In etcd, the election recipe creates a key under the election prefix attached to the leader's lease; the key's create_revision is cluster-wide monotonic, so each new leader's revision is strictly larger — that's the token. In ZooKeeper, the sequential znode number or the zxid serves the same role. The leader sends it on every write to the protected store. The store keeps max_token per resource and does UPDATE … SET data=?, max_token=? WHERE resource=? AND max_token <= ?. Zero rows updated means a newer leader exists; the old leader exits. If the protected resource is S3 or a third-party API that can't compare tokens, I make the operations idempotent with deterministic keys, and I call the lock an efficiency lock."

Why this is L6:

  • Derives the token from a real mechanism rather than inventing a counter
  • Shows the check as an atomic conditional write at the resource
  • Has a fallback for resources that can't enforce

What L7 adds:

  • Ships the check as a library for the three most common stores (Postgres, DynamoDB, the internal job store)
  • Makes the fenced-resource declaration part of the service catalog

Drill 3: Quorum Placement#

Prompt: "Your cloud region has three zones. How many nodes, where, and why?"

Staff Answer

"Five voters as 2-2-1. Any zone loss leaves at least three — a majority. With three voters as 1-1-1, a zone loss leaves two, which works, but any other disruption — a rolling upgrade, one node's disk failing — during the zone event loses quorum. Five gives maintenance headroom. Write latency: the leader needs two follower acks; with one follower in-zone and others at 1–2ms, commit is ~2–4ms including fsync. I wouldn't go to seven — the third-fastest ack becomes the fourth, and zone count, not node count, is my failure domain."

Why this is L6:

  • Reasons from failure domains to counts
  • Identifies maintenance-during-incident as the realistic double fault
  • Quantifies latency impact of adding voters

What L7 adds:

  • Standardizes cluster shapes the platform will provision
  • Ties shape to the tier of the dependent service, with sign-off for exceptions

Drill 4: Coordination Is Down#

Prompt: "The coordination cluster has no quorum for 10 minutes. What happens to the company?"

Staff Answer

"If we designed it right: running services keep serving on cached config — fail static — and existing leaders keep working until their leases expire; after that, leaderless components pause their leader-only work, like schedulers and compaction. What breaks: deploys, new leader elections, lock acquisitions, config changes. User traffic should see nothing. If user traffic is failing, that's a finding — some service put coordination on its request path, and the post-mortem is about that dependency, not the outage. During the incident: don't force-new-cluster from one survivor unless we accept losing recent writes; restore members."

Why this is L6:

  • Describes blast radius by category — data plane vs control plane
  • Treats user impact as evidence of a design violation
  • Knows the dangerous recovery command and when not to use it

What L7 adds:

  • Maintains a dependency map of which services degrade how during coordination loss, tested in game days twice a year
  • Provides a break-glass path for deploys that doesn't depend on the coordination cluster

Drill 5: The Hot Key / Herd Effect#

Prompt: "5,000 workers wait on the same lock. It's released. What happens?"

Staff Answer

"With a naive recipe — everyone watches the lock key — release triggers 5,000 notifications and 5,000 simultaneous acquire attempts: a herd. The standard recipe: each waiter creates a sequential key under the lock prefix and watches only its immediate predecessor; release wakes exactly one waiter. But I'd also question the design — 5,000 waiters on one lock is a queue, not a lock. I'd replace it with a work queue or partition the resource into 64 locks keyed by hash so contention drops 64×."

Why this is L6:

  • Knows the predecessor-watch recipe
  • Challenges the premise — reframes contention as a data modeling problem
  • Quantifies the partitioning benefit

What L7 adds:

  • Adds lock contention metrics per tenant to the platform; tenants with > 100 waiters get a design review

Drill 6: Multi-Tenant Platform#

Prompt: "30 teams want to use your coordination service. How do you keep one from hurting the rest?"

Staff Answer

"Every tenant gets a key prefix and an identity, and the access layer enforces quotas: keys (e.g., 10K), bytes (e.g., 50 MB), max value size (64 KB), write rate (e.g., 200/s), watch count. Exceed a quota and that tenant is throttled with a clear error; others are unaffected. The Kubernetes control plane and other tier-0 consumers get dedicated clusters — they don't share with feature flags. Onboarding includes a review: what's your data size, your write rate, is your lock for correctness, where's your fence."

Why this is L6:

  • Quotas enforced before consensus, per tenant
  • Separates tier-0 consumers physically
  • Onboarding questions encode the design principles

What L7 adds:

  • Chargeback or showback per tenant for write volume
  • Error budget for the platform that pauses new onboardings when burned

Drill 7: Build vs Buy#

Prompt: "Our database team wants to write their own Raft for a new storage engine. Do you approve?"

Staff Answer

"Depends on whether consensus is inside the product's data path and what they'll invest. If they need multi-Raft — thousands of groups per node — the existing coordination services don't fit and embedding a proven library like etcd's raft package is the right start: it gives the core state machine and leaves storage and transport to them. Writing Raft from scratch, I'd approve only with a TLA+ spec, deterministic simulation testing, a Jepsen-style fault suite, and a team committed for 2+ years. Otherwise we're buying a class of bugs that shows up in year three."

Why this is L6:

  • Differentiates embed vs build vs run
  • Names concrete correctness practices as the price of building
  • Sets a time horizon for ownership

What L7 adds:

  • Makes the decision a formal build-vs-buy review with a 3-year TCO and staffing commitment
  • Plans what happens if the team is reorganized — who inherits a custom consensus implementation?

Drill 8: Changing Timeouts Safely#

Prompt: "Failover takes 3 seconds. Product wants 300ms. Change the election timeout?"

Staff Answer

"Faster election timeouts mean more false elections: any GC hiccup, fsync spike or packet loss longer than the timeout triggers an election, and each election is itself a ~1× timeout write outage. With cross-zone p99 RTT of 2ms and fsync p99 of 10ms, 300ms is feasible but leaves little margin for disk stalls. I'd first ask what the 3s is costing — for coordination, writes pausing for 3s during a leader crash, which happens a few times a year, is usually fine. If it must change: roll out in stages — 1 cluster, measure leader_changes_total for a week, then more — and keep pre-vote on."

Why this is L6:

  • Explains the false-positive cost of faster failover
  • Anchors in measured latencies
  • Stages the change and defines the success metric

What L7 adds:

  • Challenges the requirement: if a user-facing product depends on 300ms coordination failover, it's using coordination on the data path — fix that instead

Drill 9: Multi-Region#

Prompt: "We're going to three regions. What happens to the coordination service?"

Staff Answer

"Default: one cluster per region, each zone-survivable, holding regional state. Most coordination is regional — a region's schedulers elect a regional leader. For the small set of truly global state — say, which region owns which tenant — I'd run one region-spread cluster with 5 voters as 2-2-1 across three regions, election timeout ~3–5s, and write rates in the tens per second. Everything else replicates global config to regional clusters asynchronously with a revision stamp. I would not stretch every tenant's cluster across regions — it multiplies write latency 30–50× for state that doesn't need it."

Why this is L6:

  • Classifies state by scope before choosing topology
  • Quantifies the cost of region-spread consensus
  • Uses async replication for global-but-tolerant config

What L7 adds:

  • Sets the policy that global consensus is an exception requiring architecture review
  • Budgets cross-region egress and the extra on-call surface

8. Deep Dive Scenarios#

Deep Dive 1: Peak-Traffic Incident — Kubernetes Control Plane Meltdown During a Scale-Up#

Context: A marketing event triggers autoscaling from 800 to 3,000 nodes in 20 minutes. The Kubernetes API server p99 goes to 30s; pod scheduling stalls; the scale-up stalls with it. The on-call escalates to you.

Questions to Surface First:

  • Is etcd the bottleneck (fsync, DB size, leader changes) or the API server (watch cache, request volume)?
  • What's generating writes — node heartbeats, pod status updates, events?
  • Is anything writing large objects?

Typical L5 Approach: Scales the API server replicas and increases etcd resources.

Staff Approach: Checks etcd first: wal_fsync_p99 is 40ms (network disk), and Event objects are 60% of writes. Moves events to a separate etcd cluster (a supported Kubernetes configuration), confirms node heartbeats use Lease objects rather than full status updates, and plans migration to local NVMe for WAL.

Principal Approach: Treats cluster size as a capacity-planned platform limit: publishes a per-cluster node ceiling (e.g., 2,000) and a multi-cluster strategy for larger footprints, and adds control-plane load tests to the pre-event readiness checklist.

Staff Approach — Full Reasoning
PhaseWhat to Do
Immediate (0–5 min)Pause the scale-up at current size; stop the feedback loop. Check fsync p99, leader changes, DB size.
Triagefsync p99 40ms on network-attached disks; event writes dominate; no leader flapping yet.
Quick fixRaise event TTL cleanup, reduce event rate from noisy controllers; resume scale-up in 200-node steps.
GuardrailsAlert on fsync p99 > 10ms and on write rate > 70% of tested capacity.
Post-mortemControl plane was never load-tested at 3,000 nodes; disk class chosen for cost.

Metrics to Watch: etcd_disk_wal_fsync_duration_seconds, etcd_server_proposals_pending, apiserver_request_duration_seconds, etcd_mvcc_db_total_size_in_bytes

Organizational Follow-up: events cluster separation as default; storage class for etcd fixed by platform, not chosen per cluster.

Ownership Question: "Who should have known 3,000 nodes would break it?" Staff answer: The platform team owns the tested ceiling and must publish it; the event planners own asking. Neither happened.

Key Takeaway: "The coordination store's disk is the control plane's heartbeat. Pay for it."

What clears the Staff bar:

  • Checks consensus-layer metrics before scaling stateless layers
  • Separates write classes by criticality
  • Turns an incident into a published capacity limit

Deep Dive 2: Silent Failure — The Lock That Never Protected Anything#

Context: An audit finds that a data-compaction job, "protected" by a ZooKeeper lock for two years, has occasionally run twice concurrently, causing 0.02% of partitions to have duplicated rows. Nobody noticed.

Questions to Surface First:

  • Does the compaction writer check any token?
  • When did concurrency happen — correlated with GC pauses, deploys, network blips?
  • What else uses the same lock library?

Typical L5 Approach: Increases session timeout from 10s to 60s so the lock doesn't expire during pauses.

Staff Approach: Recognizes that longer timeouts only shrink the window and slow failover. Adds fencing: the compaction writer passes the lock node's sequence number to the storage layer, which rejects lower tokens per partition. Adds a detector — a nightly job counting duplicate rows — and backfills fixes.

Principal Approach: Audits every consumer of the lock library: 23 services, of which 17 claim correctness, 3 have a fence. Launches the correctness_lease API with required fence declaration and sets a 2-quarter deadline for the 14 unfenced users, tracked at the engineering-leadership level.

Staff Approach — Full Reasoning
PhaseWhat to Do
ImmediateQuantify affected partitions; stop compaction if duplicates are growing.
TriageCorrelate duplicate events with gc_pause_seconds_max > session timeout: 11 of 12 matched.
Quick fixFencing check in the storage writer; dedup repair job.
Guardrailsstale_token_rejections_total metric; alert if > 0 so we see pauses that would have corrupted data.
Post-mortemThe lock library's docs said "distributed mutex" with no mention of pauses.

Metrics to Watch: lock.stale_token_rejections_total, jvm.gc_pause_seconds_max, compaction.duplicate_rows_detected

Ownership Question: "Whose bug is it — the library or the job?" Staff answer: Both. The job owner owns their resource's correctness; the library owner owns an API that implied safety it couldn't provide.

Key Takeaway: "A lock without a fence is a probability, not a guarantee."

What clears the Staff bar:

  • Rejects timeout tuning as a fix for a safety problem
  • Adds a metric that proves the fence is working
  • Looks beyond the one job to the whole library's users

Deep Dive 3: Large-Customer Onboarding — The Team That Wants to Store 40 GB#

Context: A new ML platform team wants to use the shared etcd for model registry metadata: 2 million keys, 40 GB total, 3,000 writes/s during training bursts.

Questions to Surface First:

  • What in this data actually needs consensus — probably "which model version is live", not every artifact's metadata?
  • What are the read patterns and consistency needs?
  • What happens to them if coordination is unavailable?

Typical L5 Approach: Increases the etcd quota to 64 GB and adds nodes.

Staff Approach: Declines: 40 GB exceeds practical etcd limits, would make snapshots and recovery take tens of minutes, and 3,000 writes/s would consume most of the shared cluster's budget. Proposes: metadata in PostgreSQL or DynamoDB; only the "live version pointer" per model (a few thousand keys) in coordination, with watches for rollout.

Principal Approach: Uses the request to publish the platform's acceptance criteria (size, rate, value size, use cases) and a decision guide — "coordination vs database vs cache" — so the next team self-selects correctly.

Staff Approach — Full Reasoning
PhaseWhat to Do
ScopingSplit data into "needs agreement" (pointers, leases) vs "needs storage" (metadata).
ProposalRegistry in a DB with version columns (CAS); pointers in etcd with watches.
QuotaTenant quota: 10K keys, 20 MB, 100 writes/s.
ValidationLoad test training-burst pattern against the proposed split.
Follow-upReview in 1 quarter: usage vs quota.

Metrics to Watch: tenant keys_count, bytes, write_rate; cluster db_size, proposals_pending

Ownership Question: "Who says no?" Staff answer: The platform team, against published criteria — so it's a policy decision, not a personal one.

Key Takeaway: "Consensus is for agreement, not storage."

What clears the Staff bar:

  • Separates the agreement-needing sliver from bulk data
  • Protects other tenants' budget
  • Turns a one-off no into a published policy

Deep Dive 4: Post-Mortem — Two Leaders After a Network Partition#

Context: During a partition, a primary database's failover manager promoted a replica while the old primary kept accepting writes for 90 seconds. 4,000 writes diverged. The failover manager used a consensus service for election.

Questions to Surface First:

  • Did the old primary know it lost leadership? How would it find out?
  • Did clients route writes by a cached "who's primary" value?
  • Was there any fencing at the storage or proxy layer?

Typical L5 Approach: Shortens the old primary's lease check interval.

Staff Approach: Identifies that the consensus service elected correctly — the old primary simply couldn't reach it and kept serving. Fix: the old primary must stop accepting writes when it can't renew its lease within the lease period minus margin (self-demotion), and clients or proxies must route by epoch — writes carry the primary's epoch, and replicas/proxies reject writes from lower epochs. Divergent writes need reconciliation with owners.

Principal Approach: Standardizes failover across all database fleets on one manager with STONITH-style fencing (revoke network access or storage access of the old primary) and epochs, and requires a partition test in the failover certification of every database tier.

Staff Approach — Full Reasoning
PhaseWhat to Do
ImmediateStop writes to the old primary; snapshot both sides.
Triage4,000 divergent writes; classify by table and conflict type.
RepairOwners replay or discard per table with audit trail.
FixSelf-demotion on lease loss + epoch-based write fencing at proxy + network-level fencing of old primary.
TestPartition game day: isolate primary from coordination but not clients.

Metrics to Watch: db.writes_by_epoch, failover.self_demotions_total, proxy.stale_epoch_rejections

Ownership Question: "Who owns divergent data repair?" Staff answer: Each table's owning team, with the database platform providing tooling and the incident commander setting deadlines.

Key Takeaway: "Election is agreement among voters. Failover is making the loser stop. They're different problems."

What clears the Staff bar:

  • Distinguishes a correct election from an incomplete failover
  • Fences at multiple layers (self, proxy, network)
  • Tests the exact partition shape that caused the incident

Deep Dive 5: Multi-Region Expansion — Global Leader for a Payments Scheduler#

Context: The company goes active-active in two regions. A payments settlement scheduler must have exactly one active instance globally. Someone proposes a 4-node etcd cluster, 2 per region.

Questions to Surface First:

  • What happens when the inter-region link partitions?
  • Can the scheduler be region-scoped instead (each region settles its own accounts)?
  • What's the cost of a 10-minute scheduling pause vs a duplicate run?

Typical L5 Approach: Accepts the 4-node design; argues it's symmetric.

Staff Approach: Rejects 2+2: a partition leaves both sides with 2 of 4 — no majority anywhere, so no leader at all. Needs a third site: 5 voters as 2-2-1 with the 1 in a third region (or a lightweight witness location). Better: partition the work by region so each regional scheduler owns its accounts, eliminating the global leader; plus fencing via settlement-batch IDs in the ledger.

Principal Approach: Establishes that "global singleton" requirements need architecture review, since each one creates a three-region dependency; pushes the org toward region-partitioned ownership as the default pattern.

Staff Approach — Full Reasoning
PhaseWhat to Do
ChallengeCan ownership be partitioned by region? Usually yes for settlement.
If global needed3 regions, 2-2-1; election timeout 3–5s; write rate trivial.
FencingSettlement batches carry leader token; ledger rejects stale tokens.
Failure drillPartition each region in turn; confirm exactly one leader or none.
CostThird-region witness infra; cross-region latency on leader changes.

Metrics to Watch: election.leader_region, settlement.stale_token_rejections, etcd_network_peer_round_trip_time_seconds

Ownership Question: "Who decides pause vs duplicate risk?" Staff answer: Finance and payments jointly — a pause is almost always cheaper than a duplicate settlement.

Key Takeaway: "Even voter counts across two sites buy you a partition with no leader. Always a third site."

What clears the Staff bar:

  • Does the partition math on the proposed topology
  • Removes the global singleton when possible
  • Prefers unavailability over duplication for money

9. Level Expectations Summary#

After studying this case study, you should be able to:

  • Explain why consensus guarantees agreement inside the cluster but not exclusion outside it, and close the gap with fencing tokens
  • Size and place a quorum against zones and regions, and state the write-latency cost of each shape
  • Choose a read mode (fence, CAS, ReadIndex, watch, serializable) per use case
  • Diagnose gray failures: slow disks, leader flapping, reconnect storms, quota exhaustion
  • Perform membership changes safely: learner first, one voter at a time
  • Argue build vs buy vs avoid with concrete criteria
  • Design a multi-tenant coordination platform with quotas, tiers and fail-static clients

The Bar for This Question#

Mid-level (L4): Describes leader election and majority replication correctly. Uses a lock with a TTL. Doesn't address pauses, placement or blast radius.

Senior (L5): Explains Raft well, including terms and log matching; chooses 5 nodes; mentions ZooKeeper recipes and watches. The gap: treats the service's guarantees as the system's guarantees — no fencing at resources, node counts without failure domains, follower reads described as consistent, and a tendency to build rather than buy.

Staff+ (L6): Frames the service by intent and caps its scope. Fences at resources. Places voters against real failure domains and quantifies latency. Names gray failures with metrics and runbooks. Keeps coordination off the data plane with fail-static clients. Protects tenants from each other with quotas. Knows when a database's CAS makes the whole service unnecessary. The interviewer should learn something from the answer.


10. Staff Insiders: Controversial Opinions#

10.1 "Distributed Locks" Are Mostly Misnamed Leases#

What People SayWhat They Have
"Distributed mutex"A time-bounded lease with no enforcement at the resource
"Safe lock with Redis"An efficiency lock that fails under failover and pauses
"ZooKeeper lock = correct"Correct liveness; safety only with a fence

The Staff position: Call it a lease. Say whether it's for efficiency or correctness. If correctness, name the fence.

Why this matters in interviews: Precise language signals you've seen the failure. "Lock" without qualification signals you haven't.

10.2 Most Teams Asking for Consensus Need a Conditional Write#

RequestBetter Answer
"Lock around updating this row"UPDATE … WHERE version = ?
"Elect one worker per partition"Kafka consumer groups or a lease row in the DB with a version fence
"Exactly one cron run"Unique constraint on (job, scheduled_time)

The Staff position: Put mutual exclusion where the data is. External coordination is for when no single store owns the state.

Why this matters in interviews: Reducing scope is a senior judgment signal that most candidates never show.

10.3 Seven Nodes Is Almost Always Wrong#

The Staff position: Real failure domains come in threes (zones, regions). Five voters across three domains covers the realistic double fault. Seven adds write latency and operational surface for failures that are rarely independent anyway.

Why this matters in interviews: Picking cluster size from failure domains, not from "more is safer", shows operational maturity.

10.4 The Coordination Service Is Your Biggest Single Point of Failure — Even Though It Has No Single Point of Failure#

PropertyReality
No single node is criticalTrue
No single cluster is criticalFalse — every controller, scheduler and deploy depends on it
Correlated failureA bad upgrade, full quota or overload hits all replicas at once

The Staff position: Replication protects against independent failures. Tiering, quotas, fail-static clients and staged upgrades protect against correlated ones — which are the ones that happen.

Why this matters in interviews: It reframes availability from nodes to dependencies — an L6+ lens.

10.5 You Should Never Write Your Own Consensus — Until It's Your Product#

The Staff position: Every mature system that built its own (Chubby, ZooKeeper, etcd, KRaft) had a team whose job was that system, for years. If coordination isn't your product, it isn't your project.

Why this matters in interviews: "I'd run etcd" is not a cop-out; it's the answer — provided you then show you understand what etcd can't do for you.


11. The Principal Lens (L7)#

Why L7 Sees This Problem Differently#

The Staff engineer designs a correct, well-placed coordination cluster. The Principal engineer finds that the company runs eleven of them — three ZooKeeper ensembles from the Kafka and Hadoop era, five etcd clusters for Kubernetes, a Consul cluster a team adopted for service discovery, and two Redis "lock services" nobody admits to — each with different upgrade states, different on-call owners and different ideas of what "lock" means. The L7 problem is coordination as a governed platform: how many clusters, which products, which guarantees, who is on call, and how the organization avoids a correlated failure that takes down every control plane at once.

The Org-Level Fault Line#

One shared coordination platform vs per-team clusters.

OptionWhat WorksWhat BreaksWho Pays
Per-team clustersIsolation; team autonomyN upgrade backlogs, N on-calls, inconsistent guarantees, forgotten clusters with CVEsEach team's on-call; security
One shared clusterOne owner; consistent APICorrelated blast radius; noisy tenants; one bad upgrade hits everyoneEveryone at once
Platform-managed fleet of tiered clustersOne owner, one API, isolation by tier and cellPlatform team must build provisioning, quotas, upgrade automationPlatform team (3–6 engineers)

🧭 Principal Move: "One product, one owner, many clusters. The platform team runs a fleet of identical, tiered clusters — tier-0 dedicated for Kubernetes control planes, tier-1 shared with quotas for services — and upgrades them in waves, never all at once. Teams get isolation without owning consensus."

Cost Model#

Assumptions: cloud VMs with local NVMe; fully loaded engineer ~$250K/year; figures are order-of-magnitude.

ScaleClustersInfra ($/month)HeadcountOn-call Load
Startup (1 product, 1 region)1–2 (managed Kubernetes etcd + DB CAS for locks)~$0–1K (managed)~0.25 FTECovered by cloud provider + shared rotation
Growth (30 services, 2 regions)6–10 (K8s etcd per cluster, 1–2 shared tier-1)~$5–15K1–2 FTE platform~1–3 pages/month
Enterprise (500+ services, 5+ regions)50–150 clusters in a managed fleet~$60–150K4–8 FTE coordination platformDedicated rotation; automated remediation for top-5 alerts

The pricing insight: the dominant cost is never the VMs — it's engineer time on upgrades, incidents and bespoke clusters. Consolidating 11 unmanaged clusters into a fleet typically saves more in on-call and upgrade toil than it costs in platform headcount within a year.

The 3-Year Evolution Path#

Diagram: The 3-Year Evolution Path

One-Way Doors vs Two-Way Doors#

DecisionDoor TypeReversibility Cost
Building a custom consensus implementationOne-wayA team and codebase you must staff for years; migrating off means data and client migration
Region-spread quorum for a service's stateOne-way-ishClients assume global consistency; splitting to regional ownership later is a redesign
Allowing application data in coordination storesOne-way-ishEvicting it later requires every tenant to migrate
Client API semantics (what "lock" guarantees)One-wayChanging guarantees silently breaks callers; needs a new API and migration
etcd vs ZooKeeper for a new clusterTwo-wayRecipes differ, but small data migrates in days
Timeouts, quotas, cluster sizesTwo-wayConfig, staged rollout
Proxy/access layer in front of consensusTwo-wayCan be introduced incrementally

The Standard I'd Write#

RFC-COORD-001: Coordination & Locking Standard
Status: Approved   Owner: Coordination Platform

Scope
  Any use of consensus-backed coordination (leader election, locks, leases,
  membership, configuration) by production services.

MUST
  1. Use the platform coordination service or a platform-approved managed equivalent.
     No new self-managed ZooKeeper/etcd/Consul clusters.
  2. Classify every lock as EFFICIENCY or CORRECTNESS.
  3. CORRECTNESS leases MUST pass a fencing token to a declared resource that rejects
     stale tokens; the resource is registered in the service catalog.
  4. Data-plane clients MUST fail static on cached state and expose staleness as a metric.
  5. Stay within tenant quotas: ≤ 50 MB, ≤ 64 KB per value, ≤ 200 writes/s unless approved.

SHOULD
  1. Prefer resource-native conditional writes over external locks.
  2. Use watches via the platform proxy, resuming from the last revision.
  3. Keep coordination out of synchronous user request paths.

Exceptions
  Architecture review; region-spread (global) coordination requires a latency sign-off
  from the consuming service's owner and director.

Success metrics
  - Unmanaged coordination clusters: 0 by end of year 2
  - CORRECTNESS leases without a declared fence: 0
  - Coordination-caused user-facing incidents: ≤ 1 per year
  - Fleet upgrade lag: all clusters within 1 minor version within 60 days of release

What I'd Tell the VP#

"Every deploy, scheduler and failover in the company depends on a handful of small coordination clusters — and today we run eleven of them with six different owners, two of which nobody is on call for. The risk isn't that one fails; it's that a bad upgrade or a misbehaving team takes several down at once and stops our ability to deploy the fix. I'm proposing one platform team that runs a standardized fleet with per-team limits, and a lock standard that closes the class of 'two workers did the same job' incidents we've had three times this year. It's about four engineers; it retires roughly the equivalent of two engineers' worth of scattered toil in year one and removes a correlated-outage risk we can't currently quantify."

Principal Interview Signals#

SignalWhat It Sounds Like
Counts the org's clusters"Before designing a new one, how many do we already run and who's on call for each?"
Designs against correlated failure"Replication doesn't help against a bad upgrade. Waves and cells do."
Governs semantics"The word 'lock' in our API is a contract; I'd split efficiency from correctness."
Prices toil, not VMs"The VMs cost $10K a month. The upgrades cost two engineers."
Knows when not to standardize"The storage team's embedded multi-Raft stays separate — it's their product, not a coordination use."

Staff answers that L7 interviewers find insufficient:

  • "We'll run a well-configured 5-node etcd" — correct for one cluster, silent on the ten others and on tenancy.
  • "Teams should use fencing tokens" — right principle, no mechanism to make it happen across 30 teams.
  • "We'll monitor leader changes and fsync latency" — good operations, no answer for correlated upgrade risk or deploy break-glass.

Appendices

Appendix A: Raft Mechanics in Depth#

A.1 Roles and Elections#

Diagram: A.1 Roles and Elections
  • Each node's timeout is randomized (e.g., 150–300ms in the paper; ~1s in etcd defaults) so split votes are rare.
  • A node votes at most once per term, and only for a candidate whose log is at least as up-to-date — the election restriction that guarantees a new leader has every committed entry.
  • Pre-vote (etcd, others): a would-be candidate first asks whether it could win before incrementing its term, so a partitioned node rejoining doesn't disrupt a healthy leader.
  • Check-quorum: a leader that can't hear from a majority steps down, limiting how long a partitioned leader believes it leads.

A.2 Log Replication#

Diagram: A.2 Log Replication

Commit cost = leader fsync ∥ (network RTT + follower fsync) for the fastest majority. On NVMe within a region: ~1–4ms.

A.3 Membership Changes#

Diagram: A.3 Membership Changes

Change one voter at a time (single-server changes) or use joint consensus; never add and remove multiple voters in one step. Removing the failed node first (when the rest is healthy) keeps quorum requirements lower during catch-up.

A.4 Why Paxos and Raft Are "the Same" for Interviews#

Both need a majority, both use a monotonically increasing ballot/term, both guarantee that once a value is chosen no conflicting value can be chosen. Raft packages leader election, log replication and membership into a prescribed structure; Multi-Paxos leaves them as engineering choices. Choose by implementation maturity, not by algorithm.

Appendix B: Keys, Recipes and Data Model#

/tenants/{tenant}/elections/{name}/{lease_id}     # leader = lowest create_revision
/tenants/{tenant}/locks/{name}/{lease_id}          # waiters watch predecessor
/tenants/{tenant}/members/{service}/{instance}     # attached to instance lease
/tenants/{tenant}/config/{component}               # watched; value ≤ 64 KB

Leader election recipe (etcd-style):

lease = LeaseGrant(ttl=10s); keepalive(lease) every 3s
Put(prefix + lease.id, my_addr, lease=lease)
loop:
  kvs = Get(prefix, sort=create_revision asc)
  if kvs[0].key == prefix + lease.id:
      token = kvs[0].create_revision     # fencing token
      act_as_leader(token)               # every write carries token
  else:
      Watch(kvs[my_index-1].key) until deleted

Resource-side fence:

UPDATE job_runs SET state = $2, leader_token = $3
 WHERE job_id = $1 AND leader_token <= $3;
-- 0 rows → stale leader → exit

Appendix C: Coordination Mechanisms — Quick Comparison#

MechanismSafety Under PauseLatencyFailure ModeUse For
Redis SET NX PXNo< 1msLost on failover; no tokenEfficiency locks, dedup of cheap work
Redlock (multi-Redis)No (timing assumptions)~1–5msClock and pause assumptionsRarely justified
ZooKeeper ephemeral sequentialLiveness yes, safety only with fence~2–10msHerd if all watch one nodeElections, locks with fences
etcd lease + revisionLiveness yes, safety with fence~2–5msQuota, disk latencyElections, locks, config
DB conditional writeYes, at that resource~1–5msRow contentionSingle-store exclusion
DynamoDB conditional putYes, at that item~5–10msHot partitionsLeases in serverless stacks
Kafka consumer groupsYes for partition ownership (with generation fencing)n/aRebalance stormsPartitioned consumers

Appendix D: Client Contract & Behavior#

  • Sessions: keepalive every TTL/3; on missed keepalives, the client assumes loss at TTL − margin by its monotonic clock and stops acting as leader before the server expires it.
  • Reconnect: exponential backoff with full jitter, base 100ms, cap 30s; resume watches from last revision; never re-read full state unless the revision was compacted (then re-list once).
  • Caching: last-known-good persisted locally; staleness exposed as coord.client.staleness_seconds.
  • Errors: distinguish NO_LEADER (retry later), QUOTA_EXCEEDED (tenant bug, don't retry fast), COMPACTED (re-list), STALE_TOKEN from resources (step down, never retry).

Appendix E: Observability#

Core metrics (etcd names where applicable):

  • etcd_server_has_leader, etcd_server_leader_changes_seen_total
  • etcd_disk_wal_fsync_duration_seconds p99, etcd_disk_backend_commit_duration_seconds p99
  • etcd_server_proposals_failed_total, etcd_server_proposals_pending
  • etcd_mvcc_db_total_size_in_bytes vs quota
  • etcd_network_peer_round_trip_time_seconds
  • Proxy: watchers_count, per-tenant write_rate, bytes, throttled_total
  • Clients: coord.client.staleness_seconds, lock.stale_token_rejections_total

Critical alerts:

AlertThresholdSeverity
No leaderhas_leader == 0 for 30sSev-1 page
Leader churn> 3 changes in 15mPage
WAL fsync p99> 10ms for 10mWarn; > 50ms page
DB size> 70% quota warn; > 85% pagePage
Proposals failedrising for 5mPage
Tenant throttlingsustained > 5mTicket to tenant
Stale token rejections> 0Info — proves fence works; spikes investigated

Control plane vs data plane: coordination is control plane; its outage should appear in deploy latency and controller lag, not in user error rates. A user-error-rate correlation with coordination health is itself an alertable design violation.

Appendix F: Scale Evolution#

ScaleWhat WorksWhat Breaks Next
1 team, < 100 clientsManaged etcd or DB CASNothing
10–30 teams, ~5K clientsShared tier-1 cluster + proxy + quotasNoisy tenants; upgrade coordination
100+ teams, 50K+ clientsFleet of tiered clusters, cells per regionFleet automation, upgrade waves
Consensus in data pathMulti-Raft inside the databasePlacement, hot ranges — a database problem

What you don't build on day one: a custom consensus implementation, a global cross-region cluster, a proxy layer (until ~1K clients), chargeback.

Appendix G: Multi-Tenancy, Fairness & Cost#

  • Quotas at the proxy, not the core: the consensus core has no notion of tenants; enforcement must happen before proposals are created.
  • Weighted write budgets: tier-0 tenants reserve 50% of tested write capacity; tier-1 share the rest with per-tenant caps.
  • Watch fairness: coalesce identical watches; cap per-tenant watch count (e.g., 5,000 upstream) — a tenant with 50,000 unique watches is a design review.
  • Showback: report per-tenant writes/day and bytes; the conversation "your team is 60% of our write load" resolves most noisy-neighbor issues without enforcement.
  1. Loading the index…