Hiring BarSupport

Design a Replicated Data Store — Staff-Level Case Study

Case study73 min read8 diagrams

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.

ModeTimeWhat to Read
Quick Review15 minExecutive Summary → Interview Walkthrough → Fault Lines table → Drills 1, 3, 5
Targeted Study1–2 hrsExecutive Summary → Walkthrough → Section 3 (Fault Lines) → Section 4 (Failures) → two Deep Dives
Deep Dive3+ hrsEverything, 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
TopologyHow It WorksProsCons
Single-leader, async followersLeader commits locally, ships log to followers in backgroundFast writes (1 local fsync); simple mental modelFailover can lose the last N ms of acked writes (RPO > 0); stale follower reads
Single-leader, sync / semi-syncLeader waits for ≥1 follower ack before acking clientRPO = 0 for 1 failure+1 RTT per write (1–2ms cross-AZ, 60–100ms cross-region)
Consensus-replicated (Raft/Paxos) per shardLeader needs majority ack; automatic leader election with termsLinearizable, automatic safe failover, no split brainMajority RTT on every write; leader is a per-shard bottleneck
Multi-leaderEach region accepts writes; async cross-replicationLocal-latency writes everywhere; region-independentWrite–write conflicts are guaranteed; needs a resolution policy
Leaderless (Dynamo-style)Client/coordinator writes to W of N replicas, reads from RHigh write availability; no failover stepConflicts, read repair, sloppy quorums; R+W>N is not linearizable
Chain replicationWrites go head → … → tail; reads from tailStrong consistency, high read throughput at tailLatency = 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#

BehaviorSenior (L5)Staff (L6)Principal (L7)
First movePicks a topology ("leader + 2 followers")Asks which operations need linearizability and which tolerate staleness; commits per operationAsks which classes of data the org stores and publishes a consistency menu the platform supports
Consistency"Strong consistency" or "eventual consistency" for the whole systemMixes: linearizable for money/identity writes, read-your-writes for user sessions, eventual for feedsPrices 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 gapsStandardizes 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-regionGame-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 cartsBans 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 incidentsRedraws 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#

PositionRationale
Consistency is chosen per operation, not per systemOne 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-writerConflicts 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 regionsMajority 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 sessionsTokens survive load-balancer reshuffles and region routing; stickiness breaks on the first deploy
Never last-write-wins on wall clocks for data that mattersNTP skew of 10–100ms silently discards acknowledged writes; use versions, HLCs, or CRDTs
Failover is fenced, majority-elected, and rehearsedThe failover path is the least-exercised, highest-blast-radius code you own
Anti-entropy is not optional in any AP designRead 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.

IntentConstraintStrategyFailure ModeCorrectness 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-holderUnavailable (write-blocked) for the minority side of a partitionRPO = 0, linearizable, audited
User-facing read scale (profiles, catalogs, settings)p99 read < 10ms from any region; users must see their own editsSingle 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 partitionMulti-leader or leaderless, CRDTs / version vectors, anti-entropyConflicts on every concurrent edit; merge semantics visible to usersConvergence 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 LineThe Tension
1Leader-based vs LeaderlessOne writer per key (no conflicts, failover risk) or any writer (no failover, permanent conflicts)?
2Synchronous vs Asynchronous ReplicationPay a round-trip on every write for RPO = 0, or ack fast and accept losing the tail on failover?
3Read Consistency LevelLinearizable (leader, slow, available only on majority side) vs session vs eventual (fast, stale)?
4Conflict ResolutionLWW (simple, lossy) vs version vectors + app merge (correct, product burden) vs CRDTs (automatic, limited types)?
5Failover AuthorityAutomatic (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#

Diagram: 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#

TopicThe L5 AnswerThe 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#

MetricValueWhy It Matters
Intra-AZ round trip~0.1–0.5msFloor for any replicated write
Cross-AZ round trip (same region)~1–2msCost of an in-region majority commit
Cross-region RTT: US east ↔ west~60–70msCost of one synchronous cross-region ack
Cross-region RTT: US ↔ EU~80–100msWhy global sync writes feel slow
Cross-region RTT: US ↔ APAC~150–200msWhy global consensus is reserved for tiny, critical data
Raft in-region commit p99~2–5msIncludes fsync on majority
Async replica lag, healthyp50 10–100ms, p99 < 1sThe size of your stale-read window and your RPO
Async replica lag, during backfill or vacuum10s – minutesWhy RPO must be measured, not assumed
NTP clock skew between hosts1–10ms typical, 100ms+ misconfiguredWhy wall-clock LWW loses writes
Leader lease duration5–10sBounds unavailability after leader loss
Failure detection timeout3–10s in-regionToo low = flapping elections; too high = long write outage
Cassandra default gc_grace_seconds864,000 (10 days)Repair must complete inside it or deletes resurrect
Fault tolerance of a majority quorumN=3 survives 1 loss; N=5 survives 2N=5 spread 2/2/1 across 3 AZs survives a full AZ loss, or any 2 nodes
Merkle tree repair of 1 TB replicahours, throttled to ~10–50 MB/sSchedule 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), tombstone flag
  • 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_ms so callers can see the lie they were told."


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

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

Walk the flow in 90 seconds:

  1. Key hashes to a range; the router caches the range → leader map (epoch-versioned, from etcd).
  2. Writes go to the leader, append to the Raft log, commit on 2 of 3 in-region acks (2–4ms), return a commit token.
  3. The log ships asynchronously to remote-region followers.
  4. Reads: linearizable → leader with valid lease; session → any replica whose applied index ≥ token; eventual → nearest replica.
  5. 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#

MistakeL5 Does ThisL6 Does This Instead
Storage engine rabbit hole10 min on LSM compaction vs B-tree"LSM, because write-heavy. Engine choice doesn't change the replication contract."
Partitioning over replicationLong 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-entropyStops at read repairMerkle 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:

  1. Consistency named per operation, with a reason
  2. RPO/RTO stated in numbers with a named approver
  3. Failover walked through with fencing
  4. Quorum arithmetic stated precisely (overlap, not linearizability)
  5. Anti-entropy and deletes handled explicitly

1.2 The L5 vs L6 Contrast — Visual#

Diagram: 1.2 The L5 vs L6 Contrast — Visual

1.3 The Staff Question That Cuts Through Everything#

"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 DetailWhy It MattersWhat to Ask
Durability promiseDecides sync vs async and RPO"Can we ever lose an acknowledged write? For which failure — node, AZ, region?"
Read freshness per use caseDecides where reads go"Which reads must see the latest write — everyone's or just the writer's?"
Partition behaviorDecides CP vs AP per operation"During a region partition, should the cut-off region reject writes or accept and merge later?"
Write geographyDecides home-region vs multi-leader"Where are writers? One region, or all of them with equal latency needs?"
Delete semanticsDecides tombstones, GC, compliance"Are deletes legally required to be permanent within N days?"
Data typesDecides whether CRDTs are viable"Is this counters, sets, documents, or balances?"

Staff engineers surface these. Senior engineers assume them away.

2.4 Precise Terminology#

TermPrecise MeaningCommon Misuse
LinearizableEach operation appears to take effect atomically at one instant between its call and return; all clients agree on orderUsed as a synonym for "durable"
SerializableTransactions appear to execute in some serial order — not necessarily real-time orderConfused with linearizable; they're orthogonal
Read-your-writesA session always sees its own prior writesAssumed to follow from "eventual consistency"
Monotonic readsA session never sees older data after seeing newerBroken by load-balancing across replicas with different lag
Bounded stalenessReads are at most Δ seconds or K versions behindPromised without measuring lag
RPOMaximum window of acknowledged data that can be lostStated as "zero" for async replication
RTOTime until service restored after failureMeasured from detection, hiding detection time
QuorumA subset whose size guarantees intersection with another subsetTreated as "strong consistency"
Sloppy quorumWrites accepted by any N healthy nodes, not the designated NAssumed to preserve R+W>N overlap (it doesn't)
Fencing token / epochMonotonic number that storage uses to reject stale leadersOmitted from failover design
TombstoneMarker recording a delete so replicas don't resurrect dataPurged 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.

ChoiceWhat WorksWhat BreaksWho Pays
Single leader per shard (consensus)No write conflicts; linearizable available; simple reasoning for app teamsWrite 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 partitionEvery concurrent write to the same key conflicts; uniqueness impossible without a coordinatorProduct teams (merge semantics), users (surprising merges)
Leaderless (Dynamo quorum)No failover; writes succeed if any W nodes reachableRead repair, hinted handoff, tombstone resurrection, LWW data lossApp teams (siblings), SRE (repair scheduling)
Diagram: 3.1 Fault Line 1: Leader-Based vs Leaderless

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.

ChoiceWrite LatencyRPOWho Pays
Async everything1 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 lossLatency budget; stalls if the one sync follower is slow
Majority in-region (Raft)+1–3ms0 for AZ lossPlatform (consensus ops)
Majority cross-region+60–200ms0 for region lossEvery user on every write

The Staff position is layered: majority in-region, async cross-region, and an explicit statement of the region-loss RPO.

Diagram: 3.2 Fault Line 2: Synchronous vs Asynchronous Replication

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.

LevelServed FromLatencyStalenessAvailable During Partition?Use For
LinearizableLeader with valid lease (or ReadIndex round)1–3ms in-region, +RTT cross-region0Majority side onlyUniqueness checks, balance, locks, inventory decrement
Session (RYW + monotonic)Any replica with applied index ≥ tokenLocal in ~97–99% of cases0 for own writes; ≤ lag for othersYes, with fallback to leaderUser-facing reads after edits
Bounded stalenessAny replica with lag ≤ ΔLocal≤ Δ (e.g., 10s)Yes, until lag > ΔAnalytics, dashboards, search indexing
EventualNearest replicaLocal, < 1msUnbounded in theoryYesCounters, feeds, recommendations
Diagram: 3.3 Fault Line 3: Read Consistency Level

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.

StrategyHow It WorksWhat WorksWhat BreaksWho Pays
LWW, wall clockHighest timestamp winsTrivial, no siblingsSkew of 10–100ms drops acked writes silentlyUsers (vanishing edits), undetectable
LWW, HLC / logicalHybrid logical clock, tiebreak by node IDCausality respected; bounded skewStill drops concurrent writes — just deterministicallyUsers, but explainable
Version vectors + siblingsDetect concurrency, return all versionsNothing lostApp must merge; sibling explosion if clients don'tApp teams
CRDTsTypes whose merge is commutative, associative, idempotentAutomatic convergenceLimited types; metadata overhead (tombstones in OR-sets can grow 2–10×)Platform (implementation), product (semantics like "add wins")
Single-writer routingRoute all writes for a key to its home regionNo conflictsCross-region write latency for non-home usersRemote 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_discards counter."


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.

ChoiceRTORiskWho Pays
Automatic, watchdog-based (external health checker promotes)10–60sSplit brain under partition; promotes lagging replicaCustomers (data divergence), on-call (reconciliation for days)
Automatic, consensus-based (Raft election)3–10sSafe for in-region; can't help if majority is gonePlatform (consensus correctness)
Human-gated cross-region5–30 minSlower; decision fatigueCustomers (downtime), incident commander
Automatic cross-region with RPO guard1–5 minOnly promotes if measured lag < thresholdPlatform (guard logic)
Diagram: 3.5 Fault Line 5: Failover Authority

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#

FailureDetection SignalBlast RadiusMitigationOwner
Leader crashraft.leader_elections_total spike1 shard, 3–10s writesAutomatic electionStorage platform
Paused leader / split brainraft.leader_count_per_shard > 11 shardTerm fencing, leasesStorage platform
AZ lossaz.replica_unreachableUp to 1/3 of replicasMajority survives; rebuild in other AZStorage platform + infra
Region lossregion.health, lag frozenAll home-region writesHuman-gated promote; RPO = lagIncident commander + platform
Replication lag spiralreplication.lag_seconds > 30Session reads, leader loadThrottle bulk writers, cap proxyPlatform + offending team
Resurrected deletesantientropy.last_full_repair_age_daysDeleted data returnsWipe stale nodesPlatform + privacy
LWW discardsconflict.lww_discardsConcurrent writers' dataMove data type to CRDT/single-writerProduct owning the data
Silent divergenceconsistency.sample_mismatch_rateUnknown until failoverRepair from majorityPlatform
Hot shard leadershard.leader_cpu > 80%1 shard's keysSplit range, move leadershipPlatform

5. Evaluation Rubric#

5.1 Level-Based Signals#

DimensionSenior (L5)Staff (L6)Principal (L7)
Consistency modelOne setting for the systemPer-operation tiers with latency and staleness numbersOrg-wide consistency menu with published price, p99 and RPO per tier
TopologyLeader-followerConsensus per shard in-region, async cross-region, justified by RTTsDecides which topologies the platform supports at all; retires the rest
QuorumsR+W>N = consistentOverlap ≠ linearizable; sloppy quorum and LWW gapsRemoves tunables from app teams; tiers instead of knobs
FailoverPromote a replicaTerms, fencing, leases, RPO accounting, human gate cross-regionFailover success rate as an SLO; quarterly region evacuations
ConflictsLWWPer data type: HLC-LWW, CRDT, app merge, single-writerPlatform standard bans wall-clock LWW for money; exception process
Anti-entropyRead repairMerkle repair inside tombstone GC windowDivergence detection as a first-class compliance control
OwnershipDB team owns itPlatform owns mechanics; product owns chosen anomaliesConsistency level reviewed as part of API contract; cost charged back

5.2 Strong Hire Signals#

SignalWhat 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#

SignalWhy It Misses the Bar
"Strong consistency everywhere"Ignores the cross-region latency bill; no per-operation reasoning
"Eventual consistency is fine" without defining for whomHands users anomalies nobody approved
Automatic failover with no fencingDesigns a split-brain generator
Multi-leader for usernamesDoesn't understand that uniqueness requires a single serialization point
No numbers for lag or RPOCan'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#

PhaseTimeGoal
Framing0–3 minDurability promise, freshness per use case, partition behavior; RPO/RTO numbers
Entities + API3–5 minRecord, shard, commit token; consistency as a read parameter
High-level design5–12 minConsensus in-region, async cross-region, router
Transition12 minOffer failover, read consistency, multi-writer
Deep dives12–40 minFailover + fencing → session guarantees → conflicts/anti-entropy
Wrap-up40–45 minAnomaly contract, ownership, rehearsal, what's next

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

Interviewer SignalWhat They Care AboutWhere to Go Deep
"What if the leader dies?"Operational maturityTerms, fencing, RPO accounting
"Can we read from replicas?"Consistency literacySession tokens, bounded staleness
"Make writes work in every region"Conflict reasoningCRDTs vs version vectors, uniqueness caveat
"Is R+W>N linearizable?"PrecisionOverlap vs linearizability, sloppy quorum
"How do you know replicas agree?"Silent failure awarenessMerkle anti-entropy, sampled compare
"What about GDPR deletes?"Compliance and tombstonesGC window, stale-node rebuild

6.3 What to Deliberately Skip#

TopicWhy L5 Goes HereWhat L6 Says Instead
LSM vs B-tree internalsFeels deep"LSM for write-heavy; doesn't change the replication contract."
Paxos vs Raft proofsFeels rigorous"Raft for understandability; both give the same guarantees here."
Consistent hashing ring mathFamiliar"Range partitioning with splits; see sharding. Replication is the question."
Client SDK detailsEasy to list"Client carries the commit token; that's the only SDK requirement."

6.4 Follow-Up Questions to Expect#

  1. "Is R=2, W=2, N=3 linearizable? Prove or disprove."
  2. "The leader was paused for 15 seconds. What does it do when it wakes?"
  3. "How do you give read-your-writes to a user who switches from phone to laptop?"
  4. "Two regions accept a write to the same key 5ms apart. Which wins, and does anyone find out?"
  5. "How do you detect that two replicas silently disagree?"
  6. "How long can a replica be offline before it can't safely rejoin?"
  7. "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_violations sampled 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
PhaseWhat 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.
TriageSplit lag into ship lag vs apply lag. Apply lag with idle network = CPU/disk on followers. Identify top write sources by tenant.
Quick fixPause 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.
GuardrailsLeader CPU < 75%; write p99 < 10ms. If breached, prefer returning HTTP 409 "still saving, retry" over stale reads.
Post-mortemThe 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 = 0 a precondition for any planned failover org-wide.

Staff Approach — Full Reasoning
PhaseWhat to Do
Immediate (0–5 min)Freeze further failovers. Fail back to the old leader if it is intact and fenced correctly (bump epoch).
TriageMerkle diff old vs promoted per range; list divergent keys and their last-modified times; correlate with deploy/migration history.
Quick fixRepair promoted replica from old leader's values for divergent keys; customer comms for affected accounts.
GuardrailsPre-failover check: range hashes match within the last 24h; otherwise failover requires director approval.
Post-mortemRoot 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
PhaseWhat to Do
Week 1Confirm contract language with legal; map RPO/SLA promises to an architecture; get written acceptance of the 25–35ms write latency.
Weeks 2–3Provision 5-replica groups across 3 EU regions; dedicated shard range; leader placement near the customer's primary writers.
Weeks 4–5Load test at 2× forecast; chaos test region isolation; verify writes continue with RPO=0 when one region is cut.
Week 6Shadow traffic, then cutover with rollback plan.
OngoingPer-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_discards and 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
PhaseWhat to Do
ImmediateRoute 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.
TriageInstrument discards; confirm that mobile and web writing from different regions within 30ms account for the losses.
FixOR-set CRDT for items; PN-counter for quantity; tombstone GC tied to causal stability.
GuardrailsShadow-merge old and new cart representations for 2 weeks; alert on divergence.
Post-mortemRoot 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
PhaseWhat to Do
Months 0–3Region 2 async followers; commit tokens; lag and rpo.unreplicated_writes dashboards.
Months 3–6Failover drills on low-risk shards monthly; fix what breaks (it will be config and DNS first).
Months 6–9Geo-home shards for region-2-local writers; routing layer honors per-shard home.
Months 9–12Region 3; first full evacuation game day with a declared RPO/RTO target.
OngoingQuarterly 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#

EvidenceDetail
No bound"Eventually" has no upper limit; lag can be minutes during backfills
No per-user guaranteeUsers seeing their own writes vanish is legal under eventual consistency
No merge semanticsSays 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#

EvidenceDetail
Geo-homing solves most latencyMost users write mostly their own data; home it near them
Conflicts are foreverEvery concurrent edit becomes a product-visible merge
Uniqueness is impossibleUsernames, 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#

EvidenceDetail
Rare eventReal regional losses happen maybe once every few years per provider-region
Partitions look like failuresShort partitions trigger promotions that fork history
Reconciliation costForked 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#

EvidenceDetail
Silent divergence existsNon-deterministic migrations, disk bit-rot, apply bugs
Unexercised paths rotConfig, DNS, secrets drift in the standby region
Row counts lieOnly 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#

EvidenceDetail
Knobs get turned in incidentsW=1 "temporarily" becomes permanent
Few teams can reason about sloppy quorumsMisconfigurations 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.

ChoiceWhat WorksWhat BreaksWho Pays
Central platform, tiered menuConsistent guarantees, one failover practice, shared on-call expertisePlatform becomes a bottleneck; edge-case teams feel constrainedPlatform team headcount; teams with exotic needs
Team choiceFit-for-purpose per team; velocity early8 databases × 3 replication modes; nobody rehearses failover; hiring for eachEvery on-call rotation; the company during regional disasters
Paved road + exceptionsDefaults for 85–90% of teams; review for the restException creep if review is rubber-stampedArchitecture 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.

ScaleData / ThroughputTopologyInfra $/monthHeadcountOn-call Load
Small1 TB, 10K QPS3 AZs, 1 region, managed service$3K–8K0.5 FTE (shared)~1 page/month
Medium50 TB, 200K QPSRaft in-region + 2 async regions (≈7 copies)$80K–150K (incl. ~$15K transfer)3–5 FTEDedicated rotation, ~4 pages/month
Large1 PB, 5M QPSGeo-homed shards, 3–5 regions, tiered RF$1.5M–3M15–25 FTE platform team24/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#

Diagram: The 3-Year Evolution Path

One-Way Doors vs Two-Way Doors#

DecisionDoorReversal Cost
Replication topology of an existing store (single-leader → multi-leader)One-wayData model and app merge logic change; months
Conflict-resolution semantics exposed to usersOne-wayUsers and downstream systems depend on "add wins"
Vendor-managed store with proprietary APIMostly one-wayEgress at $0.02–0.09/GB plus rewrite; PB scale = 7-figure migration
Key/partitioning schemeOne-way at scaleFull re-shard of data
Read consistency tier for an endpointTwo-wayFlag flip + capacity check
Replica count per regionTwo-wayHours to rebuild or remove
Lag alert thresholds, proxy timeoutsTwo-wayConfig change

The Standard I'd Write#

RFC: Data Replication and Consistency Standard (v1)

Scope: All persistent stores holding customer or financial data.

MUST:

  1. Every store MUST declare, per API operation, a consistency tier: strict (linearizable), session, bounded(Δ), or eventual.
  2. Every store MUST publish RPO and RTO for node, AZ, and region loss, approved by the product owner.
  3. Leader failover MUST use epoch/term fencing enforced by storage, not by a health checker.
  4. Wall-clock last-write-wins MUST NOT be used for financial data or user-generated collections.
  5. AP stores MUST complete full anti-entropy repair within 70% of the tombstone GC window.
  6. 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#

SignalWhat 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#

NWRGuaranteeTolerates for WritesTolerates for Reads
322Overlap for completed writes1 down1 down
331Overlap; fast reads0 down2 down
313Overlap; fast writes2 down0 down
311None2 down2 down
533Overlap2 down2 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#

FieldTypePurpose
keybytesRange-partitioned; prefix with tenant for isolation
valuebytesOpaque to the store
version(term, index) or HLCConditional writes, LWW ordering
vvmap node→counterOnly in multi-writer tier; detects concurrency
tombstonebool + deleted_atDeletes survive replication until GC
home_regionenumGeo-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#

MechanismDetects ConcurrencyLoses DataMetadata CostApp WorkBest For
Wall-clock LWWNoYes, silently8 bytesNoneNothing important
HLC LWWNo (orders deterministically)Yes, concurrent only12 bytesNoneOverwrite fields
Version vectors + siblingsYesNoO(writers)Merge functionDocuments with custom merge
OR-set CRDTYes (by construction)NoTombstones per removed elementNoneCarts, tags, memberships
PN-counterYesNoO(replicas) countersNoneLikes, inventory display
Single-writer routingN/ANoNoneRoutingBalances, 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 / FieldDirectionMeaning
If-Match: <version>RequestConditional write; 412 on mismatch
X-Consistency: strict / session / bounded=10s / eventualRequestTier for this read
X-Commit-TokenBothSession high-water mark
X-Served-By-RegionResponseDebugging stale reads
X-Staleness-MsResponseMeasured lag of the replica that served
Idempotency-KeyRequestSafe 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#

AlertThresholdAction
Lag warninglag_seconds > 5 for 2 minTicket; check bulk writers
Lag pagelag_seconds > 30 for 2 minPage; RPO at risk
Split brainleader_count_per_shard > 1Page immediately; freeze writes to shard
Election storm> 5 elections/hour/shardPage; likely GC or network
Repair overduerepair age > 70% of GC windowPage platform
Silent divergencedivergent_ranges > 0Incident, 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#

ScaleDesign
< 1 TB, 1 regionManaged PostgreSQL/Aurora with sync standby; read replicas
1–50 TB, 1–2 regionsConsensus-replicated sharded store or DynamoDB/Spanner; async remote reads
50 TB – 1 PB, 3+ regionsGeo-homed shards, tiered RF, dedicated platform team
> 1 PB, globalMultiple 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#

ConcernMechanism
Noisy writer causes lag for allPer-tenant write admission; bulk priority class
Tenant needs stricter tierDedicated shard groups, priced separately
Linearizable read costChargeback per leader-read; default tier is session
ResidencyTenant → home region pinning; replication allowlist per tenant
Deletion SLAPer-tenant tombstone tracking; repair-age SLO covers their ranges
  1. Loading the index…