Why This Matters#
Replication is not about copies. Every candidate knows to "add a replica." Replication is about which copy is allowed to be wrong, for how long, and who notices. A follower that is 3 seconds behind is a correct system for a product catalog and a broken one for a bank balance. A failover that loses 2 seconds of writes is a non-event for a like counter and a regulatory filing for a ledger. The mechanics — leader-follower, multi-leader, leaderless — are just three different answers to that one question.
Interviewers bring up replication because almost every design reaches it within ten minutes: "the database is a single point of failure," "reads are 95% of traffic," "we need a second region." The Senior answer names a topology. The Staff answer names the lag budget, the data-loss window on failover (RPO), the time to recover (RTO), and the person who signs off when automation promotes a replica that was missing the last few hundred writes.
The third reason is that replication is where outages hide. Most of the famous database incidents are not disk failures. They are a 40-second network blip that triggered a failover nobody could cleanly undo, a lagging replica that served stale reads to a payment flow, or a semi-synchronous setup that quietly fell back to asynchronous when its one sync replica got slow. The copy is never the hard part. Promotion, lag and conflict are.
This page compares the three replication models side by side and says when to choose each. The guarantees they produce (linearizability, session guarantees, CAP/PACELC, isolation) live in Consistency, CAP & PACELC. An end-to-end leaderless store design lives in Replicated Key-Value Store.
The 60-Second Version#
- Three models, one question. Leader-follower: one node accepts writes, others copy. Multi-leader: several nodes accept writes and must reconcile conflicts. Leaderless: clients write to N replicas and need W acks; reads ask R replicas and need
R + W > Nfor overlap. - Sync vs async is a price, not a virtue. A synchronous replica in another AZ adds ~0.5–2ms per commit; in another region it adds 60–150ms. An asynchronous replica adds nothing to the write and loses whatever it had not received when the leader dies.
- RPO for async = lag × write rate. 2s of lag at 5,000 writes/s is ~10,000 acknowledged writes that disappear on failover. Say that number out loud before choosing async.
- Failover is minutes of decisions, not one event. Detection 10–30s, election and promotion 5–30s, client reconnect and DNS/proxy update 5–60s. Typical automated in-region failover: 30s–2min. Cross-region: minutes, and often a human.
- Lag causes user-visible anomalies. Read-your-writes (my comment vanished), monotonic reads (the comment came back, then vanished again), consistent prefix (the answer appears before the question). Each has a cheap targeted fix.
- Split brain is the failure that costs days. Two nodes that both believe they are leader both accept writes. Fencing (epoch numbers checked by storage) is the fix, not better timeouts.
- Logical replication and CDC turn the database log into a product. The same change stream feeds replicas, search indexes, caches and analytics — and becomes a contract other teams depend on.
How Replication Works (for System Designers)#
Key Terms#
| Term | Meaning | Why It Matters in an Interview |
|---|---|---|
| Replication lag | Time (or log bytes) between a write committing on the leader and being visible on a replica | Sets how stale a read can be. Measure it as replica.lag_seconds, alert on p99 not average |
| RPO (recovery point objective) | How much acknowledged data you may lose in a failover | Async: ≈ lag at the moment of failure. Sync: 0 for a single failure |
| RTO (recovery time objective) | How long writes are unavailable during failover | Detection + election + promotion + client reroute |
| Replication log | The ordered stream the leader ships: WAL, binlog, Raft log, commit log | Its format decides whether you can replicate across versions, engines or into other systems |
| Quorum | The minimum number of acks (W) or responses (R) out of N replicas | R + W > N makes read and write sets overlap. It does not make the system linearizable on its own |
| Epoch / term / generation | A number that increases on every leadership change | The basis of fencing. Storage rejects writes carrying an old epoch |
| Split brain | Two nodes accepting writes as leader for the same data | Divergent histories. Recovery means reconciling by hand or discarding writes |
The Three Models Side by Side#
| Dimension | Leader-Follower | Multi-Leader | Leaderless |
|---|---|---|---|
| Who accepts writes | One leader per shard | One leader per region (or per device) | Any replica; client or coordinator fans out |
| Write latency | Leader commit + sync followers (if any) | Local leader commit only | Wait for W of N acks (slowest of the fastest W) |
| Conflicts | None by construction (serialized at leader) | Inevitable: concurrent writes to the same key in two regions | Concurrent writes to the same key; resolved by timestamps, versions or merge |
| Failure of the writer | Failover required; writes stop for 30s–2min | Other leaders keep accepting writes | No failover; writes continue while W replicas are reachable |
| Read staleness | Followers lag; leader reads are fresh | Each region sees others' writes after cross-region lag | R + W > N gives overlap; R = 1 can be stale |
| Operational pain | Failover, promotion, split brain | Conflict resolution logic, schema of merges, debugging divergent states | Repair (read repair, anti-entropy), tombstones, tuning W/R per operation |
| Typical systems | PostgreSQL, MySQL, Raft-based stores, Redis | Cross-region active-active setups, offline-first clients, collaborative editors | Cassandra, Dynamo-style stores, Riak |
| Default for | ~90% of designs | Writes that must succeed in every region during a partition | Very high write availability with key-level, mergeable data |
For most production systems: leader-follower per shard, with a synchronous or quorum-acknowledged copy inside the region and asynchronous copies across regions. Reach for multi-leader or leaderless only when a named requirement forces it.
The Synchronicity Spectrum#
| Mode | Leader Acks the Client When | Added Write Latency | RPO on Leader Loss | What Breaks |
|---|---|---|---|---|
| Asynchronous | Local commit only | 0 | ≈ lag (ms to seconds; minutes under load) | Acknowledged writes vanish on failover |
| Semi-synchronous | ≥1 replica has received the change (written to its log) | +1 RTT (0.5–2ms cross-AZ) | 0 for one failure, if the fallback never triggered | Silent fallback to async on timeout |
| Synchronous (one named replica) | That replica has flushed to disk | +1 RTT + replica fsync | 0 for one failure | Replica stall = write stall |
| Quorum-synchronous | A majority (or ANY k of n) has flushed | +RTT to the k-th fastest replica | 0 while a majority survives | Needs ≥3 copies; cross-region quorum costs 60–150ms per write |
| Synchronous apply | Replica has applied and made it visible to reads | + apply time on replica | 0, and replicas are read-your-writes safe | Slowest mode; a slow replica slows every commit |
Physical vs Logical Replication#
| Physical (WAL / block shipping) | Logical (row changes / statements) | |
|---|---|---|
| What ships | Byte-level page changes | "Row X in table T changed from A to B" |
| Replica must be | Same engine, same major version, whole database | Any consumer that understands the row format; can be a subset of tables |
| Use for | HA standbys, read replicas, fast failover | Zero-downtime major-version upgrades, cross-engine migration, feeding search/cache/warehouse (CDC) |
| Pitfalls | No partial replication; a corrupt page ships too | DDL often not replicated; sequences and large transactions need care; consumer lag holds log on the source |
Statement-based replication (shipping the SQL text) is the third, older option. It breaks on non-deterministic statements (NOW(), RAND(), auto-increment races) and is rarely the right default today.
Change Data Capture Is Replication#
CDC reads the database's logical change stream and publishes it, usually to Kafka, for other systems to consume: a search index, a cache invalidator, a warehouse, another service's read model. It is replication to a different shape of store. That framing matters because CDC inherits every replication problem — lag, ordering per key, at-least-once delivery, and what happens when the consumer falls behind and the source has to retain log for it. A logical replication slot that nobody is reading can fill the leader's disk. Treat it as a replica with an owner, an alert and a lag SLO. Pipeline mechanics live in Batch & Stream Pipelines.
Core Strategies#
Strategy 1: Leader-Follower, Asynchronous#
One leader takes every write and streams its log to followers. Followers serve reads. The leader acks the client as soon as it has committed locally.
write(k, v):
leader.wal.append(k, v); leader.fsync() # ~0.1–1ms on SSD
ack(client) # client sees success here
for f in followers: f.ship_async(wal_position) # arrives 1–100ms later
read(k, session):
if session.needs_fresh: return leader.read(k)
return any_follower_with(lag < 1s).read(k) # route by measured lag
When to use: Read-heavy workloads where a few hundred milliseconds of staleness is acceptable and losing the last second of writes on a rare failover is an accepted business risk: catalogs, profiles, content, analytics read models. It is also the default for cross-region copies, because a synchronous cross-region commit puts 60–150ms on every write. Read scaling patterns built on this live in Read-Heavy Systems.
Failure mode: Lossy failover. The leader dies with 1.5s of writes not yet shipped. The promoted follower never saw them. The client got a success for each. When the old leader comes back, its extra writes conflict with new ones written at the same keys or auto-increment IDs. The fix is not "faster replication." It is deciding — before the incident — whether automation may promote a replica that is behind, and by how much.
Strategy 2: Leader-Follower with a Synchronous Quorum#
The leader waits for at least one (or k of n) replicas to confirm before acking. Majority-based consensus (Raft, Paxos) is the general version: a write is committed once a majority of the group has it in its log, and any new leader must come from that majority.
write(k, v):
leader.wal.append(k, v)
acks = wait_for(k_of_n = 1 of [replica_az_b, replica_az_c], timeout = 2s)
if acks < 1:
# the critical decision: stall writes, or degrade to async?
return policy == "durability_first" ? ERROR_UNAVAILABLE : ack_with_alert()
ack(client)
# PostgreSQL: synchronous_standby_names = 'ANY 1 (az_b, az_c)'
When to use: Any data where an acknowledged write must survive the loss of one node or one AZ: payments, orders, account state, inventory. Inside a region the cost is one cross-AZ round trip (~0.5–2ms), which is almost always affordable. Name ANY 1 of 2 (or a majority of 3), never "one specific replica," so a single slow replica does not stall every commit.
Failure mode: Silent degradation. A semi-synchronous setup that falls back to asynchronous after a timeout is, during the incident that matters, an asynchronous setup. If the fallback happens without an alert, your stated RPO of zero is false and nobody knows. Alert on replication.sync_replicas_connected < required and on every fallback event.
🎯 Staff Move: "Inside the region I'll require an ack from any one of two standbys in other AZs before we tell the client the order is placed. That costs about 1–2ms per commit and makes single-node failover lossless. Across regions I'll replicate asynchronously and publish the lag, because a synchronous cross-region commit would add ~70ms to every checkout."
Strategy 3: Multi-Leader#
Each region (or device) has its own leader that accepts writes locally and replicates them to the other leaders asynchronously. Every write is fast and local. The price is that two regions can modify the same record concurrently, and the system must decide what the merged result is.
write(k, v) in region R:
version = (hlc_now(), R) # hybrid logical clock + region tiebreak
local_leader.commit(k, v, version)
replicate_async(k, v, version) → other regions
on_receive(k, v_remote, ver_remote):
v_local, ver_local = store.get(k)
if concurrent(ver_local, ver_remote):
store.put(k, resolve(k, v_local, v_remote)) # LWW, merge, CRDT, or flag for a human
elif ver_remote > ver_local:
store.put(k, v_remote)
| Conflict Strategy | Works For | Silently Loses |
|---|---|---|
| Last-writer-wins (timestamp) | Overwrite-style fields: display name, settings | The concurrent write that lost; clock skew can pick the older one |
| Merge function (union, max, sum) | Sets, counters, flags | Nothing for commutative data; wrong for "remove" without tombstones |
| CRDTs | Collaborative text, counters, sets, maps (Collaborative Documents) | Nothing, at the cost of metadata growth and harder reasoning |
| Home region per record | Most user-owned data: route writes for user X to X's home | Write availability for X during a partition of X's home |
| Surface to the user | Documents, calendars, file sync (Cloud File Sync) | User time |
When to use: Writes must succeed in every region during a cross-region partition, or clients work offline. Before choosing it, check whether home-region routing (single leader per record, many leaders overall) meets the requirement. It usually does, with no conflict logic. The multi-region version of this decision is in Multi-Region.
Failure mode: Conflicts nobody designed. An auto-increment ID, a uniqueness constraint ("one username"), or a balance check ("don't go negative") cannot be enforced by two leaders that do not talk to each other synchronously. These invariants need a single owner per key, or they will break during the first long partition.
Strategy 4: Leaderless Quorums#
The client (or a coordinator node) sends each write to all N replicas for the key and waits for W acks. Reads ask R replicas and take the newest version. Replicas that missed a write catch up through read repair, hinted handoff and background anti-entropy.
N, W, R = 3, 2, 2 # R + W > N → read and write sets overlap
write(k, v):
ver = new_version(k) # timestamp or vector clock
send_to(replicas(k), k, v, ver)
await acks >= W else ERROR # slowest of the fastest 2
read(k):
resp = await responses >= R from replicas(k)
newest = max_by_version(resp)
repair_stale(resp, newest) # read repair in the background
return newest
| (N, W, R) | Tolerates for Writes | Tolerates for Reads | Use For |
|---|---|---|---|
| (3, 2, 2) | 1 replica down | 1 replica down | Balanced default; overlap guaranteed |
| (3, 1, 1) | 2 down | 2 down | Fast, can read stale or lose a write held by one node |
| (3, 3, 1) | 0 down | 2 down | Read-heavy, fresh reads, writes fragile |
| (5, 3, 3) | 2 down | 2 down | Two-AZ-failure tolerance, higher latency |
When to use: Very high write rates, key-value or wide-row access, a requirement that writes never wait for a failover, and data that can tolerate or merge concurrent updates. Time series, event logs, carts, activity feeds. Cassandra and DynamoDB (internally) are the standard references. The full end-to-end design is in Replicated Key-Value Store.
Failure mode: Believing R + W > N means linearizable. Sloppy quorums (writing to a substitute node during failures), concurrent writes resolved by last-writer-wins, and a read racing a partially applied write all break it. Also: deletes. A tombstone that is purged before every replica has seen it lets the deleted value resurrect through repair.
Strategy 5: Logical Replication and CDC as a First-Class Replica#
Publish the row-level change stream with a stable schema, ordered per key, and let consumers build their own copies.
source DB → logical decoding slot "orders_cdc"
→ connector (at-least-once, ordered per primary key)
→ topic orders.changes (partitioned by order_id, 7-day retention)
→ consumers: search indexer, cache invalidator, warehouse loader
consumer.apply(event):
if event.lsn <= store.last_applied_lsn(event.key): return # idempotent
store.upsert(event.key, event.after)
When to use: Keeping search, caches and analytics in sync without dual writes; zero-downtime migrations (replicate old → new, verify, cut over); feeding other teams without giving them database credentials. Search Engine is the classic consumer.
Failure mode: An abandoned slot. The consumer stops, the source must retain its log until the slot advances, and the leader's disk fills over hours or days. Alert on replication.slot_retained_bytes and give every slot a named owner and a maximum retention after which it is dropped.
Lag and Failover: The Hard Sub-Problem#
Choosing a model takes one sentence. Living with it means two problems: lag, which makes reads wrong while everything is healthy, and failover, which makes writes wrong when something is not.
What Lag Looks Like to a User#
| Anomaly | What the User Sees | Cause | Targeted Fix |
|---|---|---|---|
| Read-your-writes violation | "I posted a comment and it's gone" | Write went to leader, next read hit a lagging follower | Return the commit LSN/version to the client; route reads to a replica at or past it, else to the leader |
| Monotonic reads violation | Comment appears, refresh, gone, refresh, back | Successive reads hit replicas with different lag | Pin a session to one replica (hash on user ID), or carry the "highest seen" version |
| Consistent prefix violation | A reply appears before the message it answers | Partitions replicate independently; reader sees partition B ahead of A | Keep causally related writes in one partition, or carry causal dependencies |
| Stale read after failover | Balance looks like it did 2 seconds ago, permanently | Lossy promotion: the new leader never received the last writes | Sync or quorum replication for that data class |
These are the session guarantees described in Consistency, CAP & PACELC. The replication-side lesson is that they are cheap to provide per session and expensive to provide globally. Provide them where a user would notice.
Where Lag Comes From#
| Source | Typical Lag | Signature |
|---|---|---|
| Normal in-region streaming | 1–50ms | Flat, low, below the alert threshold |
| Cross-region streaming | RTT (60–150ms) + apply time | Floor equals the network RTT |
| Single-threaded apply on the replica | Seconds to minutes | Lag climbs during write bursts while the leader is fine |
Long-running transaction or large batch (UPDATE of 50M rows) | Minutes | Lag jumps when the big transaction commits, then drains |
| Replica serving heavy analytics queries | Seconds to minutes | Replay paused or slowed by query conflicts; lag tracks the report schedule |
| Schema change on the leader | Minutes to hours | Lag starts exactly when the migration does |
🎯 Staff Insight: "Average lag is a vanity metric. I alert on p99 lag per replica, and the read router takes any replica over 1 second out of rotation automatically. Otherwise the load balancer happily sends 30% of reads to the one replica that is two minutes behind."
The Anatomy of a Failover#
t=0 Leader stops responding (host crash, kernel hang, network partition).
t=+10s Health checks miss 3 consecutive probes at 3s intervals → leader suspected.
t=+15s Coordinator (Raft majority, Orchestrator, Patroni + etcd) confirms with peers.
t=+20s Pick the most up-to-date replica. Compare replication positions.
t=+25s Fence the old leader: bump epoch, revoke its lease, block it at the proxy.
t=+30s Promote the replica; it starts accepting writes in the new epoch.
t=+35–90s Clients rediscover the leader (proxy reconfig, service discovery, DNS TTL).
RTO ≈ 30s–2min. RPO = 0 if sync/quorum, else ≈ lag at t=0.
The step most designs skip is t=+25s. A leader that is merely partitioned, not dead, is still accepting writes from clients on its side. Without fencing you now have two leaders.
Split Brain and Fencing#
Timeouts cannot tell "dead" from "slow" or "unreachable from here." So any automated failover will eventually promote a replica while the old leader is still alive. The design question is what stops the old leader's writes from landing.
| Mechanism | How It Works | Gap |
|---|---|---|
| Leader lease | Leader may write only while it holds a lease (e.g., 10s) from a consensus store; new leader waits for expiry | Relies on bounded clock drift and the old leader checking its lease before every write |
| Epoch fencing at storage | Every write carries the leader's epoch; storage or replicas reject anything below the current epoch | Needs storage that checks; the strongest option |
| STONITH (power off the old node) | Kill the old leader through out-of-band control before promotion | Needs a working out-of-band path; feels crude, works |
| Quorum-based commit | A leader that cannot reach a majority cannot commit anything | Built into Raft/Paxos; the reason to prefer them over home-made failover |
Fencing tokens in depth are covered in Distributed Locking and Distributed Consensus.
🎯 Staff Move: "Automatic failover inside a region, because the majority is inside the region and fencing is cheap. Cross-region promotion requires a human, a check of the replication position, and confirmation that the old primary is fenced. A 40-second network blip should never be able to move our primary across a continent on its own."
When NOT to Add a Synchronous Replica#
- Cross-region, on the hot write path. 60–150ms on every commit, and a region-to-region link degradation becomes a write outage everywhere. Use async cross-region plus an explicit RPO.
- With exactly one sync target and no fallback plan. A single sync standby turns its stall into your stall. Use
ANY 1 of 2or a majority of 3. - For data you can rebuild. Caches, search indexes, derived feeds and analytics read models can be recomputed from the source of truth. Paying 1–2ms per write to protect them is waste.
- For high-volume, low-value writes. Clickstream, view counts and last-seen timestamps can tolerate losing a second. Per-transaction durability settings let the ledger and the clickstream share a database with different guarantees.
- When the real risk is region loss. A sync replica in the next AZ does nothing if the whole region goes down. Be honest about which failure you are buying protection against.
When NOT to Go Multi-Leader#
- When home-region routing works. If each record has a natural owner region (a user, a tenant, a store), single leader per record gives local writes for most traffic with zero conflict logic.
- When the data has invariants. Uniqueness, non-negative balances, inventory counts and sequential IDs cannot be enforced by independent leaders.
- When nobody owns the merge rules. Conflict resolution is product logic. If no team will own "what happens when two admins edit the same setting in two regions," last-writer-wins will decide for them, silently.
- When the requirement is read latency, not write availability. Async read replicas in every region solve read latency. Multi-leader only solves "writes must succeed during a partition."
Visual Guide#
Synchronous vs Asynchronous Commit#
Choosing a Replication Model#
Leaderless Write and Read Overlap (N=3, W=2, R=2)#
Any 2 replicas a reader picks must include at least one of the 2 that took the write. That is all R + W > N promises. It says nothing about a write that is still in flight, or two concurrent writes.
Failover With and Without Fencing#
Logical Replication Feeding Many Copies#
Implementation Patterns#
Read Routing by Replication Position#
on write: resp.headers["X-Min-Version"] = commit_lsn
on read: min = req.headers["X-Min-Version"] or session.last_seen_lsn
replica = pick(r for r in replicas if r.replayed_lsn >= min and r.lag < 1s)
return replica or leader # fall back, never serve older than min
This gives read-your-writes and monotonic reads per session without sending all reads to the leader. In practice 95%+ of reads still hit replicas, because most sessions have not written recently.
Failure Scenario: The Lossy Cross-Region Promotion#
t=0 A transit link between us-east and us-west degrades for 45 seconds.
t=+15s Failover tooling in us-west no longer sees the us-east primary. Lag at cut: 1.8s.
t=+25s Tooling promotes the us-west replica. Apps in us-west start writing there.
t=+45s Link recovers. us-east primary has ~2,400 writes the new primary never received.
us-west has ~1,100 new writes, some to the same rows and ID ranges.
t=+3min Every app in us-east now pays ~70ms per query to reach the new primary. p99 triples.
t=+10min Incident commander faces the choice: fail back and discard 1,100 writes,
or stay and discard 2,400. Neither is acceptable. Writes are paused.
t=+20h Both write sets reconciled by hand from binlogs. Service restored.
Detection: db.primary.region changed; replication.position_gap_at_promotion > 0; app db.query.latency.p99 jumping by one cross-region RTT.
Blast radius: every service on that cluster, in both regions, for hours.
Mitigation: stop writes on both sides; take the side with more critical writes as truth; replay the other side's writes from the log after review.
Prevention: cross-region promotion requires a human and a zero (or approved) position gap; fence the old primary before promoting; keep apps and their primary in the same region.
Owner: the database platform team owns the failover policy; the data owner (payments, orders) owns the "accept N seconds of loss" sign-off.
Operational Reality Matrix#
| Failure | Detection Signal | Blast Radius | Mitigation | Owner |
|---|---|---|---|---|
| Replica lag spike | replica.lag_seconds p99 > 1s | Stale reads on that replica | Remove from read pool; find the big transaction or query | Database platform |
| Sync replica lost, silent fallback to async | replication.sync_replicas_connected < required; fallback counter | RPO quietly becomes nonzero | Restore or add a sync replica; page on fallback | Database platform |
| Split brain | Two nodes report role=primary; epoch mismatch errors | Every writer on that shard | Fence the stale leader; freeze writes; reconcile | Database platform + data owner |
| Abandoned CDC slot | replication.slot_retained_bytes growing; disk free < 20% | Leader disk full → write outage | Drop or advance the slot after owner sign-off | Consumer team owns the slot; platform enforces limits |
| Lossy failover | replication.position_gap_at_promotion > 0 | Lost acknowledged writes; possible duplicates and ID reuse | Reconcile from the old leader's log | Data owner signs off |
| Multi-leader conflict storm | replication.conflicts_resolved per minute 10× baseline | Silently overwritten user data | Pin hot keys to a home region; review merge rules | Product team owning the data |
The Numbers in Context#
| Number | Value | What It Means for Your Design |
|---|---|---|
| In-region async lag (healthy) | ~1–50ms | Read-after-write on a replica fails for a few ms. Route the writer's next read carefully, leave everyone else on replicas |
| Cross-region async lag | RTT (60–150ms) + apply; seconds under load | Remote regions are always at least one RTT behind |
| Cross-AZ sync commit cost | +0.5–2ms per commit | Affordable for almost every OLTP write |
| Cross-region sync commit cost | +60–150ms per commit | One synchronous continent hop per write is a product decision |
| Failure detection | 10–30s (3 missed heartbeats at 3–10s) | Faster detection means more false failovers on network blips |
| In-region automated failover (RTO) | ~30s–2min | Budget it in the SLO: a few failovers a year fits 99.95% |
| Async RPO | lag × write rate: 2s × 5,000/s = ~10,000 writes | Say the count, not the seconds, to the data owner |
| Quorum (3, 2, 2) | Survives 1 replica loss for reads and writes | 3 copies across 3 AZs is the standard floor |
| Six copies, 4/6 write, 3/6 read | Survives an AZ plus one more node | What it costs to tolerate AZ + 1 failures with quorums |
| Replica read capacity | Each replica ≈ leader's read capacity; writes are replayed on every replica | Replicas scale reads, never writes |
| Retention for a lagging logical slot | Write rate × outage time: 50MB/s × 6h ≈ 1TB | One forgotten consumer can fill the primary's disk overnight |
How This Shows Up in Interviews#
Scenario 1: "The database is a single point of failure. Fix it."#
Do not stop at "add a replica." Say: "Two standbys in other AZs, commit waits for any one of them, so single-node or single-AZ loss loses no acknowledged writes and costs ~1–2ms per commit. Failover is automated through a consensus store with epoch fencing, about 30–60 seconds of write unavailability. Region loss is a separate decision: an async replica in a second region with a stated RPO of a few seconds and a human-approved promotion."
Scenario 2: "Reads are 95% of traffic and the primary is at 80% CPU." (Full Walkthrough)#
Step 1 — Confirm what's hot. "First I check it's reads, not a few expensive queries or writes. If it's 95% simple reads, replicas are the right lever. If it's three bad queries, an index is."
Step 2 — Add async replicas. "Three read replicas, async, in-region. Each one handles roughly what the primary handles for reads, so we take the primary from ~80% to ~25% CPU."
Step 3 — Classify reads by staleness tolerance. "Product pages, search results and other users' profiles go to replicas with a 1-second lag cap. A user's own just-written data and anything in checkout reads from the primary, or from a replica past the user's last write position."
Step 4 — Protect against lag. "The read router removes any replica over 1s of p99 lag. If all replicas are lagging, reads fail over to the primary with a concurrency cap so the primary doesn't fall over too."
Step 5 — Name the limit. "Replicas don't scale writes. When write load reaches ~60% of the primary, the next step is sharding, and that's a different conversation."
Step 6 — Owners. "The database platform team owns replica health and the lag SLO. Each product team owns the classification of its reads, because only they know which ones a user would notice."
Why this is a Staff answer: it classifies reads by who notices staleness, protects the primary from the fallback path, states that replication does not scale writes, and assigns the read classification to the team that owns the product consequence.
Scenario 3: "We need to be active in two regions."#
This tests whether you jump to multi-leader. "Which requirement: low read latency in both regions, surviving a region loss, or writes that succeed during a partition? The first two need async replicas and a promotion plan. Only the third needs multi-leader, and most of it is solved by giving each user a home region. I'd only accept true multi-leader for data that merges cleanly."
Scenario 4: "After a failover, some customers were charged twice and some orders vanished."#
This tests failover forensics. "Classic lossy async promotion. Orders acked by the old leader in its last second never reached the new leader, so they 'vanished'. Clients retried, and the new leader reused ID ranges, so some charges duplicated. Short term: reconcile from the old leader's binlog. Long term: payments get quorum-synchronous commit, idempotency keys on every charge, and no promotion with a nonzero position gap without sign-off." Link the idempotency piece to Payment Processing.
Advanced Patterns#
| Pattern | How It Works | When to Use |
|---|---|---|
| Chain replication | Writes enter the head, flow down a chain, ack from the tail; reads served by the tail | Strong reads with simple failure handling, when chain latency is acceptable |
| Witness / tie-breaker replica | A cheap node in a third site votes in elections but stores little or no data | Two data sites that still need a majority to avoid split brain |
| Delayed replica | A replica that deliberately applies the log 1–24h behind | Recovery from a bad DELETE or a corrupting deploy faster than a backup restore |
| Cascading replication | Replicas replicate from other replicas, not the leader | Many replicas or a remote region without loading the leader's network |
| Per-transaction durability | The ledger commit waits for a sync replica, the clickstream insert does not | Mixed-value data in one database |
| Log-structured storage replication | Ship only the redo log to a quorum of storage nodes and let storage build pages | Cloud databases that separate compute from storage |
| Zero-downtime migration via logical replication | Replicate old → new, verify row counts and checksums, flip writes, keep reverse replication briefly | Major-version upgrades, engine migrations, re-sharding |
Beyond Staff: The Principal View#
Why L7 Sees This Problem Differently#
A Staff engineer picks the right replication mode for their database and writes the failover runbook. A Principal engineer notices that the company runs 400 databases, each with a replication setting someone chose on the day it was created, and that nobody can say what the company loses if us-east-1 disappears for an hour. Some ledgers are async across regions; some caches are synchronously replicated for no reason; three teams have automatic cross-region promotion and nobody has tested what happens when it fires during a blip. At L7, replication stops being a topology choice and becomes a data-loss policy: every data class has a stated RPO and RTO, the platform enforces the settings that meet them, and a named person approves every lossy failover.
🧭 Principal Move: "I don't want 400 teams choosing replication modes. I want five data classes, each with a published RPO and RTO, and a platform that sets the replication mode from the class. Then 'can we lose two seconds of this?' is answered once, in a policy the business signed, not at 3 a.m. by whoever holds the pager."
The Org-Level Fault Line#
One managed database platform with policy-driven replication vs each team running its own databases.
| Option | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Each team runs its own | Autonomy; teams tune for their workload | Inconsistent RPO; untested failovers; automation that promotes across regions on a blip | The data owner, after the incident; customers whose writes were lost |
| Central DBA team approves every change | Consistent, reviewed settings | Bottleneck; tickets take weeks; teams route around it with self-hosted stores | Product velocity |
| Platform with data classes (Tier 0 ledger → Tier 4 derived) | Replication mode, failover authority and backups derived from the class; self-service within the class | Classification arguments; exceptions process needed; platform becomes tier-0 | Platform team (6–15 engineers); a one-time classification effort by every team |
The Principal position: classify data, not databases. The class sets the floor (sync in-region, async cross-region with lag SLO, human-approved cross-region promotion for Tier 0). Teams can exceed the floor freely and go below it only through an exception that expires.
Cost Model#
Assumptions: managed-database pricing in the style of the large clouds; a replica costs about the same as the primary; cross-region transfer ~$0.02/GB; fully loaded engineer ~$25K/month; each database instance ~$1–3K/month at the mid-size tier.
| Scale | Fleet | Replication Spend | Platform Headcount | On-Call Load | The Expensive Mistake |
|---|---|---|---|---|---|
| Startup | 5 databases, 1 region, 1 standby each | ~$5–10K/month (standbys double instance cost) | 0.5 engineer | A failover a quarter, mostly managed by the cloud | No cross-region copy at all; region loss = restore from backup, hours of RTO |
| Growth | 80 databases, 2 regions, 2 in-region + 1 cross-region replica each | ~$250–400K/month; cross-region transfer ~$10–20K/month | 4–6 engineers (~$125K/month) | Dedicated rotation; lag and slot incidents weekly | Async replication for ledgers because "it's the default" |
| Large | 1,000+ databases, 3–5 regions | ~$4–8M/month; replicas are ~50–65% of database spend | 12–20 engineers | Tier-0 rotation; quarterly region-evacuation game days | Over-replicating Tier 3–4 data: 5 copies of rebuildable caches cost ~$1M/month |
The lever most orgs miss: data classification pays for itself at Large scale. Dropping derived and rebuildable stores from 3 copies to 2 (and from cross-region sync to async or none) typically frees 15–25% of database spend, which funds the platform team several times over.
The 3-Year Evolution Path#
One-Way Doors vs Two-Way Doors#
| Decision | Door | Reversal Cost |
|---|---|---|
| Sync vs async for one database | Two-way | Config change and a restart; minutes |
| Read routing and lag caps | Two-way | Router config |
| Choosing leaderless or multi-leader for a dataset | One-way | Application logic now assumes conflicts and merges; going back means a migration and new invariants |
| Exposing a CDC stream to other teams | Mostly one-way | The change-event schema becomes a public contract; consumers pin it for years |
| Automatic cross-region promotion | Two-way to switch off, one-way once it fires badly | A single bad promotion can mean a day of reconciliation |
| Primary region placement | One-way | Data gravity; every dependent service's latency is built around it |
| Conflict resolution rule (LWW vs merge) | One-way per dataset | Data already overwritten by LWW cannot be recovered |
The Standard I'd Write#
RFC: Data Replication and Failover Standard (v1)
Scope: All production stores of record. Caches and derived stores follow the Tier 3–4 rules.
MUST:
- Every store declares a data class: Tier 0 (money, legal), Tier 1 (user-created content), Tier 2 (operational state), Tier 3 (derived, rebuildable), Tier 4 (ephemeral).
- Tier 0–1: commit waits for a quorum or
ANY 1 of ≥2replicas in other AZs. RPO = 0 for single-AZ loss.- Tier 0–2: an asynchronous copy in a second region with p99 lag ≤ 5s, alerted.
- Cross-region promotion of Tier 0–1 MUST be human-approved by the data owner, after confirming the replication gap and fencing the old primary.
- Every logical replication slot or CDC consumer has a named owning team, a lag SLO and a maximum retention after which the platform may drop it.
SHOULD: route reads by replication position for session guarantees; keep apps in the same region as their primary; run a failover drill per Tier 0–1 cluster each quarter.
Exceptions: Filed with the database platform team; approved by the data owner's director; expire after 2 quarters.
Success metrics: zero unapproved lossy failovers per year; 100% of Tier 0–1 stores drilled each quarter with measured RTO; replication spend on Tier 3–4 stores down 20% within a year.
What I'd Tell the VP#
"Every database keeps extra copies, but today each team decides how up to date those copies are and whether a machine can switch to them on its own. That means we can't answer how much customer data we'd lose if a region went down. I want us to sort our data into five classes, decide with the business how much loss each class can tolerate, and have our database platform enforce it. For payments and orders the answer will be 'none,' which costs about a millisecond per write. For data we can rebuild we'll keep fewer copies and save money. The risk is a few months of classification work across teams. The payoff is that the next regional outage is a drill, not a data-loss incident."
Principal Interview Signals#
| Signal | What It Sounds Like |
|---|---|
| Turns modes into policy | "I'd define RPO and RTO per data class and derive the replication mode from the class, not from each team's preference." |
| Owns lossy-failover authority | "Automation may promote within a region. Crossing a region with a nonzero gap needs the data owner's yes." |
| Prices the copies | "Replicas are over half our database bill. A third of them protect data we could rebuild in an hour." |
| Treats CDC as a contract | "Once five teams consume this change stream, its schema is an API. It gets versioning and an owner." |
| Plans the drill, not just the design | "Every Tier 0 cluster fails over once a quarter on a weekday afternoon, so the first real failover isn't the first one." |
Staff answers that L7 interviewers find insufficient:
- "We'll use synchronous replication for everything important." — Correct for one database; doesn't say who decides what is important or what it costs fleet-wide.
- "Automatic failover with a 30-second timeout." — Right in a region; silent on who approves a lossy cross-region promotion.
- "We'll publish changes with CDC." — Ignores that the stream becomes a contract and an abandoned slot can fill the primary's disk.
How Real Companies Built It#
These are public, documented examples.
GitHub: The October 2018 Cross-Region Failover#
On 21 October 2018, a 43-second loss of connectivity between GitHub's US East Coast network hub and its primary East Coast data center led its failover tooling (Orchestrator) to promote database primaries in the US West Coast. When connectivity returned, the East Coast primaries held writes that had not replicated west (954 writes on the most affected cluster), the West Coast primaries had taken new writes, and applications in the East could not tolerate a cross-country round trip on most database calls. GitHub chose to protect data integrity, and the incident lasted 24 hours and 11 minutes while data was restored and reconciled (GitHub blog).
Staff insight: The partition lasted under a minute; the damage lasted a day. The lesson is not "turn off automation." It is that cross-region promotion is a business decision about lost writes and latency, and the tooling should require a human, a position check and fencing before it crosses a region.
Amazon Aurora: Six Copies and an AZ+1 Quorum#
Aurora's storage layer replicates each database volume six ways across three Availability Zones, with a write quorum of four and a read quorum of three. AWS describes this as tolerating the loss of an entire AZ without losing write availability, and an AZ plus one more failure without losing data; volumes are split into 10 GB segments so a lost copy can be re-replicated quickly (AWS Database Blog).
Staff insight: The number of copies follows from the failure you want to survive. A 2-of-3 quorum loses write availability when one AZ goes down and another node fails at the same time; tolerating "AZ + 1" pushes you to six copies. Say the failure model first, then the copy count.
PostgreSQL: Synchronous Commit, Quorum Standbys#
PostgreSQL lets a commit wait for standbys at several levels: remote_write (received and written to the standby's OS), on (flushed to the standby's WAL on disk) and remote_apply (replayed and visible to queries). synchronous_standby_names accepts FIRST k (…) for priority-based or ANY k (…) for quorum-based standbys, and the documentation warns that commits may never complete if a required synchronous standby crashes and no other is available (PostgreSQL docs).
Staff insight: Synchronous replication is a dial, not a switch. ANY 1 (a, b) gives zero-loss single failover without making one standby a single point of failure, and the durability level can be set per transaction. More on these dials in PostgreSQL.
MySQL: Semi-Synchronous Replication Falls Back to Async#
MySQL's semi-synchronous replication makes the source wait until at least one replica acknowledges that it has written and flushed the transaction's events to its relay log. If no replica acknowledges within the configured timeout, the source reverts to asynchronous replication, and returns to semi-synchronous once a replica catches up (MySQL docs).
Staff insight: The fallback is a deliberate availability choice, and it means your "zero data loss" claim holds only while the fallback is not active. Monitor and page on it, or you will discover your real RPO during the failover that needed it.
Staff Calibration#
What Staff Engineers Say (That Seniors Don't)#
| Concept | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| Durability | "We'll add a replica for HA" | "Commit waits for any 1 of 2 cross-AZ standbys: zero acked-write loss on one failure for ~1–2ms per commit" | "RPO is set per data class in a signed policy; the platform derives the mode, so nobody re-decides it per database" |
| Read scaling | "Send reads to replicas" | "Replicas with a 1s lag cap; the writer's next reads go by commit position; replicas don't scale writes" | "Read classification is owned by product teams, lag SLOs by the platform. The contract between them is the position token" |
| Failover | "Automatic failover" | "Automatic in region with epoch fencing, ~30–60s RTO; cross-region needs a human and a zero position gap" | "Lossy failover authority is named per data class and drilled quarterly; the first real failover is never the first one" |
| Multi-region writes | "Go multi-master" | "Home-region routing first; multi-leader only for mergeable data, with a named conflict rule" | "Active-active is a product commitment with a permanent conflict-resolution owner. I'd price that headcount before agreeing" |
| CDC | "Stream changes to Kafka" | "Logical slot with an owner, ordered per key, idempotent consumers, a retention cap" | "The change stream is a versioned public API across teams; schema changes go through review like any other contract" |
| Cost | "Replicas are cheap" | "Each replica is a full instance; three copies triples storage and compute" | "Replicas are over half of database spend; dropping copies of rebuildable data funds the platform team" |
Why "Failover" separates levels
The Senior answer is correct: automatic failover reduces downtime. The Staff answer notices that failover is the moment data is lost, separates in-region (safe to automate with fencing) from cross-region (needs a position check and a human), and states RTO and RPO in numbers. The Principal answer notices that the person who should approve losing 954 writes is not the on-call engineer and not the database team, but the owner of the data — and makes that a standing policy with a drill, so the decision is made in daylight before the incident.
Why "Multi-region writes" separates levels
"Multi-master" sounds like the strongest answer, which is why it is the most common trap. Staff engineers ask which requirement forces it and usually find that home-region routing meets it. Principal engineers recognize that a conflict-resolution rule is product logic that someone must own forever, and that choosing last-writer-wins silently decides which customer's data is discarded.
Common Interview Traps#
- "Add a replica" without a mode. Sync or async, how many, where. Say the RPO and the added write latency.
- Treating replicas as write scaling. Every replica replays every write. Replicas scale reads; sharding scales writes.
- Ignoring the writer's own reads. The user who just wrote is the one who notices lag. Route their reads by position.
- Assuming
R + W > Nis linearizable. It guarantees overlap, not ordering of concurrent writes or atomicity of in-flight ones. - Automatic cross-region promotion. Fast detection plus no fencing equals split brain on the next network blip.
- One named sync standby. Its stall becomes a write outage, or a silent fallback to async.
- Forgetting the CDC slot. An abandoned consumer can fill the leader's disk.
- Multi-leader for data with invariants. Unique usernames and non-negative balances need one owner per key.
Practice Drill#
Prompt: "Our orders database runs in us-east with two async read replicas in the same region and one async replica in eu-west. Last month the primary's host died, we promoted a replica in 40 seconds, and support later found ~300 orders customers had been charged for but that no longer existed. Leadership wants this to never happen again, and the EU team wants local writes. What do you do?"
Staff Answer
The missing orders are the asynchronous RPO: lag at the moment of failure times the write rate. At ~150 orders/s and ~2s of lag under peak load, ~300 lost orders is exactly what this setup promises. So the fix is the replication mode, not the failover speed. (1) Move the two in-region replicas to standbys in separate AZs and make commit wait for ANY 1 of them; that costs ~1–2ms per order write and makes single-node and single-AZ failover lossless. (2) Put failover behind a consensus store with epoch fencing so a partitioned old primary cannot keep accepting writes; target RTO 30–60s. (3) Charges and orders must be idempotent end to end: an idempotency key on the charge, and the order written before capture, so a lost order and a retried charge can be reconciled automatically. (4) Reconcile last month's ~300 orders from the old primary's log. (5) For the EU ask, I would not go multi-leader on orders: inventory and payment state have invariants two leaders can't enforce. Instead, EU reads come from the eu-west replica with a lag SLO of p99 ≤ 2s, and EU writes go to us-east (~80ms extra per order write, acceptable for checkout) until the business has a reason to give EU customers a home region of their own. Owners: database platform owns the replication config and failover tooling; the orders team owns idempotency and reconciliation; the head of payments signs off on the RPO for this data class.
Why this is L6:
- Diagnoses the incident as async RPO with a number, instead of blaming the failover.
- Uses
ANY 1 of 2cross-AZ, with the latency cost stated, rather than a single named sync replica. - Adds fencing and idempotency so the next failure is safe even if replication misbehaves.
- Refuses multi-leader for data with invariants and offers a cheaper answer to the EU requirement.
What L7 adds:
- Turns "orders" into a Tier 0 data class with a written RPO of zero for single-AZ loss and a stated cross-region RPO, so every store holding money gets the same treatment.
- Makes lossy cross-region promotion a decision owned by the head of payments, rehearsed in a quarterly game day.
- Prices the change:
1–2ms per write and one extra standby ($2–3K/month for this cluster) against the refund, support and trust cost of ~300 phantom charges a failover.
Where This Appears#
- Replicated Key-Value Store — designing a leaderless replicated store end to end: quorums, repair, conflict resolution
- Multi-Region Architecture — async cross-region replicas, home-region routing and regional failover
- Consensus Service — Raft and Paxos as quorum-synchronous replication with safe leader election
- Database Sharding — replication inside each shard, and why replicas don't scale writes
- Search Engine — CDC as replication into a differently shaped store
- Payment Processing — why money needs zero-RPO commit plus idempotency
- Distributed Locking — fencing tokens that stop a deposed leader's writes
- Distributed Caching — replicated caches and when not to replicate rebuildable data synchronously
Related Foundations & Patterns: Consistency, CAP & PACELC · Partitioning · Estimation on a Whiteboard · Scaling Reads · Batch & Stream Pipelines · Distributed Coordination · Graceful Degradation: Fail Open or Closed
Related Technologies: PostgreSQL · Cassandra · DynamoDB · Kafka · ZooKeeper & etcd · Redis