Technologies referenced in this case study: Cassandra · DynamoDB · PostgreSQL · ZooKeeper & etcd · Kafka
Related: Consistency Models · Consistent Hashing · Sharding · Distributed Consensus · Database Sharding · Degraded Mode
How to Use This Case Study#
Organized for interview use first, reference second. Read front-to-back once, then return to individual sections before a loop.
| Mode | Time | What to Read |
|---|---|---|
| Quick Review | 15 min | Executive Summary → Interview Walkthrough → Fault Lines table → Drills 1, 3, 5 |
| Targeted Study | 1–2 hrs | Executive Summary → Walkthrough → Section 3 (Fault Lines) → Section 4 (Failures) → two Deep Dives |
| Deep Dive | 3+ hrs | Everything, including The Principal Lens and appendices (quorum math, conflict resolution, anti-entropy) |
What is a Replicated Data Store? — Why interviewers pick this topic
A replicated data store keeps multiple copies of each piece of data on different machines — often in different availability zones or regions — so that the loss of one machine, rack, zone or region does not lose data or stop service. Every serious database does it: PostgreSQL streaming replication, MySQL semi-sync, Cassandra's tunable quorums, DynamoDB's three-AZ writes, Spanner's Paxos groups.
The hard part is not making copies. The hard part is deciding what a reader is allowed to see while the copies disagree, and who wins when two writers disagree.
Before vs After — the "vanishing profile edit" incident:
Without an explicit consistency contract:
t=0: User updates shipping address in us-east (leader)
t=+5ms: Write acked after local commit only (async replication)
t=+40ms: User's next page load is routed to a us-west follower (lag: 900ms)
t=+40ms: Page shows the OLD address. User re-submits the edit.
t=+2s: Order placed. Fulfillment reads from a follower that is 30s behind
after a replication stall. Package ships to the old address.
t=+3 days: Support ticket. Nobody can reproduce it. "Eventual consistency."
With an explicit contract (read-your-writes + leader-read for money paths):
t=0: Same edit. Response carries a commit token (LSN 88121934).
t=+40ms: Follower sees token > applied LSN → waits up to 50ms or proxies to leader.
t=+40ms: User sees the new address. No re-submit.
t=+2s: Order service reads with 'linearizable' flag → leader. Correct address.
Replication lag is a dashboard number, not a customer-facing bug.
Why interviewers reach for this question: It is the purest test of whether a candidate can reason about anomalies as product decisions. Every candidate knows "leader-follower" and "quorum." Few can say which anomaly a given design exposes, which user sees it, and who signed off on that.
Mechanics Refresher: Replication Topologies
| Topology | How It Works | Pros | Cons |
|---|---|---|---|
| Single-leader, async followers | Leader commits locally, ships log to followers in background | Fast writes (1 local fsync); simple mental model | Failover can lose the last N ms of acked writes (RPO > 0); stale follower reads |
| Single-leader, sync / semi-sync | Leader waits for ≥1 follower ack before acking client | RPO = 0 for 1 failure | +1 RTT per write (1–2ms cross-AZ, 60–100ms cross-region) |
| Consensus-replicated (Raft/Paxos) per shard | Leader needs majority ack; automatic leader election with terms | Linearizable, automatic safe failover, no split brain | Majority RTT on every write; leader is a per-shard bottleneck |
| Multi-leader | Each region accepts writes; async cross-replication | Local-latency writes everywhere; region-independent | Write–write conflicts are guaranteed; needs a resolution policy |
| Leaderless (Dynamo-style) | Client/coordinator writes to W of N replicas, reads from R | High write availability; no failover step | Conflicts, read repair, sloppy quorums; R+W>N is not linearizable |
| Chain replication | Writes go head → … → tail; reads from tail | Strong consistency, high read throughput at tail | Latency = chain length; chain reconfiguration is subtle |
For most production systems: single leader per shard, replicated with Raft/Paxos (or semi-sync) inside a region, async followers across regions, plus session guarantees for reads. Reach for leaderless or multi-leader only when a named requirement — writes must succeed during a region partition — forces it.
Executive Summary
If you only read one section, read this. Every decision in the case study follows from the contrast and the fault lines below.
What This Interview Actually Tests#
A replicated data store is not a "how do I copy bytes" question. Everyone can draw a leader and two followers.
It is an anomaly-budgeting question that tests:
- Whether you name the consistency guarantee per operation instead of per system
- Whether you know which anomaly each topology exposes (stale read, lost update, write skew, resurrected delete)
- Whether you can state an RPO and RTO in numbers and name who signed off on them
- Whether failover is designed as the most dangerous routine operation you run — because it is
The key insight: Replication is a contract about which lies the system is allowed to tell, to whom, for how long. Staff engineers write that contract down; Senior engineers discover it in a post-mortem.
The L5 vs L6 Contrast — Start Here#
| Behavior | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| First move | Picks a topology ("leader + 2 followers") | Asks which operations need linearizability and which tolerate staleness; commits per operation | Asks which classes of data the org stores and publishes a consistency menu the platform supports |
| Consistency | "Strong consistency" or "eventual consistency" for the whole system | Mixes: linearizable for money/identity writes, read-your-writes for user sessions, eventual for feeds | Prices each tier ($/GB, p99, RPO) so product teams choose with their own budget |
| Quorums | "N=3, W=2, R=2, so it's consistent" | Knows R+W>N gives overlap, not linearizability; names sloppy quorum and concurrent-write gaps | Standardizes on consensus-backed shards so 30 teams never have to reason about quorum edge cases |
| Failover | "Promote a replica automatically" | Fencing tokens, lease expiry, RPO accounting of lost writes, human gate for cross-region | Game-days failover quarterly; tracks failover success rate as an org SLO; kills manual runbooks |
| Conflicts | "Last write wins" | Picks per data type: LWW for idempotent overwrite, CRDT for counters/sets, app merge for carts | Bans wall-clock LWW for financial data in the platform standard; audits exceptions |
| Ownership | "The DB team runs it" | Platform owns replication + failover; product owns consistency choice per API and its incidents | Redraws boundary: consistency level is part of the API contract, reviewed like a schema change |
Why "first move" separates levels
L5: Starts from topology. "Leader-follower, async replication, read from followers to scale." It is a reasonable design — and it silently commits the product to stale reads and lossy failover without anyone deciding that.
L6: Starts from operations. "Which reads must reflect the latest write? Balance checks and username uniqueness, yes. Profile page after my own edit, yes for me, not for others. Follower counts, no. That gives me three consistency tiers, and the topology falls out of the strictest one per shard."
L7: Recognizes the organization is about to have this conversation 40 times. "I don't want every team re-deriving quorum math. The platform should offer three named tiers — strict, session, eventual — each with a published p99, RPO, and price, and the choice goes in the API spec."
Why "quorums" separates levels
L5: "R + W > N means every read sees the latest write." This is the most common confident-but-wrong statement in the interview.
L6: "R + W > N guarantees the read set overlaps the write set for completed writes. It does not give linearizability: a read concurrent with an in-flight write can return new then old from different coordinators; sloppy quorums with hinted handoff break overlap entirely during partitions; and last-write-wins on skewed clocks can discard an acknowledged write. If I need linearizability, I use consensus per shard, not quorum arithmetic."
L7: Treats quorum tuning as an org hazard. "Every tunable is a footgun that a team will set to ONE during an incident and forget. The platform exposes tiers, not knobs."
Why "failover" separates levels
L5: "Health check fails, promote the most up-to-date replica." Correct shape, missing the two ways it kills you: the old leader is not actually dead (split brain), and the new leader is missing acked writes (data loss).
L6: "Failover requires three things: the old leader must be fenced — a monotonically increasing epoch that storage rejects if stale; the new leader must be chosen by a majority, not by a watchdog; and we must account for writes acked by the old leader but not replicated. With async cross-region replication that's up to the lag at failure time — I'll quote it as RPO ≤ 5s p99 and product has to sign off."
L7: "The failure I worry about isn't one failover. It's that we do 2 a year, so nobody's practiced it. I'd mandate quarterly regional evacuations and publish the failover success rate."
The Staff Positions#
| Position | Rationale |
|---|---|
| Consistency is chosen per operation, not per system | One user action spans reads with different staleness tolerance; a global setting over-pays or under-protects |
| Single leader per shard until a requirement forces multi-writer | Conflicts become impossible by construction; multi-leader buys write latency at the cost of a permanent conflict-resolution tax |
| Consensus (Raft/Paxos) inside a region, async across regions | Majority RTT inside a region is 1–3ms; cross-region majority is 60–200ms and rarely worth it for every write |
| Read-your-writes via commit tokens, not sticky sessions | Tokens survive load-balancer reshuffles and region routing; stickiness breaks on the first deploy |
| Never last-write-wins on wall clocks for data that matters | NTP skew of 10–100ms silently discards acknowledged writes; use versions, HLCs, or CRDTs |
| Failover is fenced, majority-elected, and rehearsed | The failover path is the least-exercised, highest-blast-radius code you own |
| Anti-entropy is not optional in any AP design | Read repair only fixes hot data; cold data diverges forever without Merkle-tree repair on a schedule |
The Three Intents#
Three intents produce three incompatible stores. Name one before drawing a box.
| Intent | Constraint | Strategy | Failure Mode | Correctness Bar |
|---|---|---|---|---|
| System of record (balances, orders, identity) | No acked write may be lost; no two users may both "win" | Consensus-replicated shards, sync majority commit, linearizable reads from leader/lease-holder | Unavailable (write-blocked) for the minority side of a partition | RPO = 0, linearizable, audited |
| User-facing read scale (profiles, catalogs, settings) | p99 read < 10ms from any region; users must see their own edits | Single leader per shard, async regional followers, session guarantees (RYW, monotonic reads) | Other users see stale data for ≤ replication lag (p99 < 1s) | Session consistency; RPO ≤ seconds, signed off |
| Always-writable (carts, presence, collaborative state, counters) | Writes succeed in every region even during partition | Multi-leader or leaderless, CRDTs / version vectors, anti-entropy | Conflicts on every concurrent edit; merge semantics visible to users | Convergence guaranteed; no acked write silently dropped |
🎯 Staff Move: "I'll design for user-facing read scale with a system-of-record core: single leader per shard, Raft inside the home region, async followers in other regions, and read-your-writes via commit tokens. The few operations that need linearizability — uniqueness checks, balance mutations — go to the leader. I'm deliberately not building multi-leader; if product needs writes during a regional partition, that's a separate conversation about which data types can merge."
The Five Fault Lines#
| # | Fault Line | The Tension |
|---|---|---|
| 1 | Leader-based vs Leaderless | One writer per key (no conflicts, failover risk) or any writer (no failover, permanent conflicts)? |
| 2 | Synchronous vs Asynchronous Replication | Pay a round-trip on every write for RPO = 0, or ack fast and accept losing the tail on failover? |
| 3 | Read Consistency Level | Linearizable (leader, slow, available only on majority side) vs session vs eventual (fast, stale)? |
| 4 | Conflict Resolution | LWW (simple, lossy) vs version vectors + app merge (correct, product burden) vs CRDTs (automatic, limited types)? |
| 5 | Failover Authority | Automatic (fast RTO, split-brain risk) vs human-gated (slow RTO, safe)? Who decides the region is dead? |
In the Wild: Real Production Systems#
Why this section belongs here: Naming the real design and the real reason it was built that way shows you have studied operational reality, not a textbook diagram.
Amazon Dynamo — Always-Writable Shopping Cart#
Amazon's 2007 Dynamo paper describes a leaderless, consistent-hashing store built so that the shopping cart would always accept writes, even during failures. It used tunable N/R/W quorums, sloppy quorums with hinted handoff to keep writing when preferred replicas were down, vector clocks to detect concurrent versions, and Merkle trees for anti-entropy. Conflicting cart versions were returned to the application, which merged them — the paper notes that deleted items could reappear as a result.
Staff insight: Dynamo is the canonical example of choosing an intent first. "Add to cart must never fail" was a business decision, and the resurrected-item anomaly was the explicitly accepted cost. Cassandra and DynamoDB are descendants, but DynamoDB (the service) hides most of this behind a single-leader-per-partition design — which tells you how much operational pain the original model carried.
Google Spanner — Paying Latency for External Consistency#
Spanner replicates each data split with Paxos across zones or regions and uses TrueTime — GPS and atomic clocks exposing an explicit uncertainty interval, typically a few milliseconds — to assign commit timestamps. Transactions "commit-wait" out the uncertainty before becoming visible, giving external consistency (linearizability for transactions) globally.
Staff insight: Spanner shows linearizability is purchasable, and shows the price: commit-wait of roughly the clock uncertainty (single-digit ms) plus a cross-region Paxos round (tens of ms) on multi-region configurations. When an interviewer says "can't we just have strong consistency everywhere?", the answer is "yes — here is the latency bill and the hardware bill."
Facebook TAO / Memcache — Regional Leader with Read-Your-Writes Patches#
Facebook's published TAO and memcache papers describe a design where one region is the leader for a given shard's writes, other regions hold async replicas and caches, and writes are forwarded to the leader region. To avoid users seeing their own writes disappear, the memcache design used "remote markers" to route a user's reads for recently written keys to the leader region until replication caught up.
Staff insight: At the largest scale, the answer was not stronger replication; it was a session guarantee patched on top of async replication for the one anomaly users actually notice — not seeing their own write. That is the exact Staff default in this case study.
What Interviewers Probe#
| After You Say... | They Will Ask... | (What They're Evaluating) |
|---|---|---|
| "N=3, W=2, R=2" | "Is that linearizable? What about concurrent reads during the write?" | Whether you know overlap ≠ linearizability |
| "Async replication to followers" | "Leader dies with 800ms of lag. What happened to those writes? Who was told?" | RPO literacy and ownership |
| "Automatic failover" | "The old leader was just GC-paused for 12 seconds. It wakes up. Now what?" | Fencing, epochs, split brain |
| "Last write wins" | "Two servers with 50ms clock skew write 10ms apart. Which one survives?" | Clock reasoning; silent data loss |
| "Read from followers" | "User edits their bio and refreshes. What do they see?" | Session guarantees |
| "Multi-region active-active" | "Two regions accept the same username. Now what?" | Whether you know uniqueness needs a single serialization point |
| "Read repair handles divergence" | "What about keys nobody reads for 6 months?" | Anti-entropy and tombstone/GC interaction |
System Architecture Overview#
Reading the diagram: Writes always go to the shard leader in the home region and commit when 2 of 3 in-region replicas persist them (2–4ms). Remote regions receive the log asynchronously (lag p99 < 1s). Reads carry the session's commit token: a remote follower serves the read only if it has applied at least that LSN, otherwise it waits briefly or proxies to the leader. Linearizable reads go to the leader holding a valid lease. The failover controller owns epochs in etcd; anti-entropy catches the divergence read repair never sees.
Quick-Reference: The 30-Second Cheat Sheet#
| Topic | The L5 Answer | The L6 Answer — Say This |
|---|---|---|
| Topology | "Leader + followers" | "Single leader per shard, Raft in-region, async cross-region. Multi-writer only for types that merge." |
| Consistency | "Strong" or "eventual" | "Per operation: linearizable for uniqueness and money, read-your-writes for sessions, eventual for aggregates." |
| Quorums | "R+W>N = consistent" | "Overlap, not linearizability. Sloppy quorums and LWW break it. If I need linearizable, I use consensus." |
| Failover | "Promote a replica" | "Majority-elected, fenced by epoch, RPO accounted. Cross-region failover is human-gated with a 5-minute decision SLO." |
| Conflicts | "Last write wins" | "LWW only for idempotent overwrites with HLC timestamps. CRDTs for counters/sets. App merge for carts. Never wall-clock LWW on money." |
| Divergence | "Read repair" | "Read repair for hot keys, Merkle anti-entropy for everything, and it must finish inside the tombstone GC window." |
Key Numbers Worth Memorizing#
| Metric | Value | Why It Matters |
|---|---|---|
| Intra-AZ round trip | ~0.1–0.5ms | Floor for any replicated write |
| Cross-AZ round trip (same region) | ~1–2ms | Cost of an in-region majority commit |
| Cross-region RTT: US east ↔ west | ~60–70ms | Cost of one synchronous cross-region ack |
| Cross-region RTT: US ↔ EU | ~80–100ms | Why global sync writes feel slow |
| Cross-region RTT: US ↔ APAC | ~150–200ms | Why global consensus is reserved for tiny, critical data |
| Raft in-region commit p99 | ~2–5ms | Includes fsync on majority |
| Async replica lag, healthy | p50 10–100ms, p99 < 1s | The size of your stale-read window and your RPO |
| Async replica lag, during backfill or vacuum | 10s – minutes | Why RPO must be measured, not assumed |
| NTP clock skew between hosts | 1–10ms typical, 100ms+ misconfigured | Why wall-clock LWW loses writes |
| Leader lease duration | 5–10s | Bounds unavailability after leader loss |
| Failure detection timeout | 3–10s in-region | Too low = flapping elections; too high = long write outage |
Cassandra default gc_grace_seconds | 864,000 (10 days) | Repair must complete inside it or deletes resurrect |
| Fault tolerance of a majority quorum | N=3 survives 1 loss; N=5 survives 2 | N=5 spread 2/2/1 across 3 AZs survives a full AZ loss, or any 2 nodes |
| Merkle tree repair of 1 TB replica | hours, throttled to ~10–50 MB/s | Schedule it; it competes with foreground I/O |
Interview Walkthrough
The most common mistake: Candidates spend 20 minutes on partitioning schemes and storage engines (LSM vs B-tree) and never reach the question the interviewer cares about — what users observe while replicas disagree, and what happens during failover. Compress the basics to 10–12 minutes.
Phase 1: Requirements & Framing (2–3 minutes)#
State the functional scope in one sentence:
"A key-value / document store with get, put, conditional put, and delete, replicated across three availability zones in a home region and readable from two other regions."
Then spend the time on the non-functional contract — this is where the design is decided:
"Three questions shape everything. First, what's the durability promise — can we lose any acknowledged write? Second, which reads must be fresh? Third, must writes succeed in a region that's partitioned from the home region? I'll assume: zero acked-write loss for a region-internal failure, a few seconds of RPO for losing a whole region, read-your-writes for users, and writes may fail on the minority side of a partition."
Commit to numbers:
"Targets: write p99 < 10ms in the home region, read p99 < 5ms from local region replicas, 99.99% read availability, 99.95% write availability, RPO = 0 for AZ loss, RPO ≤ 5s for region loss, RTO ≤ 30s for leader loss and ≤ 15 min for region evacuation."
🎯 Staff Move: Say the RPO and RTO in seconds and say who has to sign off on them. "RPO ≤ 5s on region loss means a customer can lose their last few seconds of edits in a regional disaster. Product and legal sign that, not the database team."
Phase 2: Core Entities & API (1–2 minutes)#
- Record:
key,value,version(per-key monotonic or HLC),tombstoneflag - Shard / Range: key range → replica set → leader,
epoch - Commit token:
(shard_id, log_index)returned on every write
PUT /kv/{key} body, If-Match: <version> → 200 { version, commit_token }
GET /kv/{key}?consistency=linearizable|session|eventual
X-Commit-Token: <token> → 200 { value, version, served_by, staleness_ms }
DELETE /kv/{key} If-Match: <version> → 200 { commit_token }
🎯 Staff Move: "Consistency level is a parameter on the read API and a documented field in each calling team's API spec. It's not a cluster-wide config. And every response returns
staleness_msso callers can see the lie they were told."
Phase 3: High-Level Architecture (≤5 minutes)#
Walk the flow in 90 seconds:
- Key hashes to a range; the router caches the range → leader map (epoch-versioned, from etcd).
- Writes go to the leader, append to the Raft log, commit on 2 of 3 in-region acks (2–4ms), return a commit token.
- The log ships asynchronously to remote-region followers.
- Reads:
linearizable→ leader with valid lease;session→ any replica whose applied index ≥ token;eventual→ nearest replica. - Anti-entropy runs in the background; failover controller owns epochs.
🎯 Staff Move: "This is the reasonable design. What makes it production-grade is three things: what a read sees during lag, what happens during failover, and how divergence is detected when nobody reads the data. Which do you want first?"
Phase 4: Transition to Depth (1 minute)#
"The topology is table stakes. The decisions that matter are: (1) where the linearizability boundary is and what it costs, (2) how failover avoids split brain and accounts for lost writes, and (3) what we do if product later demands writes during a region partition. I'd start with failover — it's where replicated stores actually lose data."
Phase 5: Deep Dives (25–30 minutes)#
Pattern for each: state the tradeoff → pick a position → quantify → name who pays.
Deep dive A: Failover without split brain (6–8 min)
"Leader loss detection: followers stop receiving heartbeats; after the election timeout — 3–5 seconds in-region — a follower with an up-to-date log starts an election for term T+1 and wins with a majority. The old leader, if it's merely paused, still believes it's leader. Two defenses: leader leases — the old leader stops serving linearizable reads when its lease (shorter than the election timeout, minus a clock-drift margin) expires; and fencing — every write carries the term, and replicas reject terms lower than the highest they've seen. A GC-paused leader that wakes up at term 41 cannot commit, because no majority accepts term 41 anymore."
"In-region RPO is zero: the new leader has every majority-committed entry by construction. Cross-region is different: if the entire home region is lost, remote followers are up to lag behind. At p99 lag < 1s we lose ≤ 1s of writes; during a replication stall we could lose minutes. That's why cross-region promotion is human-gated: someone must decide that losing the tail is better than waiting."
Deep dive B: Read-your-writes across regions (5–7 min)
"Every write returns (shard, log_index). The client library stores the max per shard in the session cookie. A remote follower receiving a read with token 88,121,934 checks its applied index. If it's caught up, serve. If not: wait up to 50ms, then proxy to the leader region (+80ms). At p99 lag of 1s, about 1–3% of post-write reads take the slow path — we measure it as session.token_wait_ms and session.leader_proxy_rate."
"Monotonic reads come free with the same token: the client never accepts a replica older than its high-water mark, so it can't see time go backwards when the load balancer bounces it between replicas."
Deep dive C: The quorum trap and conflict resolution (5–7 min)
"If the interviewer pushes toward leaderless: R+W>N guarantees that a read quorum intersects the write quorum of a completed write. It does not make concurrent operations linearizable, and with sloppy quorums during partitions the intersection doesn't hold at all. And once two writers can write the same key, I need a conflict policy. LWW with wall clocks silently drops the write with the smaller timestamp — with 20ms skew, a write made 15ms later can lose. For a cart I'd use an observed-remove set CRDT; for counters a PN-counter; for arbitrary documents, version vectors with siblings returned to the app for merge. Each has a product owner who must accept its semantics."
Deep dive D: Anti-entropy (3–5 min)
"Read repair fixes divergence only on keys that are read. Cold keys need scheduled Merkle-tree comparison: each replica builds a hash tree per range, compares roots, and descends only into differing subtrees — for 1 billion keys with a depth-15 tree we exchange ~32K leaf hashes to localize differences. It must complete within the tombstone GC window, or a replica that missed a delete will resurrect the data after the tombstone is purged elsewhere."
Phase 6: Wrap-Up (2–3 minutes)#
"Replication is a contract about which anomalies users can see. I put linearizability only where uniqueness and money live, session guarantees where users would notice, and eventual everywhere else — each with a published staleness number. The organizational point: consistency level is part of each team's API contract, reviewed like a schema change, and failover is rehearsed quarterly because a path we exercise twice a year is a path that doesn't work."
"What I'd build later: a multi-writer tier restricted to CRDT types for the two or three products that truly need partition-tolerant writes, and follower reads with bounded staleness (AS OF 10s) for analytics."
🎯 Staff Move: End on ownership and rehearsal, not on another component.
Common Timing Mistakes#
| Mistake | L5 Does This | L6 Does This Instead |
|---|---|---|
| Storage engine rabbit hole | 10 min on LSM compaction vs B-tree | "LSM, because write-heavy. Engine choice doesn't change the replication contract." |
| Partitioning over replication | Long consistent-hashing discussion | "Range partitioning with split/merge; consistent hashing is fine too. The interesting part is per-range replication." |
| One global consistency setting | "We'll use strong consistency" | Per-operation tiers with numbers |
| Failover as a bullet point | "Automatic failover" | Walks detection → election → fencing → RPO accounting → who is told |
| No anti-entropy | Stops at read repair | Merkle repair inside the GC window |
| No numbers | "Low replication lag" | "p99 < 1s, alert at 5s, RPO budget 5s" |
1. The Staff Lens#
1.1 Why This Problem Exists in Staff Interviews#
Replicated stores sit under every other system the company runs. A Senior engineer can make one correct; a Staff engineer decides which anomalies the whole product is allowed to expose and makes sure somebody owns each one. The question also reveals operational scar tissue fast: anyone who has lived through a failover that lost data answers differently from someone who has read about it.
The five behaviors interviewers listen for:
- Consistency named per operation, with a reason
- RPO/RTO stated in numbers with a named approver
- Failover walked through with fencing
- Quorum arithmetic stated precisely (overlap, not linearizability)
- Anti-entropy and deletes handled explicitly
1.2 The L5 vs L6 Contrast — Visual#
1.3 The Staff Question That Cuts Through Everything#
"The leader in us-east just stopped responding. Walk me through the next 60 seconds. Which writes are lost, which reads are wrong, who finds out, and how does the old leader learn it's no longer leader?"
A candidate who answers with terms, leases, fencing, and a number for lost writes has operated one of these. A candidate who says "we promote a replica" has not.
2. Problem Framing & Intent#
2.1 The Three Intents — Explained#
System of record → consensus, linearizable, write-unavailable on the minority
- Constraint: an acknowledged write is never lost; invariants (uniqueness, non-negative balance) never violated
- Mechanism: Raft/Paxos per shard, majority commit, leader leases, conditional writes (
If-Match) - Failure mode: the side of a partition without a majority rejects writes — a deliberate CP choice
- Who pays: users on the minority side see errors; the business pays in availability, not correctness
User-facing read scale → leader + async followers + session guarantees
- Constraint: local-region reads under 5ms; a user must see their own edits; others may lag
- Mechanism: commit tokens, monotonic reads, bounded-staleness follower reads
- Failure mode: stale reads for other users during lag; RPO of seconds on region loss
- Who pays: product owns the "others see it up to 1s late" decision; SRE owns the lag alert
Always-writable → multi-leader / leaderless + merge
- Constraint: writes accepted locally in every region, even when regions can't talk
- Mechanism: CRDTs, version vectors, hinted handoff, Merkle anti-entropy
- Failure mode: concurrent updates produce conflicts; resolution semantics become product behavior
- Who pays: product and end users — "the item I removed came back" is a merge semantic, not a bug
2.2 When NOT to Build a Replicated Data Store#
- When a managed service already fits. DynamoDB, Spanner, Aurora, Cosmos DB, and managed Cassandra encode years of failover hardening. Building your own replication layer is a 5–10 engineer, multi-year commitment. See Build vs Buy.
- When the data is derivable. Caches, search indexes, and materialized views can be rebuilt from a source of truth; replicate the source, not the derivative.
- When single-node + backups meets the RPO. A PostgreSQL primary with one sync standby and PITR backups gives RPO ≈ 0 and RTO of minutes for most internal tools. Don't build Raft for a 50 GB admin database.
- When you need global linearizable writes across continents for everything. That is not a replication design; that is a product requirement to challenge. Cross-continent consensus costs 150–200ms per write.
2.3 What the Interviewer Leaves Underspecified#
| Omitted Detail | Why It Matters | What to Ask |
|---|---|---|
| Durability promise | Decides sync vs async and RPO | "Can we ever lose an acknowledged write? For which failure — node, AZ, region?" |
| Read freshness per use case | Decides where reads go | "Which reads must see the latest write — everyone's or just the writer's?" |
| Partition behavior | Decides CP vs AP per operation | "During a region partition, should the cut-off region reject writes or accept and merge later?" |
| Write geography | Decides home-region vs multi-leader | "Where are writers? One region, or all of them with equal latency needs?" |
| Delete semantics | Decides tombstones, GC, compliance | "Are deletes legally required to be permanent within N days?" |
| Data types | Decides whether CRDTs are viable | "Is this counters, sets, documents, or balances?" |
Staff engineers surface these. Senior engineers assume them away.
2.4 Precise Terminology#
| Term | Precise Meaning | Common Misuse |
|---|---|---|
| Linearizable | Each operation appears to take effect atomically at one instant between its call and return; all clients agree on order | Used as a synonym for "durable" |
| Serializable | Transactions appear to execute in some serial order — not necessarily real-time order | Confused with linearizable; they're orthogonal |
| Read-your-writes | A session always sees its own prior writes | Assumed to follow from "eventual consistency" |
| Monotonic reads | A session never sees older data after seeing newer | Broken by load-balancing across replicas with different lag |
| Bounded staleness | Reads are at most Δ seconds or K versions behind | Promised without measuring lag |
| RPO | Maximum window of acknowledged data that can be lost | Stated as "zero" for async replication |
| RTO | Time until service restored after failure | Measured from detection, hiding detection time |
| Quorum | A subset whose size guarantees intersection with another subset | Treated as "strong consistency" |
| Sloppy quorum | Writes accepted by any N healthy nodes, not the designated N | Assumed to preserve R+W>N overlap (it doesn't) |
| Fencing token / epoch | Monotonic number that storage uses to reject stale leaders | Omitted from failover design |
| Tombstone | Marker recording a delete so replicas don't resurrect data | Purged before all replicas saw it |
🎯 Staff Insight: If the interviewer says "strong consistency," ask: "Do you mean linearizable single-key operations, serializable multi-key transactions, or just that users see their own writes? Those are three different systems with three different latency bills."
3. The Fault Lines#
Each fault line has a technical axis and an organizational axis. The organizational one — who approved the anomaly, who gets paged when users see it — is what the interviewer is scoring.
3.1 Fault Line 1: Leader-Based vs Leaderless#
The tension: A single leader per key makes conflicts impossible but makes failover a dangerous event. Leaderless removes failover but makes conflicts a permanent, everyday event.
| Choice | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Single leader per shard (consensus) | No write conflicts; linearizable available; simple reasoning for app teams | Write unavailability for 3–10s on leader loss; leader is per-shard throughput ceiling (~10–50K writes/s) | Platform on-call (elections, hot leaders) |
| Multi-leader (per region) | Local-latency writes everywhere; region survives partition | Every concurrent write to the same key conflicts; uniqueness impossible without a coordinator | Product teams (merge semantics), users (surprising merges) |
| Leaderless (Dynamo quorum) | No failover; writes succeed if any W nodes reachable | Read repair, hinted handoff, tombstone resurrection, LWW data loss | App teams (siblings), SRE (repair scheduling) |
L6 answer: "Single leader per shard is my default because it moves conflict resolution from every product team to one place — the log. I accept a 3–10 second write outage per shard on leader loss, which is rare and bounded. I'll only go multi-writer for data types that merge mathematically, and I'll say out loud that uniqueness constraints — usernames, seat assignments — can never be multi-writer."
When to deviate: Offline-first clients (mobile, collaborative editing — see Collaborative Editing), presence and counters across regions, or a business mandate that a region must keep taking writes while isolated.
🧭 Principal Move: "I'd make multi-writer an opt-in tier that only accepts registered CRDT types. If your data isn't one of those types, the platform won't let you replicate it multi-leader. That removes an entire class of incidents from 30 teams at once."
3.2 Fault Line 2: Synchronous vs Asynchronous Replication#
The tension: Every synchronous acknowledgment is a round trip added to every write. Every asynchronous one is data you can lose.
| Choice | Write Latency | RPO | Who Pays |
|---|---|---|---|
| Async everything | 1 local fsync (~0.5–1ms) | Replication lag at failure (ms → minutes) | Customers (lost writes), legal (if regulated) |
| Semi-sync (1 follower ack) | +1 cross-AZ RTT (~1–2ms) | 0 for single-node loss | Latency budget; stalls if the one sync follower is slow |
| Majority in-region (Raft) | +1–3ms | 0 for AZ loss | Platform (consensus ops) |
| Majority cross-region | +60–200ms | 0 for region loss | Every user on every write |
The Staff position is layered: majority in-region, async cross-region, and an explicit statement of the region-loss RPO.
L6 answer: "Commit on 2 of 3 in the home region: RPO zero for losing any single AZ, 2–4ms p99. Cross-region is async, so a full regional loss can drop up to the current lag. I'll alert at 5 seconds of lag, page at 30, and report rpo.unreplicated_writes so the RPO is a measured number, not a hope."
When to deviate: Regulated ledgers where RPO=0 for region loss is a legal requirement — then pay cross-region majority for that small dataset only (typically < 1% of total data), e.g., using a 5-replica group across 3 regions with the leader near the writers.
🎯 Staff Insight: "Semi-sync with a single designated follower is a trap. If that follower stalls, either writes stall or the system silently degrades to async — MySQL semi-sync does exactly this after a timeout. Majority of N ≥ 3 degrades gracefully; one designated follower does not."
3.3 Fault Line 3: Read Consistency Level#
The tension: Fresh reads come from the leader (one region, one node, capacity-limited). Scalable reads come from followers (stale). Session guarantees sit between them.
| Level | Served From | Latency | Staleness | Available During Partition? | Use For |
|---|---|---|---|---|---|
| Linearizable | Leader with valid lease (or ReadIndex round) | 1–3ms in-region, +RTT cross-region | 0 | Majority side only | Uniqueness checks, balance, locks, inventory decrement |
| Session (RYW + monotonic) | Any replica with applied index ≥ token | Local in ~97–99% of cases | 0 for own writes; ≤ lag for others | Yes, with fallback to leader | User-facing reads after edits |
| Bounded staleness | Any replica with lag ≤ Δ | Local | ≤ Δ (e.g., 10s) | Yes, until lag > Δ | Analytics, dashboards, search indexing |
| Eventual | Nearest replica | Local, < 1ms | Unbounded in theory | Yes | Counters, feeds, recommendations |
L6 answer: "Three tiers exposed as an API parameter. The default for user-facing reads is session, which means users never see their own write vanish and never see time go backward. Only operations that enforce an invariant pay for linearizable."
Who pays: Product pays in "another user may see my change up to 1s late." SRE pays in keeping lag inside the published bound. The platform pays in leader capacity for linearizable reads — budget roughly 5–10% of reads to be linearizable; if it's 50%, the leader is the bottleneck and you need to revisit.
🎯 Staff Move: "Linearizable reads via leader lease are only safe if the lease is shorter than the election timeout minus the maximum clock drift. With a 9s election timeout and a 500ms drift bound, I'd set the lease at 8s. If the platform can't bound drift, use Raft ReadIndex — one heartbeat round per read batch — instead of leases."
3.4 Fault Line 4: Conflict Resolution#
The tension: The moment two replicas accept writes for the same key independently, you must pick a winner or a merge. Every choice is a product decision disguised as a database setting.
| Strategy | How It Works | What Works | What Breaks | Who Pays |
|---|---|---|---|---|
| LWW, wall clock | Highest timestamp wins | Trivial, no siblings | Skew of 10–100ms drops acked writes silently | Users (vanishing edits), undetectable |
| LWW, HLC / logical | Hybrid logical clock, tiebreak by node ID | Causality respected; bounded skew | Still drops concurrent writes — just deterministically | Users, but explainable |
| Version vectors + siblings | Detect concurrency, return all versions | Nothing lost | App must merge; sibling explosion if clients don't | App teams |
| CRDTs | Types whose merge is commutative, associative, idempotent | Automatic convergence | Limited types; metadata overhead (tombstones in OR-sets can grow 2–10×) | Platform (implementation), product (semantics like "add wins") |
| Single-writer routing | Route all writes for a key to its home region | No conflicts | Cross-region write latency for non-home users | Remote users (+80–200ms) |
L6 answer: "Per data type. Profile fields: LWW on HLC — overwrite semantics, concurrent edits within a second are rare and the loser is explainable. Cart: OR-set CRDT, add-wins, which means a remove concurrent with an add loses — product signs off on that. Counters: PN-counter. Balances: never multi-writer; single-writer routed to the home region."
🎯 Staff Insight: "'Last write wins' is a euphemism for 'some acknowledged writes are silently discarded.' If you choose it, name the discard rate. For 1M writes/day to the same keys from two regions with ~20ms skew, the number of concurrent-within-skew writes is the number you're throwing away — measure it with a
conflict.lww_discardscounter."
3.5 Fault Line 5: Failover Authority#
The tension: Automatic failover gives fast RTO and risks split brain or promoting a replica that is missing data. Human-gated failover is safe and slow — and humans at 3am make bad calls under pressure.
| Choice | RTO | Risk | Who Pays |
|---|---|---|---|
| Automatic, watchdog-based (external health checker promotes) | 10–60s | Split brain under partition; promotes lagging replica | Customers (data divergence), on-call (reconciliation for days) |
| Automatic, consensus-based (Raft election) | 3–10s | Safe for in-region; can't help if majority is gone | Platform (consensus correctness) |
| Human-gated cross-region | 5–30 min | Slower; decision fatigue | Customers (downtime), incident commander |
| Automatic cross-region with RPO guard | 1–5 min | Only promotes if measured lag < threshold | Platform (guard logic) |
L6 answer: "In-region: automatic, consensus-driven, because a majority vote cannot produce two leaders in the same term. Cross-region: human-gated with a 5-minute decision SLO and a pre-computed lost-write estimate on the dashboard, because promoting the remote region is a one-way decision that forks history — when the old region returns, its unreplicated tail must be reconciled or discarded."
🧭 Principal Move: "GitHub's 2018 incident is the case study: a 43-second network partition led automation to promote a primary in another region, and the writes on both sides meant about 24 hours of degraded service while data was reconciled. The lesson isn't 'disable automation' — it's that cross-region promotion must check replication position and require the old side to be fenced first."
4. Failure Modes & Operational Reality#
4.1 Split Brain After a Paused Leader#
t=0: Leader L1 (term 41) enters 12s stop-the-world GC pause
t=+5s: Followers' election timeout fires; F2 wins term 42 with majority
t=+5.1s: New writes go to F2 (term 42). Router map updated to epoch 42.
t=+12s: L1 wakes. Still believes it is leader for term 41.
t=+12.01s: Stale router instance (cache TTL 30s) sends a write to L1.
L1 appends locally, sends AppendEntries(term 41).
Followers reply: term 42 > 41 → reject. L1 steps down.
Client gets error → retries → routed to F2. No divergence.
Without fencing (watchdog failover, no terms):
t=+12.01s: L1 accepts and acks the write locally. Two leaders.
t=+14min: Reconciliation discovers 3,400 conflicting rows.
Detection: raft.leader_count_per_shard > 1 (should be impossible — page immediately), raft.term_changes_per_hour, router.stale_epoch_rejections.
Blast radius: One shard per paused leader; with fencing, zero data divergence; without, every write in the overlap window.
Mitigation: Terms on every message; storage rejects stale terms; leases shorter than election timeout.
Prevention: GC tuning (pause p99 < 200ms), lease + fencing verified in chaos tests (Jepsen-style) before every major release.
Owner: Storage platform team.
4.2 Replication Lag Spiral#
t=0: Analytics team launches a backfill: 40K writes/s to 200 shards
t=+2min: eu-west followers' apply thread saturates; lag climbs 0.4s → 8s
t=+4min: Session reads from EU exceed token; 60% proxy to us-east (+85ms)
t=+5min: us-east leaders now serve 3× read load; CPU 85%; write p99 3ms → 40ms
t=+7min: Write latency slows apply further (shared disk). Lag 45s.
t=+8min: Page: replication.lag_seconds > 30 on 140 shards
Detection: replication.lag_seconds (warn 5s, page 30s), session.leader_proxy_rate (warn at 5%), leader CPU.
Blast radius: All cross-region session reads; eventually all writes via leader overload.
Mitigation: Throttle the backfill (per-tenant write rate limit); cap proxy rate — when > 20% of session reads would proxy, return stale-but-flagged data with staleness_ms or degrade the feature. See Degraded Mode.
Prevention: Bulk writes go through a separate low-priority class with admission control; parallel apply per shard.
Owner: Platform owns throttling; the analytics team owns their backfill schedule.
4.3 Resurrected Deletes (Zombie Data)#
t=0: User deletes account data (GDPR request). Tombstone written W=2 of 3.
t=0: Replica C is down for a disk replacement — misses the tombstone.
t=+11 days: gc_grace (10 days) expires; A and B compact the tombstone away.
t=+12 days: C rejoins. Anti-entropy compares A vs C: C has the row, A has nothing.
Repair copies the row BACK to A and B. The data is alive again.
t=+40 days: Privacy audit finds the "deleted" user's records. Regulatory exposure.
Detection: antientropy.last_full_repair_age_days per range (must be < gc_grace), node.downtime_seconds > gc_grace → block rejoin.
Blast radius: Every delete made while a replica was offline longer than the GC window.
Mitigation: A node offline longer than gc_grace must be wiped and rebuilt, never rejoined.
Prevention: Full repair cycle every 7 days on a 10-day grace; automation that refuses stale rejoin.
Owner: Storage platform; privacy/legal is a stakeholder in the SLO.
4.4 Silent Divergence (The Failure Nobody Pages For)#
A replica's disk returns a corrupted block; checksums at the page level catch it on read — but only for data that's read. A bug in a schema migration writes different defaults on leader vs followers for 2 hours. Nothing is down. Nothing alerts. Three months later a failover promotes that follower and customer-visible values change.
Detection: Periodic cross-replica checksum of ranges (antientropy.divergent_ranges, target 0), row-count comparisons per shard, sampled read-compare (read 0.1% of keys from two replicas and compare, emit consistency.sample_mismatch_rate).
Owner: Platform; the finding triggers an incident even with no user impact.
🎯 Staff Insight: "The dangerous replication failures are silent. I'd rather have a nightly job that says 'range 4412 differs between replicas' than find out during failover."
4.5 Operational Reality Matrix#
| Failure | Detection Signal | Blast Radius | Mitigation | Owner |
|---|---|---|---|---|
| Leader crash | raft.leader_elections_total spike | 1 shard, 3–10s writes | Automatic election | Storage platform |
| Paused leader / split brain | raft.leader_count_per_shard > 1 | 1 shard | Term fencing, leases | Storage platform |
| AZ loss | az.replica_unreachable | Up to 1/3 of replicas | Majority survives; rebuild in other AZ | Storage platform + infra |
| Region loss | region.health, lag frozen | All home-region writes | Human-gated promote; RPO = lag | Incident commander + platform |
| Replication lag spiral | replication.lag_seconds > 30 | Session reads, leader load | Throttle bulk writers, cap proxy | Platform + offending team |
| Resurrected deletes | antientropy.last_full_repair_age_days | Deleted data returns | Wipe stale nodes | Platform + privacy |
| LWW discards | conflict.lww_discards | Concurrent writers' data | Move data type to CRDT/single-writer | Product owning the data |
| Silent divergence | consistency.sample_mismatch_rate | Unknown until failover | Repair from majority | Platform |
| Hot shard leader | shard.leader_cpu > 80% | 1 shard's keys | Split range, move leadership | Platform |
5. Evaluation Rubric#
5.1 Level-Based Signals#
| Dimension | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| Consistency model | One setting for the system | Per-operation tiers with latency and staleness numbers | Org-wide consistency menu with published price, p99 and RPO per tier |
| Topology | Leader-follower | Consensus per shard in-region, async cross-region, justified by RTTs | Decides which topologies the platform supports at all; retires the rest |
| Quorums | R+W>N = consistent | Overlap ≠ linearizable; sloppy quorum and LWW gaps | Removes tunables from app teams; tiers instead of knobs |
| Failover | Promote a replica | Terms, fencing, leases, RPO accounting, human gate cross-region | Failover success rate as an SLO; quarterly region evacuations |
| Conflicts | LWW | Per data type: HLC-LWW, CRDT, app merge, single-writer | Platform standard bans wall-clock LWW for money; exception process |
| Anti-entropy | Read repair | Merkle repair inside tombstone GC window | Divergence detection as a first-class compliance control |
| Ownership | DB team owns it | Platform owns mechanics; product owns chosen anomalies | Consistency level reviewed as part of API contract; cost charged back |
5.2 Strong Hire Signals#
| Signal | What It Sounds Like |
|---|---|
| Per-operation consistency | "Username uniqueness is linearizable; bio reads are session; follower counts are eventual." |
| Precise quorum reasoning | "R+W>N gives overlap for completed writes. Concurrent reads can still flip-flop, and sloppy quorums break it." |
| Fencing named unprompted | "The old leader is rejected by term, not by a health check." |
| RPO in numbers with an approver | "Region loss costs up to 5s of writes — product and legal signed that." |
| Deletes considered | "Repair must finish inside gc_grace or deletes come back." |
5.3 Lean No-Hire Signals#
| Signal | Why It Misses the Bar |
|---|---|
| "Strong consistency everywhere" | Ignores the cross-region latency bill; no per-operation reasoning |
| "Eventual consistency is fine" without defining for whom | Hands users anomalies nobody approved |
| Automatic failover with no fencing | Designs a split-brain generator |
| Multi-leader for usernames | Doesn't understand that uniqueness requires a single serialization point |
| No numbers for lag or RPO | Can't be held to a contract |
5.4 Common False Positives#
- Recites Raft in detail: Knowing AppendEntries fields ≠ knowing where linearizability is needed.
- Name-drops CAP: "It's CP" says nothing about latency (PACELC) or which operations.
- Draws five regions active-active: Complexity is not a signal; ask what happens to a uniqueness check.
- Knows Cassandra consistency levels: Tuning knobs without a product rationale is Senior.
6. Interview Flow & Pivots#
6.1 Typical 45-Minute Shape#
| Phase | Time | Goal |
|---|---|---|
| Framing | 0–3 min | Durability promise, freshness per use case, partition behavior; RPO/RTO numbers |
| Entities + API | 3–5 min | Record, shard, commit token; consistency as a read parameter |
| High-level design | 5–12 min | Consensus in-region, async cross-region, router |
| Transition | 12 min | Offer failover, read consistency, multi-writer |
| Deep dives | 12–40 min | Failover + fencing → session guarantees → conflicts/anti-entropy |
| Wrap-up | 40–45 min | Anomaly contract, ownership, rehearsal, what's next |
6.2 How Interviewers Pivot — And What They're Testing#
| Interviewer Signal | What They Care About | Where to Go Deep |
|---|---|---|
| "What if the leader dies?" | Operational maturity | Terms, fencing, RPO accounting |
| "Can we read from replicas?" | Consistency literacy | Session tokens, bounded staleness |
| "Make writes work in every region" | Conflict reasoning | CRDTs vs version vectors, uniqueness caveat |
| "Is R+W>N linearizable?" | Precision | Overlap vs linearizability, sloppy quorum |
| "How do you know replicas agree?" | Silent failure awareness | Merkle anti-entropy, sampled compare |
| "What about GDPR deletes?" | Compliance and tombstones | GC window, stale-node rebuild |
6.3 What to Deliberately Skip#
| Topic | Why L5 Goes Here | What L6 Says Instead |
|---|---|---|
| LSM vs B-tree internals | Feels deep | "LSM for write-heavy; doesn't change the replication contract." |
| Paxos vs Raft proofs | Feels rigorous | "Raft for understandability; both give the same guarantees here." |
| Consistent hashing ring math | Familiar | "Range partitioning with splits; see sharding. Replication is the question." |
| Client SDK details | Easy to list | "Client carries the commit token; that's the only SDK requirement." |
6.4 Follow-Up Questions to Expect#
- "Is R=2, W=2, N=3 linearizable? Prove or disprove."
- "The leader was paused for 15 seconds. What does it do when it wakes?"
- "How do you give read-your-writes to a user who switches from phone to laptop?"
- "Two regions accept a write to the same key 5ms apart. Which wins, and does anyone find out?"
- "How do you detect that two replicas silently disagree?"
- "How long can a replica be offline before it can't safely rejoin?"
- "What's your RPO for losing us-east entirely, and who approved it?"
7. Active Drills#
Drill 1: The Opening#
Prompt: "Design a replicated key-value store for our product."
Staff Answer
"Before topology, three questions: can we ever lose an acknowledged write, and for which failure — node, AZ, or region? Which reads must see the latest write, and for whom? And must a region cut off from the others keep accepting writes? I'll assume: no loss for node or AZ failure, up to 5 seconds of loss for a full-region disaster with product sign-off, read-your-writes for users, and writes may fail on the minority side of a partition. That gives me single leader per shard, Raft across three AZs in a home region, async followers in two remote regions, commit tokens for session consistency, and linearizable reads only for invariant-enforcing operations. I'll go deep on failover, then read consistency, then what changes if product later needs partition-tolerant writes."
Why this is L6:
- Turns "replicated" into three concrete contract questions before drawing anything
- States RPO per failure domain and names the approver
- Picks the topology from the strictest requirement rather than by habit
What L7 adds:
- "This is the fourth team this year asking for this. I'd rather give them a tier from the platform menu than a bespoke design."
- Frames the RPO decision in dollars: expected loss × regional-disaster probability vs cost of cross-region sync
Drill 2: The Quorum Question#
Prompt: "N=3, W=2, R=2. Is that strongly consistent?"
Staff Answer
"It guarantees that any read quorum overlaps any completed write quorum, so a read after a completed write will see at least one replica with that write. It is not linearizable. Three counterexamples: (1) A write reaches replica A only and is still in flight. Reader 1 reads A and B, sees the new value. Reader 2, starting after Reader 1 finishes, reads B and C and sees the old value — new then old, a linearizability violation. Read repair before returning fixes some of this at the cost of a write on the read path. (2) Sloppy quorum: during a partition, W=2 can be satisfied by fallback nodes outside the key's preference list, so the read quorum may not intersect at all. (3) LWW with clock skew: a write with a lower timestamp arriving later is discarded even though it was acknowledged. If the product needs linearizable, I use consensus per shard — not a bigger quorum."
Why this is L6:
- Precise about what quorums guarantee instead of reciting the inequality
- Gives constructive counterexamples an interviewer can verify
- Moves to the right tool (consensus) instead of tuning knobs
What L7 adds:
- "I wouldn't expose R and W to application teams at all — someone will set W=1 during an incident and forget. Tiers, not knobs."
- Notes the testing investment: Jepsen-style linearizability checks in CI for the platform
Drill 3: Make Read-Your-Writes Concrete#
Prompt: "A user edits their profile on their phone, then opens the laptop. What do they see, and how?"
Staff Answer
"Every write returns a commit token (shard, log_index). Tokens are stored server-side against the user ID in a small session store with a 60-second TTL — not in the device cookie — so they follow the user across devices. On read, the API layer attaches the user's max token for the shard. A follower serves only if applied_index ≥ token; otherwise it waits up to 50ms, then proxies to the leader. The 60-second TTL works because lag p99 is under 1s — after 60s any replica is caught up unless lag is pathological, in which case the lag alert has already fired. Cost: one session-store lookup (~0.5ms) per read for users who wrote in the last minute, and 1–3% of their reads proxied at +80ms."
Why this is L6:
- Solves the cross-device case that cookie-based tokens miss
- Bounds the mechanism with a TTL derived from measured lag
- Prices the guarantee in latency and proxy rate
What L7 adds:
- Makes session consistency a platform primitive other stores (cache, search index) honor too — otherwise the user sees their write in the DB but not in search
- Defines the org metric:
session.guarantee_violationssampled from client-side checks, reported per product
Drill 4: The Leader Went Quiet#
Prompt: "The leader for shard 7 stops heartbeating. Walk me through the first 30 seconds."
Staff Answer
"t=0 last heartbeat. Followers wait for the randomized election timeout, 3–5s. t≈4s: follower B becomes candidate for term 42, requests votes; C grants it because B's log is at least as up-to-date. B is leader at t≈4.1s and immediately appends a no-op entry to commit prior-term entries. The router learns the new leader via etcd watch within ~100ms, bumping the epoch. Writes during 0–4s fail or time out; clients retry with the same idempotency key. If the old leader A was only paused, its lease expired at t≈3s — before B could be elected — so A stopped serving linearizable reads. When A wakes and sends term 41 messages, they're rejected, and A steps down. RPO: zero, because B had every majority-committed entry. Alerts: raft.leader_elections_total increments; a single election is a log line, 5 per hour on one shard pages."
Why this is L6:
- Walks a timeline with real timeouts and the ordering constraint lease < election timeout
- Remembers the no-op commit and idempotent retries
- Distinguishes a normal event from a page-worthy pattern
What L7 adds:
- Tracks election rate across the fleet as a leading indicator of GC, network, or noisy-neighbor issues org-wide
- Owns the chaos-testing program that proves this timeline on every release
Drill 5: Hot Shard#
Prompt: "One celebrity account's shard leader is at 95% CPU. What do you do?"
Staff Answer
"First, reads or writes? If reads: most of them don't need the leader. Move them to session or eventual consistency on followers — three in-region followers triple read capacity immediately, and remote-region followers add more. If linearizable reads are forced, batch them with ReadIndex so one heartbeat round serves many reads. If writes: a single key's writes serialize on one leader, and splitting the range doesn't help a single key. Options are changing the data model — e.g., sharding a counter into 16 sub-keys summed on read — or accepting a per-key write rate limit around 5–10K/s. Separately, move other ranges' leadership off that node so the hot key isn't co-located with 200 other shards' leadership."
Why this is L6:
- Separates read heat from write heat — different fixes
- Knows range splitting doesn't fix a single hot key
- Uses leadership placement as a lever
What L7 adds:
- Asks why the data model allows a single-key write hotspot in the first place; adds a hot-key review to the platform onboarding checklist
- Sets a published per-key write ceiling so product designs around it rather than discovering it
Drill 6: Multi-Region Writes Request#
Prompt: "Product wants EU users to write with EU latency, even if the transatlantic link fails."
Staff Answer
"Two options before multi-leader. First, re-home: if EU users mostly write EU-owned data, make the EU region the home for those shards — geo-partitioned leadership gives local-latency writes with single-leader semantics. That covers most cases. Second, for data that is genuinely shared across regions and must accept writes during a partition, we need merge semantics. I'd go through the data types: counters → PN-counter, tags/sets → OR-set, simple fields → HLC-LWW with discard metrics. Anything with a uniqueness or non-negative invariant — usernames, inventory, balances — stays single-writer, and during a partition the EU side returns a clear error for those operations. I'd want product to sign off on 'removed item reappears if removed concurrently with an add' before shipping the OR-set."
Why this is L6:
- Tries geo-partitioned single-leader first — cheaper than multi-leader
- Refuses multi-writer for invariant data and says what users will see
- Turns merge semantics into explicit product sign-off
What L7 adds:
- Quantifies: multi-writer tier costs ~2× platform engineering and carries permanent support load; limit it to named data types via the platform standard
- Notes data-residency law may force EU-home anyway — legal becomes a stakeholder in the topology
Drill 7: Build vs Buy#
Prompt: "Should we build this or use a managed database?"
Staff Answer
"Buy unless one of three things is true: our consistency or placement requirements aren't offered (e.g., per-operation session tokens across regions), the managed cost at our scale exceeds roughly 3× the fully loaded cost of a team to run our own, or we have a strategic reason to control the failure behavior. For a typical product — DynamoDB global tables or Spanner or Aurora Global gives us single-digit-ms in-region writes, managed failover, and point-in-time recovery. Building a replicated store correctly is a 6–10 engineer team for 2+ years before it's trustworthy, plus a permanent on-call. I'd start managed and design the API so the consistency tier is ours, not the vendor's, to keep the exit door open."
Why this is L6:
- Gives crisp criteria, including a cost multiple
- Names realistic team size and time-to-trust
- Protects the exit via an abstraction at the contract layer
What L7 adds:
- Evaluates vendor lock-in as a one-way door: data egress cost (~$0.02–0.09/GB) at PB scale is a seven-figure migration
- Considers the org's ability to hire and retain storage engineers — the real constraint for building
Drill 8: Changing Consistency Without an Outage#
Prompt: "We want to move the orders service from follower reads to linearizable reads. How?"
Staff Answer
"Capacity first: linearizable reads all land on leaders. If orders is 20K reads/s spread over 3 replicas, leaders now take the full 20K plus writes — check leader headroom per shard. Then shadow: for 1% of reads, issue both a follower and a leader read and log mismatches (consistency.shadow_mismatch_rate) — that tells us how often users were seeing stale orders, which justifies the change. Roll out per endpoint behind a flag: 1% → 10% → 50% → 100% over a week, watching leader CPU and p99. Rollback is a flag flip. The orders team owns the decision and the latency impact; platform owns leader capacity."
Why this is L6:
- Checks the capacity consequence before the correctness benefit
- Uses shadow reads to quantify the anomaly being fixed
- Staged rollout with a named owner and a one-flag rollback
What L7 adds:
- Uses the shadow data across all teams to publish "stale-read rate by service" — turns an invisible risk into an org dashboard
- Charges leader-read capacity back to the requesting team so linearizable reads aren't free
Drill 9: Cost#
Prompt: "Replication is 40% of our storage bill. Cut it."
Staff Answer
"Break the bill into copies × bytes × tier. Typical: RF=3 in-region × 3 regions = 9 full copies. Levers: (1) Remote regions don't need RF=3 for the data they only read — RF=2 in remote regions, or a single remote copy plus fast rebuild from the home region, cuts copies from 9 to 5–7. (2) Cold data: move ranges untouched for 90 days to erasure-coded object storage (~1.5× overhead instead of 3×). (3) Witness replicas: in a Raft group, a log-only witness participates in quorum without storing the full dataset. (4) Compression at the log shipping layer cuts cross-region transfer, which at $0.02/GB is often larger than storage. Each lever changes RTO — rebuilding a remote region from 1 copy takes hours — so the region RTO has to be re-signed."
Why this is L6:
- Decomposes the cost instead of guessing
- Knows witness replicas and erasure coding as concrete levers
- Ties each saving to the reliability promise it weakens
What L7 adds:
- Tiers replication by data class org-wide (gold/silver/bronze), with the bill visible to the owning team
- Treats the cross-region transfer contract with the cloud vendor as a negotiable line item
Drill 10: Multi-Region Expansion#
Prompt: "We're adding ap-southeast. What changes?"
Staff Answer
"Reads: an async follower set in ap-southeast, lag budget p99 < 1.5s given ~180ms RTT to us-east. Session reads from APAC that miss their token proxy to us-east at +180ms, so the proxy rate matters more — I'd raise the wait-before-proxy to 150ms there. Writes: APAC users writing to us-east-home shards pay ~180ms. If APAC write volume is material, geo-home the APAC-owned shards in APAC. Bootstrapping: seed the new followers from a snapshot plus log catch-up, not a full stream, and throttle to protect home-region leaders. Anti-entropy cost grows linearly with replicas; schedule repair per region pair. Failover: APAC should not be a promotion target for us-east shards unless its lag is the lowest — the RPO guard handles that."
Why this is L6:
- Recalibrates timeouts to the new RTT instead of copy-pasting
- Considers bootstrapping load on the leaders
- Re-evaluates homing and failover targets
What L7 adds:
- Asks whether data residency law in APAC markets forces local-home data, which changes the topology more than latency does
- Adds the new region to the quarterly evacuation drill before it takes production writes
8. Deep Dive Scenarios#
Deep Dive 1: Peak Traffic — The Lag Spiral on Launch Day#
Context: A major product launch drives write volume to 3× normal. Remote-region replication lag climbs to 40 seconds. EU users are complaining that their changes "disappear" and support tickets are spiking. The on-call escalates to you.
Questions to Surface First:
- Are session tokens being honored, or is the EU path falling back to eventual reads under load?
- Is lag caused by ship (network), apply (CPU/disk), or leader-side backlog?
- What fraction of session reads are now proxying to us-east, and what is that doing to leader CPU?
- Is any non-launch workload (backfill, reindex) contributing?
Typical L5 Approach: Scales up the EU replicas, increases apply threads, and monitors lag until it recovers. Correct mechanics, but misses why users saw their writes disappear — lag alone shouldn't cause that if read-your-writes is working.
Staff Approach: Treats "writes disappear" as a session-guarantee violation, not a lag problem. Discovers the proxy path had a circuit breaker that, above 30% proxy rate, fell back to local eventual reads silently. Fixes the degraded mode to return an explicit "saving…" state instead of stale data, throttles non-critical writers, and adds apply parallelism.
Principal Approach: Asks why the platform lets a degraded mode silently break a published guarantee. Mandates that every consistency tier has a declared degraded behavior (error, wait, or flagged-stale) visible in the API response, and introduces launch-readiness reviews that include replication headroom for any launch forecast > 2× baseline.
Staff Approach — Full Reasoning
| Phase | What to Do |
|---|---|
| Immediate (0–5 min) | Check session.leader_proxy_rate and session.fallback_to_eventual_total. If fallback > 0, users are being shown stale own-writes — that's the incident. |
| Triage | Split lag into ship lag vs apply lag. Apply lag with idle network = CPU/disk on followers. Identify top write sources by tenant. |
| Quick fix | Pause the reindex job; raise follower apply parallelism from 1 to 8 per shard (independent key ranges); raise proxy circuit threshold temporarily with leader CPU watched. |
| Guardrails | Leader CPU < 75%; write p99 < 10ms. If breached, prefer returning HTTP 409 "still saving, retry" over stale reads. |
| Post-mortem | The silent fallback violated a published contract. Every degraded mode gets a metric and a product-approved behavior. |
Metrics to Watch: replication.lag_seconds{region}, replication.apply_queue_depth, session.leader_proxy_rate, session.fallback_to_eventual_total, shard.leader_cpu, support.tickets_tagged_missing_edit.
Organizational Follow-up: Launch-readiness template gains a "replication headroom" line. Bulk jobs get a priority class that auto-pauses when lag > 5s.
Ownership Question: "Who decides that showing stale data is better than showing an error?" Staff answer: The product owner of the surface, in advance, recorded in the API spec. The platform implements whichever degraded mode was chosen and alerts when it activates. On-call should never be making that product decision at 3am.
Key Takeaway: "Lag is a capacity problem. Users seeing their own writes vanish is a contract violation. Page on the second, not the first."
What clears the Staff bar:
- Separates lag from guarantee violation
- Finds the silent degraded mode
- Assigns the degraded-mode decision to product, in advance
Deep Dive 2: Silent Divergence Found During Failover#
Context: A planned failover of shard group 12 to a new AZ completes cleanly. Within an hour, customers report account settings reverting to values from two months ago. No alert fired during the failover.
Questions to Surface First:
- Did the promoted replica actually have the same data as the old leader? When was it last verified?
- Was there a schema migration or bug window in the past two months affecting only some replicas?
- Is the old leader still intact? Can we diff?
- How many other shard groups have never been cross-verified?
Typical L5 Approach: Fails back to the old leader, restores the missing values from backup, and adds a pre-failover check comparing row counts. Fixes the incident; row counts would not have caught value-level divergence.
Staff Approach: Fails back, then diffs old leader vs promoted replica by range using Merkle hashes to find every divergent key. Finds a migration two months ago that applied a default only on the leader (statement-based replication of a non-deterministic function). Adds continuous checksum comparison and blocks failover to any replica whose range hashes haven't matched within 24 hours.
Principal Approach: Treats "replicas can silently disagree" as an org-level integrity risk, not one team's bug. Bans non-deterministic statement-based replication in the platform, requires row-based or log-based replication, and makes
antientropy.divergent_ranges = 0a precondition for any planned failover org-wide.
Staff Approach — Full Reasoning
| Phase | What to Do |
|---|---|
| Immediate (0–5 min) | Freeze further failovers. Fail back to the old leader if it is intact and fenced correctly (bump epoch). |
| Triage | Merkle diff old vs promoted per range; list divergent keys and their last-modified times; correlate with deploy/migration history. |
| Quick fix | Repair promoted replica from old leader's values for divergent keys; customer comms for affected accounts. |
| Guardrails | Pre-failover check: range hashes match within the last 24h; otherwise failover requires director approval. |
| Post-mortem | Root cause: non-deterministic migration under statement replication. Fix: row-based replication, migration linting. |
Metrics to Watch: antientropy.divergent_ranges, antientropy.last_verified_age_hours{replica}, consistency.sample_mismatch_rate, failover.precheck_failures.
Organizational Follow-up: Migration review checklist adds "deterministic across replicas?"; platform CI rejects migrations using NOW()/RAND() in statement-replicated paths.
Ownership Question: "Who owned verifying that replicas matched?" Staff answer: Nobody — that's the finding. Replication correctness was assumed. The platform team now owns a continuous verification SLO, and failover tooling enforces it.
Key Takeaway: "A replica you've never verified is a backup you've never restored. Verify continuously or find out during failover."
What clears the Staff bar:
- Uses value-level verification, not row counts
- Finds the process root cause (non-deterministic migration)
- Makes verification a failover precondition
Deep Dive 3: Large Customer Onboarding With a Residency Requirement#
Context: Sales signs an enterprise customer 4× larger than any existing tenant. Contract terms: data must stay in the EU, RPO = 0 for region loss, and 99.99% write availability. They go live in 6 weeks.
Questions to Surface First:
- RPO=0 for region loss requires synchronous cross-region replication — within the EU, which regions? What's the RTT between them?
- 99.99% write availability plus RPO=0 across regions is a CAP-shaped promise — who approved it?
- Will this tenant share shards with others? What's their per-key write hotspot profile?
- Does the contract have penalties? What's the dollar exposure?
Typical L5 Approach: Provisions EU shards for the tenant, enables synchronous replication to a second EU region, load-tests, and ships. Delivers the letter of the contract, but under a regional partition the tenant's writes block — breaking 99.99% writes — and nobody told sales.
Staff Approach: Designs a 5-replica Raft group across 3 EU regions (e.g., Frankfurt, Ireland, Paris; 10–25ms RTT) so any single region can fail with RPO=0 and writes continue — the majority survives. Write p99 becomes ~25–35ms; the customer is told. Dedicated shards for isolation. Surfaces to sales that "99.99% writes and RPO=0" is only achievable with a 3-region footprint and costs ~2.2× a standard tenant.
Principal Approach: Treats this as a pricing and contract-governance gap. Creates a "sovereign strict" tier in the platform menu with a list price, a published latency, and a pre-approved architecture, and requires that sales contracts reference tiers rather than raw RPO/SLA numbers — so the next deal can't promise a combination the platform can't deliver.
Staff Approach — Full Reasoning
| Phase | What to Do |
|---|---|
| Week 1 | Confirm contract language with legal; map RPO/SLA promises to an architecture; get written acceptance of the 25–35ms write latency. |
| Weeks 2–3 | Provision 5-replica groups across 3 EU regions; dedicated shard range; leader placement near the customer's primary writers. |
| Weeks 4–5 | Load test at 2× forecast; chaos test region isolation; verify writes continue with RPO=0 when one region is cut. |
| Week 6 | Shadow traffic, then cutover with rollback plan. |
| Ongoing | Per-tenant SLO dashboard; contract SLA reporting from the same metrics. |
Metrics to Watch: tenant.write_p99{tenant}, raft.commit_latency{group}, region.isolation_drill_result, tenant.sla_availability_30d.
Organizational Follow-up: Sales engineering gets a tier catalog; custom RPO/SLA terms require platform sign-off before signature.
Ownership Question: "Who is accountable if the SLA is missed?" Staff answer: The platform owns meeting the published tier SLO; the account team owns having sold a tier that exists. When a contract promises something outside the catalog, the approver of that exception owns the gap.
Key Takeaway: "RPO=0 for region loss and high write availability together require at least three regions. Say the latency and the price before the contract is signed."
What clears the Staff bar:
- Translates contract language into topology and latency
- Knows 2 regions can't give both RPO=0 and availability under partition
- Pushes the decision back upstream to sales/legal
Deep Dive 4: Post-Mortem — The Cart That Lost Items#
Context: After a multi-region rollout of the cart service to active-active, customers report items vanishing from carts at a rate of ~0.3% of carts per day. The team uses LWW with server wall-clock timestamps.
Questions to Surface First:
- What's the clock skew distribution across regions?
- How often do two regions write the same cart within the skew window (mobile + web, retries)?
- Is "vanishing" a lost add, or a resurrected remove?
- What's the revenue impact of a lost cart item?
Typical L5 Approach: Tightens NTP, adds a node-ID tiebreak, and moves on. Reduces the rate but concurrent writes within the remaining skew are still silently discarded.
Staff Approach: Measures
conflict.lww_discardsand confirms they match the report rate. Replaces whole-cart LWW with an OR-set CRDT of line items (add-wins), quantity as a per-item PN-counter, and product sign-off that a concurrent remove + add results in the item present. Lost adds drop to zero.
Principal Approach: Uses this as the forcing function to publish the platform's conflict-resolution standard: wall-clock LWW is prohibited for any user-generated collection; multi-writer tier accepts only registered CRDT types; any team enabling multi-region writes goes through a design review with a data-type checklist.
Staff Approach — Full Reasoning
| Phase | What to Do |
|---|---|
| Immediate | Route each cart's writes to a single home region (by user) as a stopgap — conflicts go to zero, remote users pay +80ms on cart writes. |
| Triage | Instrument discards; confirm that mobile and web writing from different regions within 30ms account for the losses. |
| Fix | OR-set CRDT for items; PN-counter for quantity; tombstone GC tied to causal stability. |
| Guardrails | Shadow-merge old and new cart representations for 2 weeks; alert on divergence. |
| Post-mortem | Root cause: data type with add/remove semantics replicated with overwrite semantics. |
Metrics to Watch: conflict.lww_discards, cart.item_lost_reports, crdt.tombstone_bytes_per_cart, cart.write_p99{region}.
Organizational Follow-up: Conflict-resolution checklist for every multi-writer adoption; product sign-off template for merge semantics.
Ownership Question: "Who should have caught this before launch?" Staff answer: The design review. The data type (a set) and the replication policy (overwrite) were mismatched on paper; no amount of testing in a single region would have shown it.
Key Takeaway: "LWW on a collection is a bug with a probability. Replicate sets as sets."
What clears the Staff bar:
- Quantifies the discard rate and ties it to the symptom
- Picks CRDT semantics per field, not per object
- Uses single-homing as a fast, safe stopgap
Deep Dive 5: Multi-Region Expansion — From One Region to Three#
Context: The company runs everything in us-east with 3-AZ replication. Leadership wants a second and third region within 12 months for disaster recovery and latency.
Questions to Surface First:
- Is the goal DR (RPO/RTO for regional disaster), latency for remote users, or data residency? Each implies a different topology.
- Which services can tolerate async RPO, and which need synchronous?
- Has the org ever done a regional evacuation? What breaks outside the database — DNS, config, secrets, queues?
- What's the budget? Three regions is ~2.5–3× infrastructure cost.
Typical L5 Approach: Sets up async replicas of every database in two new regions, adds DNS failover, and declares DR done. Nothing proves failover works, and every service's RPO is whatever the lag happens to be.
Staff Approach: Phases it: (1) read replicas in region 2 with measured lag and session tokens; (2) monthly failover drills of non-critical shards; (3) geo-homing of shards whose writers are local; (4) region 3. Each service declares its tier and RPO; the dashboard shows real lag against declared RPO.
Principal Approach: Frames the program as an org capability, not a database project. Builds the evacuation playbook across all dependencies (identity, config, queues, caches), sets a target "region evacuation in < 30 min" validated by quarterly game days, and ties every team's tier declaration to their incident budget. Stops any service without a declared tier from being deployed to the new region.
Staff Approach — Full Reasoning
| Phase | What to Do |
|---|---|
| Months 0–3 | Region 2 async followers; commit tokens; lag and rpo.unreplicated_writes dashboards. |
| Months 3–6 | Failover drills on low-risk shards monthly; fix what breaks (it will be config and DNS first). |
| Months 6–9 | Geo-home shards for region-2-local writers; routing layer honors per-shard home. |
| Months 9–12 | Region 3; first full evacuation game day with a declared RPO/RTO target. |
| Ongoing | Quarterly game days; failover success rate reported to leadership. |
Metrics to Watch: replication.lag_seconds{region}, failover.drill_success_rate, failover.drill_duration_minutes, rpo.declared_vs_measured.
Organizational Follow-up: Every service owner declares consistency tier and RPO; exceptions tracked by the platform; drills are calendar events, not heroics.
Ownership Question: "Who owns 'region evacuation works'?" Staff answer: The platform owns the database portion. The org needs a single accountable owner — usually an SRE lead — for the end-to-end evacuation, because the database failing over is useless if config and identity don't.
Key Takeaway: "A second region is not DR until you have failed over to it on purpose. Budget the drills, not just the replicas."
What clears the Staff bar:
- Asks which of DR, latency, residency is the goal
- Phases the rollout with drills before commitments
- Declared vs measured RPO per service
9. Level Expectations Summary#
After studying this case study, you should be able to:
- State a consistency guarantee per operation and justify each with a latency number and a named owner
- Explain precisely why R+W>N is overlap, not linearizability, with a counterexample
- Walk a leader failover second by second, including terms, leases, fencing, and the no-op commit
- Quote an RPO for node, AZ, and region loss and name who approved each
- Implement read-your-writes across devices with commit tokens and bound its cost
- Pick a conflict-resolution strategy per data type and say what users will observe
- Explain why anti-entropy must finish inside the tombstone GC window
- Describe how the platform, not each team, should own these decisions at org scale
The Bar for This Question#
Mid-level (L4): Draws leader + followers, knows async vs sync replication, reads from replicas for scale. Can describe failover as "promote a replica." Rarely mentions anomalies unless prompted.
Senior (L5): Chooses a reasonable topology, knows quorum arithmetic, mentions replication lag and eventual consistency, handles failover with automation. The design works on the happy path and for single-node failures. Gaps: consistency is global rather than per operation, quorum overlap is confused with linearizability, split brain and RPO are hand-waved, deletes and silent divergence are not addressed.
Staff+ (L6): Starts from the anomaly contract. Linearizable where invariants live, session guarantees where users would notice, eventual elsewhere — each with a number. Failover is fenced, majority-elected, RPO-accounted and rehearsed. Conflicts are resolved per data type with product sign-off. Anti-entropy and tombstones are explicit. Ownership is split cleanly: platform owns mechanics, product owns chosen anomalies. The interviewer should learn something from the answer.
10. Staff Insiders: Controversial Opinions#
10.1 "Eventual Consistency" Is Not a Design — It's an Absence of One#
| Evidence | Detail |
|---|---|
| No bound | "Eventually" has no upper limit; lag can be minutes during backfills |
| No per-user guarantee | Users seeing their own writes vanish is legal under eventual consistency |
| No merge semantics | Says nothing about which concurrent write survives |
The Staff position: Replace "eventual" with a named, measured guarantee: session consistency with p99 lag < 1s, or bounded staleness ≤ 10s. If you can't state the bound, you can't alert on it.
Why this matters in interviews: Saying "eventual consistency is fine here" without a bound is the single most common L5 tell in this question.
10.2 Multi-Leader Replication Is Usually a Latency Optimization Wearing a Resilience Costume#
| Evidence | Detail |
|---|---|
| Geo-homing solves most latency | Most users write mostly their own data; home it near them |
| Conflicts are forever | Every concurrent edit becomes a product-visible merge |
| Uniqueness is impossible | Usernames, seats, balances still need a single serialization point |
The Staff position: Default to single-leader with geo-partitioned homes. Use multi-writer only for data types that merge mathematically.
Why this matters in interviews: Proposing active-active for everything reads as not having debugged a conflict in production.
10.3 Automatic Cross-Region Failover Does More Harm Than Good for Most Companies#
| Evidence | Detail |
|---|---|
| Rare event | Real regional losses happen maybe once every few years per provider-region |
| Partitions look like failures | Short partitions trigger promotions that fork history |
| Reconciliation cost | Forked history takes hours to days to reconcile manually |
The Staff position: Automate in-region (consensus makes it safe). Gate cross-region on a human with a 5-minute decision SLO and a pre-computed RPO estimate — unless you rehearse automatic cross-region failover monthly.
Why this matters in interviews: Saying "we'd make it automatic" without fencing and RPO guards is a red flag; saying "human-gated, here's why and here's the decision SLO" is a Staff-level signal.
10.4 Your Replication Is Only As Good As Your Last Verified Failover#
| Evidence | Detail |
|---|---|
| Silent divergence exists | Non-deterministic migrations, disk bit-rot, apply bugs |
| Unexercised paths rot | Config, DNS, secrets drift in the standby region |
| Row counts lie | Only value-level checksums detect divergence |
The Staff position: Continuous Merkle verification plus scheduled failovers. Treat "last successful failover > 90 days ago" as an incident-worthy risk.
Why this matters in interviews: It shows you think about the system after it's built.
10.5 Most Teams Should Never See a Quorum Knob#
| Evidence | Detail |
|---|---|
| Knobs get turned in incidents | W=1 "temporarily" becomes permanent |
| Few teams can reason about sloppy quorums | Misconfigurations are silent |
| Tiers are auditable | "Session tier" can be tested; "R=1,W=2" rarely is |
The Staff position: The platform exposes named tiers with guarantees; knobs are internal.
Why this matters in interviews: Moving from "which setting" to "which contract" is the L6 → L7 bridge.
11. The Principal Lens (L7)#
Why L7 Sees This Problem Differently#
At Staff level, the replicated data store is one system to get right. At Principal level, it is the substrate under 50–200 services, each of which will independently choose a database, a replication mode, and a set of anomalies — usually without writing any of it down. The Principal problem is not "design the replication"; it is "make sure the company's aggregate anomaly exposure is chosen, priced, and tested, and that the number of distinct replication stacks the org must operate stays small enough to operate well."
The Org-Level Fault Line#
One storage platform with a consistency menu vs every team choosing its own database.
| Choice | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Central platform, tiered menu | Consistent guarantees, one failover practice, shared on-call expertise | Platform becomes a bottleneck; edge-case teams feel constrained | Platform team headcount; teams with exotic needs |
| Team choice | Fit-for-purpose per team; velocity early | 8 databases × 3 replication modes; nobody rehearses failover; hiring for each | Every on-call rotation; the company during regional disasters |
| Paved road + exceptions | Defaults for 85–90% of teams; review for the rest | Exception creep if review is rubber-stamped | Architecture review board time |
🧭 Principal Move: "I'd support three stores on the paved road — a consensus-replicated KV/SQL store, a managed wide-column store for write-heavy AP workloads, and object storage — each with named tiers. Anything else requires an exception with its own on-call and its own failover drill evidence."
Cost Model#
Assumptions: cloud list-ish prices, NVMe-backed instances at ~$0.10–0.15/GB-month effective for replicated storage, cross-region transfer ~$0.02/GB, fully loaded engineer $300K/year ($25K/month). Figures are order-of-magnitude.
| Scale | Data / Throughput | Topology | Infra $/month | Headcount | On-call Load |
|---|---|---|---|---|---|
| Small | 1 TB, 10K QPS | 3 AZs, 1 region, managed service | $3K–8K | 0.5 FTE (shared) | ~1 page/month |
| Medium | 50 TB, 200K QPS | Raft in-region + 2 async regions (≈7 copies) | $80K–150K (incl. ~$15K transfer) | 3–5 FTE | Dedicated rotation, ~4 pages/month |
| Large | 1 PB, 5M QPS | Geo-homed shards, 3–5 regions, tiered RF | $1.5M–3M | 15–25 FTE platform team | 24/7 follow-the-sun, game days quarterly |
The step from small to medium is dominated by copies (3 → 7–9) and transfer; the step to large is dominated by people — the platform team is often a bigger line item than the storage premium over a managed service.
The 3-Year Evolution Path#
One-Way Doors vs Two-Way Doors#
| Decision | Door | Reversal Cost |
|---|---|---|
| Replication topology of an existing store (single-leader → multi-leader) | One-way | Data model and app merge logic change; months |
| Conflict-resolution semantics exposed to users | One-way | Users and downstream systems depend on "add wins" |
| Vendor-managed store with proprietary API | Mostly one-way | Egress at $0.02–0.09/GB plus rewrite; PB scale = 7-figure migration |
| Key/partitioning scheme | One-way at scale | Full re-shard of data |
| Read consistency tier for an endpoint | Two-way | Flag flip + capacity check |
| Replica count per region | Two-way | Hours to rebuild or remove |
| Lag alert thresholds, proxy timeouts | Two-way | Config change |
The Standard I'd Write#
RFC: Data Replication and Consistency Standard (v1)
Scope: All persistent stores holding customer or financial data.
MUST:
- Every store MUST declare, per API operation, a consistency tier:
strict(linearizable),session,bounded(Δ), oreventual.- Every store MUST publish RPO and RTO for node, AZ, and region loss, approved by the product owner.
- Leader failover MUST use epoch/term fencing enforced by storage, not by a health checker.
- Wall-clock last-write-wins MUST NOT be used for financial data or user-generated collections.
- AP stores MUST complete full anti-entropy repair within 70% of the tombstone GC window.
- Replication lag and unreplicated-write metrics MUST be emitted in the standard names.
SHOULD: use the paved-road stores; run a failover drill at least quarterly; verify replica checksums before planned failovers.
Exceptions: Filed with the storage architecture group; require an owning on-call rotation and drill evidence; expire after 12 months.
Success metrics: % of services with declared tiers (target 100% in 2 quarters); failover drill success rate ≥ 95%; declared-vs-measured RPO breaches = 0; number of distinct replication stacks in production (target ≤ 3).
What I'd Tell the VP#
Our data is copied across multiple machines and locations, but today nobody can tell you how much we'd lose in a regional outage or how long we'd be down — the answer differs per team and has never been tested end to end. I want every service to declare that promise in plain numbers, and I want us to prove it by failing over on purpose each quarter. Most of our data doesn't need expensive global guarantees; we'll reserve those for money and identity and save on the rest. This needs a platform team of about four people over the next year and roughly a 15% increase in storage spend for the second region. The payoff is that a regional outage becomes a 30-minute planned procedure rather than a multi-day incident.
Principal Interview Signals#
| Signal | What It Sounds Like |
|---|---|
| Menu, not knobs | "Teams pick a tier with a published price and SLO; they don't set R and W." |
| Priced tradeoffs | "The second region is ~$15K/month in transfer alone at 50 TB churn — here's what it buys in RPO." |
| Rehearsal as a control | "Failover success rate is an org SLO. A drill we skip is a risk we accepted." |
| Upstream governance | "Sales contracts reference tiers, not raw RPO numbers." |
| Knows when not to standardize | "The ML feature store stays outside the platform; it's derivable and rebuilt nightly." |
Staff answers that L7 interviewers find insufficient:
- A flawless single-store design that ignores the other 40 stores in the org with different guarantees
- "We'll make failover automatic" without an org-wide drill program or success metric
- Naming who pays without pricing it — "product accepts staleness" with no dollar or SLO attached
Appendices
Appendix A: Consistency Mechanics in Depth#
A.1 Leader Leases for Linearizable Reads#
on_become_leader(term):
lease_expiry = now_monotonic() + LEASE - MAX_DRIFT # e.g. 9s - 0.5s
append(no_op, term) # commit prior-term entries
on_heartbeat_majority_ack(sent_at):
lease_expiry = sent_at + LEASE - MAX_DRIFT
read_linearizable(key):
if now_monotonic() < lease_expiry and commit_index_applied():
return local_read(key)
else:
idx = read_index_round() # one heartbeat round to a majority
wait_until_applied(idx)
return local_read(key)
Why it's right: a new leader cannot be elected until followers' election timeouts expire, which is longer than the lease. Why it's wrong if misconfigured: if MAX_DRIFT understates real clock drift, a paused-then-resumed leader can serve a stale read inside what it thinks is a valid lease.
A.2 Quorum Arithmetic#
| N | W | R | Guarantee | Tolerates for Writes | Tolerates for Reads |
|---|---|---|---|---|---|
| 3 | 2 | 2 | Overlap for completed writes | 1 down | 1 down |
| 3 | 3 | 1 | Overlap; fast reads | 0 down | 2 down |
| 3 | 1 | 3 | Overlap; fast writes | 2 down | 0 down |
| 3 | 1 | 1 | None | 2 down | 2 down |
| 5 | 3 | 3 | Overlap | 2 down | 2 down |
Overlap holds only with strict quorums on the designated preference list; sloppy quorums break it.
A.3 Hybrid Logical Clocks#
send_or_local_event():
pt = physical_now()
l_new = max(l, pt)
c = (c + 1) if l_new == l else 0
l = l_new
return (l, c)
receive(msg_l, msg_c):
pt = physical_now()
l_new = max(l, msg_l, pt)
if l_new == l == msg_l: c = max(c, msg_c) + 1
elif l_new == l: c = c + 1
elif l_new == msg_l: c = msg_c + 1
else: c = 0
l = l_new
HLCs preserve causality (if A happened-before B, ts(A) < ts(B)) while staying within clock skew of physical time — which makes LWW deterministic and causally safe, but still lossy for truly concurrent writes.
Appendix B: Data Model and Keys#
B.1 Record Layout#
| Field | Type | Purpose |
|---|---|---|
key | bytes | Range-partitioned; prefix with tenant for isolation |
value | bytes | Opaque to the store |
version | (term, index) or HLC | Conditional writes, LWW ordering |
vv | map node→counter | Only in multi-writer tier; detects concurrency |
tombstone | bool + deleted_at | Deletes survive replication until GC |
home_region | enum | Geo-homing; routes writes |
B.2 Commit Token Format#
base64(shard_id:u32 | term:u32 | log_index:u64) — 16 bytes, opaque to clients, comparable per shard. Session store holds max token per (user, shard) with 60s TTL.
Appendix C: Conflict Resolution — Quick Comparison#
| Mechanism | Detects Concurrency | Loses Data | Metadata Cost | App Work | Best For |
|---|---|---|---|---|---|
| Wall-clock LWW | No | Yes, silently | 8 bytes | None | Nothing important |
| HLC LWW | No (orders deterministically) | Yes, concurrent only | 12 bytes | None | Overwrite fields |
| Version vectors + siblings | Yes | No | O(writers) | Merge function | Documents with custom merge |
| OR-set CRDT | Yes (by construction) | No | Tombstones per removed element | None | Carts, tags, memberships |
| PN-counter | Yes | No | O(replicas) counters | None | Likes, inventory display |
| Single-writer routing | N/A | No | None | Routing | Balances, uniqueness |
C.1 Anti-Entropy With Merkle Trees#
build_tree(range, depth=15): # 32,768 leaves
for each key in range:
leaf[hash(key) % 2^depth] ^= hash(key, version, value)
internal nodes = hash(children)
repair(replica_a, replica_b, range):
if a.root == b.root: return
descend only into differing children
for each differing leaf: stream keys in that leaf bucket, reconcile by version
Throttle to 10–50 MB/s per node; stagger per range so each range is fully repaired every 7 days on a 10-day GC window.
Appendix D: API Contract and Client Behavior#
| Header / Field | Direction | Meaning |
|---|---|---|
If-Match: <version> | Request | Conditional write; 412 on mismatch |
X-Consistency: strict / session / bounded=10s / eventual | Request | Tier for this read |
X-Commit-Token | Both | Session high-water mark |
X-Served-By-Region | Response | Debugging stale reads |
X-Staleness-Ms | Response | Measured lag of the replica that served |
Idempotency-Key | Request | Safe retry across failover |
Clients retry writes on leader-change errors with jittered backoff (base 50ms, cap 2s) and the same idempotency key. Without jitter, a leader election on a hot shard produces a synchronized retry storm at the new leader.
Appendix E: Observability#
E.1 Core Metrics#
replication.lag_seconds{shard,region} # applied-time lag
replication.lag_bytes{shard,region} # log backlog
rpo.unreplicated_writes{shard,region} # writes not yet in remote region
raft.leader_elections_total{shard}
raft.leader_count_per_shard # must be 1
session.leader_proxy_rate
session.token_wait_ms
conflict.lww_discards{table}
antientropy.divergent_ranges
antientropy.last_full_repair_age_days
consistency.sample_mismatch_rate
E.2 Critical Alerts#
| Alert | Threshold | Action |
|---|---|---|
| Lag warning | lag_seconds > 5 for 2 min | Ticket; check bulk writers |
| Lag page | lag_seconds > 30 for 2 min | Page; RPO at risk |
| Split brain | leader_count_per_shard > 1 | Page immediately; freeze writes to shard |
| Election storm | > 5 elections/hour/shard | Page; likely GC or network |
| Repair overdue | repair age > 70% of GC window | Page platform |
| Silent divergence | divergent_ranges > 0 | Incident, even with no user impact |
E.3 Control Plane vs Data Plane#
The data plane (leaders, followers, router) must keep serving if the control plane (etcd shard map, failover controller) is unavailable — routers cache the map and followers keep replicating. The control plane being down should block rebalancing and cross-region promotion, never reads and writes.
Appendix F: Scale Evolution#
F.1 What Works at Each Scale#
| Scale | Design |
|---|---|
| < 1 TB, 1 region | Managed PostgreSQL/Aurora with sync standby; read replicas |
| 1–50 TB, 1–2 regions | Consensus-replicated sharded store or DynamoDB/Spanner; async remote reads |
| 50 TB – 1 PB, 3+ regions | Geo-homed shards, tiered RF, dedicated platform team |
| > 1 PB, global | Multiple stores by workload; CRDT tier for partition-tolerant writes |
F.2 What You Don't Build on Day One#
- Multi-leader replication
- Automatic cross-region failover
- Custom consensus implementation (use etcd/Raft libraries or a managed store)
- Per-key consistency configuration (per-operation tiers are enough)
- Global linearizable writes
Appendix G: Multi-Tenancy and Cost Fairness#
| Concern | Mechanism |
|---|---|
| Noisy writer causes lag for all | Per-tenant write admission; bulk priority class |
| Tenant needs stricter tier | Dedicated shard groups, priced separately |
| Linearizable read cost | Chargeback per leader-read; default tier is session |
| Residency | Tenant → home region pinning; replication allowlist per tenant |
| Deletion SLA | Per-tenant tombstone tracking; repair-age SLO covers their ranges |