Hiring BarSupport

Replication Fundamentals

Foundation40 min read6 diagrams

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 > N for 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#

TermMeaningWhy It Matters in an Interview
Replication lagTime (or log bytes) between a write committing on the leader and being visible on a replicaSets 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 failoverAsync: ≈ lag at the moment of failure. Sync: 0 for a single failure
RTO (recovery time objective)How long writes are unavailable during failoverDetection + election + promotion + client reroute
Replication logThe ordered stream the leader ships: WAL, binlog, Raft log, commit logIts format decides whether you can replicate across versions, engines or into other systems
QuorumThe minimum number of acks (W) or responses (R) out of N replicasR + W > N makes read and write sets overlap. It does not make the system linearizable on its own
Epoch / term / generationA number that increases on every leadership changeThe basis of fencing. Storage rejects writes carrying an old epoch
Split brainTwo nodes accepting writes as leader for the same dataDivergent histories. Recovery means reconciling by hand or discarding writes

The Three Models Side by Side#

DimensionLeader-FollowerMulti-LeaderLeaderless
Who accepts writesOne leader per shardOne leader per region (or per device)Any replica; client or coordinator fans out
Write latencyLeader commit + sync followers (if any)Local leader commit onlyWait for W of N acks (slowest of the fastest W)
ConflictsNone by construction (serialized at leader)Inevitable: concurrent writes to the same key in two regionsConcurrent writes to the same key; resolved by timestamps, versions or merge
Failure of the writerFailover required; writes stop for 30s–2minOther leaders keep accepting writesNo failover; writes continue while W replicas are reachable
Read stalenessFollowers lag; leader reads are freshEach region sees others' writes after cross-region lagR + W > N gives overlap; R = 1 can be stale
Operational painFailover, promotion, split brainConflict resolution logic, schema of merges, debugging divergent statesRepair (read repair, anti-entropy), tombstones, tuning W/R per operation
Typical systemsPostgreSQL, MySQL, Raft-based stores, RedisCross-region active-active setups, offline-first clients, collaborative editorsCassandra, Dynamo-style stores, Riak
Default for~90% of designsWrites that must succeed in every region during a partitionVery 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#

ModeLeader Acks the Client WhenAdded Write LatencyRPO on Leader LossWhat Breaks
AsynchronousLocal commit only0≈ 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 triggeredSilent fallback to async on timeout
Synchronous (one named replica)That replica has flushed to disk+1 RTT + replica fsync0 for one failureReplica stall = write stall
Quorum-synchronousA majority (or ANY k of n) has flushed+RTT to the k-th fastest replica0 while a majority survivesNeeds ≥3 copies; cross-region quorum costs 60–150ms per write
Synchronous applyReplica has applied and made it visible to reads+ apply time on replica0, and replicas are read-your-writes safeSlowest mode; a slow replica slows every commit

Physical vs Logical Replication#

Physical (WAL / block shipping)Logical (row changes / statements)
What shipsByte-level page changes"Row X in table T changed from A to B"
Replica must beSame engine, same major version, whole databaseAny consumer that understands the row format; can be a subset of tables
Use forHA standbys, read replicas, fast failoverZero-downtime major-version upgrades, cross-engine migration, feeding search/cache/warehouse (CDC)
PitfallsNo partial replication; a corrupt page ships tooDDL 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 StrategyWorks ForSilently Loses
Last-writer-wins (timestamp)Overwrite-style fields: display name, settingsThe concurrent write that lost; clock skew can pick the older one
Merge function (union, max, sum)Sets, counters, flagsNothing for commutative data; wrong for "remove" without tombstones
CRDTsCollaborative text, counters, sets, maps (Collaborative Documents)Nothing, at the cost of metadata growth and harder reasoning
Home region per recordMost user-owned data: route writes for user X to X's homeWrite availability for X during a partition of X's home
Surface to the userDocuments, 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 WritesTolerates for ReadsUse For
(3, 2, 2)1 replica down1 replica downBalanced default; overlap guaranteed
(3, 1, 1)2 down2 downFast, can read stale or lose a write held by one node
(3, 3, 1)0 down2 downRead-heavy, fresh reads, writes fragile
(5, 3, 3)2 down2 downTwo-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#

AnomalyWhat the User SeesCauseTargeted Fix
Read-your-writes violation"I posted a comment and it's gone"Write went to leader, next read hit a lagging followerReturn the commit LSN/version to the client; route reads to a replica at or past it, else to the leader
Monotonic reads violationComment appears, refresh, gone, refresh, backSuccessive reads hit replicas with different lagPin a session to one replica (hash on user ID), or carry the "highest seen" version
Consistent prefix violationA reply appears before the message it answersPartitions replicate independently; reader sees partition B ahead of AKeep causally related writes in one partition, or carry causal dependencies
Stale read after failoverBalance looks like it did 2 seconds ago, permanentlyLossy promotion: the new leader never received the last writesSync 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#

SourceTypical LagSignature
Normal in-region streaming1–50msFlat, low, below the alert threshold
Cross-region streamingRTT (60–150ms) + apply timeFloor equals the network RTT
Single-threaded apply on the replicaSeconds to minutesLag climbs during write bursts while the leader is fine
Long-running transaction or large batch (UPDATE of 50M rows)MinutesLag jumps when the big transaction commits, then drains
Replica serving heavy analytics queriesSeconds to minutesReplay paused or slowed by query conflicts; lag tracks the report schedule
Schema change on the leaderMinutes to hoursLag 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.

MechanismHow It WorksGap
Leader leaseLeader may write only while it holds a lease (e.g., 10s) from a consensus store; new leader waits for expiryRelies on bounded clock drift and the old leader checking its lease before every write
Epoch fencing at storageEvery write carries the leader's epoch; storage or replicas reject anything below the current epochNeeds storage that checks; the strongest option
STONITH (power off the old node)Kill the old leader through out-of-band control before promotionNeeds a working out-of-band path; feels crude, works
Quorum-based commitA leader that cannot reach a majority cannot commit anythingBuilt 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 2 or 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#

Diagram: Synchronous vs Asynchronous Commit

Choosing a Replication Model#

Diagram: Choosing a Replication Model

Leaderless Write and Read Overlap (N=3, W=2, R=2)#

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

Diagram: Failover With and Without Fencing

Logical Replication Feeding Many Copies#

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

FailureDetection SignalBlast RadiusMitigationOwner
Replica lag spikereplica.lag_seconds p99 > 1sStale reads on that replicaRemove from read pool; find the big transaction or queryDatabase platform
Sync replica lost, silent fallback to asyncreplication.sync_replicas_connected < required; fallback counterRPO quietly becomes nonzeroRestore or add a sync replica; page on fallbackDatabase platform
Split brainTwo nodes report role=primary; epoch mismatch errorsEvery writer on that shardFence the stale leader; freeze writes; reconcileDatabase platform + data owner
Abandoned CDC slotreplication.slot_retained_bytes growing; disk free < 20%Leader disk full → write outageDrop or advance the slot after owner sign-offConsumer team owns the slot; platform enforces limits
Lossy failoverreplication.position_gap_at_promotion > 0Lost acknowledged writes; possible duplicates and ID reuseReconcile from the old leader's logData owner signs off
Multi-leader conflict stormreplication.conflicts_resolved per minute 10× baselineSilently overwritten user dataPin hot keys to a home region; review merge rulesProduct team owning the data

The Numbers in Context#

NumberValueWhat It Means for Your Design
In-region async lag (healthy)~1–50msRead-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 lagRTT (60–150ms) + apply; seconds under loadRemote regions are always at least one RTT behind
Cross-AZ sync commit cost+0.5–2ms per commitAffordable for almost every OLTP write
Cross-region sync commit cost+60–150ms per commitOne synchronous continent hop per write is a product decision
Failure detection10–30s (3 missed heartbeats at 3–10s)Faster detection means more false failovers on network blips
In-region automated failover (RTO)~30s–2minBudget it in the SLO: a few failovers a year fits 99.95%
Async RPOlag × write rate: 2s × 5,000/s = ~10,000 writesSay the count, not the seconds, to the data owner
Quorum (3, 2, 2)Survives 1 replica loss for reads and writes3 copies across 3 AZs is the standard floor
Six copies, 4/6 write, 3/6 readSurvives an AZ plus one more nodeWhat it costs to tolerate AZ + 1 failures with quorums
Replica read capacityEach replica ≈ leader's read capacity; writes are replayed on every replicaReplicas scale reads, never writes
Retention for a lagging logical slotWrite rate × outage time: 50MB/s × 6h ≈ 1TBOne 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#

PatternHow It WorksWhen to Use
Chain replicationWrites enter the head, flow down a chain, ack from the tail; reads served by the tailStrong reads with simple failure handling, when chain latency is acceptable
Witness / tie-breaker replicaA cheap node in a third site votes in elections but stores little or no dataTwo data sites that still need a majority to avoid split brain
Delayed replicaA replica that deliberately applies the log 1–24h behindRecovery from a bad DELETE or a corrupting deploy faster than a backup restore
Cascading replicationReplicas replicate from other replicas, not the leaderMany replicas or a remote region without loading the leader's network
Per-transaction durabilityThe ledger commit waits for a sync replica, the clickstream insert does notMixed-value data in one database
Log-structured storage replicationShip only the redo log to a quorum of storage nodes and let storage build pagesCloud databases that separate compute from storage
Zero-downtime migration via logical replicationReplicate old → new, verify row counts and checksums, flip writes, keep reverse replication brieflyMajor-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.

OptionWhat WorksWhat BreaksWho Pays
Each team runs its ownAutonomy; teams tune for their workloadInconsistent RPO; untested failovers; automation that promotes across regions on a blipThe data owner, after the incident; customers whose writes were lost
Central DBA team approves every changeConsistent, reviewed settingsBottleneck; tickets take weeks; teams route around it with self-hosted storesProduct 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 classClassification arguments; exceptions process needed; platform becomes tier-0Platform 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.

ScaleFleetReplication SpendPlatform HeadcountOn-Call LoadThe Expensive Mistake
Startup5 databases, 1 region, 1 standby each~$5–10K/month (standbys double instance cost)0.5 engineerA failover a quarter, mostly managed by the cloudNo cross-region copy at all; region loss = restore from backup, hours of RTO
Growth80 databases, 2 regions, 2 in-region + 1 cross-region replica each~$250–400K/month; cross-region transfer ~$10–20K/month4–6 engineers (~$125K/month)Dedicated rotation; lag and slot incidents weeklyAsync replication for ledgers because "it's the default"
Large1,000+ databases, 3–5 regions~$4–8M/month; replicas are ~50–65% of database spend12–20 engineersTier-0 rotation; quarterly region-evacuation game daysOver-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#

Diagram: The 3-Year Evolution Path

One-Way Doors vs Two-Way Doors#

DecisionDoorReversal Cost
Sync vs async for one databaseTwo-wayConfig change and a restart; minutes
Read routing and lag capsTwo-wayRouter config
Choosing leaderless or multi-leader for a datasetOne-wayApplication logic now assumes conflicts and merges; going back means a migration and new invariants
Exposing a CDC stream to other teamsMostly one-wayThe change-event schema becomes a public contract; consumers pin it for years
Automatic cross-region promotionTwo-way to switch off, one-way once it fires badlyA single bad promotion can mean a day of reconciliation
Primary region placementOne-wayData gravity; every dependent service's latency is built around it
Conflict resolution rule (LWW vs merge)One-way per datasetData 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:

  1. 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).
  2. Tier 0–1: commit waits for a quorum or ANY 1 of ≥2 replicas in other AZs. RPO = 0 for single-AZ loss.
  3. Tier 0–2: an asynchronous copy in a second region with p99 lag ≤ 5s, alerted.
  4. 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.
  5. 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#

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

ConceptSenior (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 > N is 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 2 cross-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#

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

  1. Loading the index…