System Design
X1

Replication

Keep copies of the same data on several machines. Copies survive a failed machine and serve more reads, but each copy is a little late.

Not startedSaved in this browser only.
  1. 1One leader takes writes and streams its log. Standbys replay it in order, a little late.
  2. 2Asynchronous replicas can lose the last commits on failover. A synchronous standby trades commit latency for zero loss.
  3. 3Replica reads can be stale. Fix read-your-writes with a token or a leader read after a write.
  4. 4Quorums need R + W > N, and a sloppy quorum breaks that rule.
X1
    A

    A stale read

    recorded from a Postgres standby
    Userleadertakes writesreplicaserves readsname = 'Grace'ok, committedreload: SELECT name'Ada' (old)WAL applied 300 ms later✕ The user does not see the change they just made
    • Replication copies every write from the leader to its replicas.
    • Copies give failover, read capacity and copies near users.
    • The cost: a copy lags, and a failover can lose what it never received.

    A replica is always a little behind the leader. A read there can miss a change the same user just made.

    B

    Replication methods

    which for which need
    methodon leader failurestatus
    One leader, asynchronous standbysThe last commits can be lost.Small loss accepted
    One leader, quorum commit: ANY 1 of 2 standbysNo commit lost if a synchronous standby survives.Approved
    One leader, every standby synchronousNo loss, but one slow standby stops all writes.Not approved
    Multi-leader, one leader per regionWrites go on; concurrent edits conflict.Multi-region writes
    Leaderless quorum (N, R, W)No failover step at all.R + W > N
    Replay SQL statements on each copynow() and random() differ per copy.Not approved

    Postgres streams its write-ahead log: the replica replays the exact bytes, not the statements.

    I start with one leader and asynchronous standbys, and I add a synchronous standby when losing the last commits on failover is not acceptable.

    C

    When a commit returns

    synchronous_commit, per transaction
    the commit returns at the marked pointleadersynchronous standbynetworkin leadermemoryflushed toleader diskstandby wroteit to its OSflushed tostandby diskreplayed:reads see itofflocalremote_writeonremote_apply
    levelleader restartstandby OS crashstandby readlocalhost
    offLoses 600 msMay loseStale0.08 ms
    localKeptMay loseStale0.14 ms
    remote_writeKeptMay loseStale0.35 ms
    onKeptKeptStale0.36 ms
    remote_applyKeptKeptFresh0.35 ms

    Guarantees from the Postgres 16 documentation, with a synchronous standby set. off can lose up to 3 × wal_writer_delay: 600 ms at the default 200 ms. Latency: median of 500 commits, primary and standby on one laptop; across regions add one network round trip.

    I choose the durability level per transaction. Payments wait for a synchronous standby to flush; logs only wait for the local disk.

    D

    The request path

    click a step; its path lights up
    Clientapp or browserService + routerknows the leaderLeadertakes all writesFailover managerleader lock in etcdStandby 1hot standby, readsStandby 2hot standby, reads

    Step 1: Write

    • Send every write to the leader.
    • The commit returns at the synchronous_commit level: local disk, or a standby too.

    If it fails

    The leader is down. Writes fail until failover promotes a standby (step 5).

    Writes go to the leader. Reads go to standbys, except a read that must see the user's own write, which waits for a token or goes to the leader.

    E

    Capabilities used

    what each tool gives you
    toolcapabilitywhat it gives this designalso used for
    PostgresStreaming replication of the WALEach standby is an exact copy, replayed in commit order.Point-in-time recovery from archived WAL
    PostgresHot standbyRead-only queries on a standby while it replays.Reports off the primary
    Postgressynchronous_standby_names: FIRST k or ANY kA commit waits for k standbys. ANY 1 of 2 survives one slow standby.Zero-loss failover
    Postgressynchronous_commit, settable per transactionDurability per write: payments wait for a standby, logs do not.Bulk loads with off
    Postgrespg_current_wal_insert_lsn(), pg_last_wal_replay_lsn()A version token: the replica knows if it has a given commit.Causal reads across services
    Postgrespg_stat_replication: write_lag, flush_lag, replay_lagLag per standby. The router skips a standby that lags too much.Alerts
    Postgresrecovery_min_apply_delayA standby that stays behind on purpose. The lab used it to make lag visible.Undo a bad DELETE
    PostgresReplication slotsCap the size The leader keeps WAL until the standby has it. A dead standby fills the disk.Logical decoding
    Postgrespg_promote()Turns a standby into the leader.Planned switchover
    RedisAsynchronous replicas; WAIT numreplicas timeoutLimit WAIT counts acknowledgements but does not make Redis strongly consistent.Read scaling for caches
    CassandraConsistency level per query: ONE, QUORUM, LOCAL_QUORUM, ALLR and W chosen per read and write. Hinted handoff and read repair heal copies.Multi-region writes
    DynamoDBConsistentRead on a readChoose a strongly consistent read; the default read can be stale.Serverless key-value
    Patroni · etcdA leader key with a time to liveOne leader at a time. The holder renews it; a new leader is promoted when it expires.Any single-leader failover

    Postgres gives me ordered WAL streaming, per-transaction durability, and LSN functions for read-your-writes. Leaderless stores give me per-query consistency levels instead.

    F

    Quorum reads and writes

    pseudo code
    leaderless quorumpseudo code
    write(value, W):
      send value to all N replicas
      wait for W acks1                         // the write set
    
    read(R):
      ask all N replicas
      wait for R replies2                      // the read set
      RETURN the reply with the newest version3
    
    R + W > N4  =>  every read set meets every write set
    1. 1The write succeeds after W replicas store it. The others get it later, or never if they are down.
    2. 2Any R replicas. The read cannot know which ones have the write.
    3. 3Each value carries a version. One fresh reply is enough.
    4. 4R + W > N means the sets must overlap: there are only N replicas to choose from.
    Tested source Go: every case, enumerated
    Go: every case, enumeratedgo
    // Outcome is what a write with W acknowledgements, then a read with R replies, can return,
    // over every choice of which W nodes acknowledge and which R replicas reply.
    type Outcome struct {
      WriteOK bool // W nodes were reachable
      ReadOK  bool // R replicas were reachable
      Cases   int  // pairs (write set, read set)
      Stale   int  // pairs where no reply carries the write: the read returns the old value
    }
    
    // Enumerate tries every write set and every read set. A read sees the write when one of its R
    // replies comes from a node in the write set; it returns the newest version among the replies.
    func Enumerate(n, r, w int, p Pattern) Outcome {
      wt, rt := targets(n, p)
      o := Outcome{WriteOK: len(wt) >= w, ReadOK: len(rt) >= r}
      if !o.WriteOK || !o.ReadOK {
        return o
      }
      for _, acks := range subsets(wt, w) {
        for _, reads := range subsets(rt, r) {
          o.Cases++
          if !overlaps(acks, reads) {
            o.Stale++
          }
        }
      }
      return o
    }
    

    A write waits for W acknowledgements and a read for R replies. With R plus W greater than N, the two sets always share a replica, and the read takes the newest version.

    G

    Why the sets overlap

    N = 3
    N = 3, W = 2, R = 2write: 2 acksv2r0v2r1v1r2read: 2 replies✓ the read sees v2R + W = 4 > 3: the sets share replica 1N = 3, W = 1, R = 1write: 1 ackv2r0v1r1v1r2read: 1 reply✕ the read returns v1R + W = 2 ≤ 3: the read can miss
    N, W, Rwrites survivereads survivereads fresh
    3, 2, 21 down1 downApproved
    3, 3, 10 down2 downApproved
    3, 1, 32 down0 downApproved
    3, 1, 12 down2 downStale in 6 of 9
    5, 3, 32 down2 downApproved

    Writes survive N − W failed replicas; reads survive N − R.

    N = 3 with W = 2 and R = 2 survives one replica down for both reads and writes, and every read is fresh.

    H

    Try it: the quorum simulator

    recorded from the lab simulation
    N replicas
    W acks for a write
    R replies for a read
    what goes wrong

    All replicas are up. The read arrives before the replicas outside the write set apply it.

    R + W = 4 > N = 3Every read set meets every write set.
    r0v2ack
    r1v1reply
    r2v2ackreply
    This read sees v2: one reply comes from a node in the write set.
    run 1 of 6 · seeded

    Over all 9 choices of write set and read set, the read returns the old value in 0 (0%).

    Every N from 1 to 5, every R and W, four patterns: 220 settings. Each enumerates all write sets and read sets.

    I check two things: can each operation still gather enough replicas, and does every read set meet every write set.

    I

    Lag anomalies and their fixes

    what the user sees
    time, ms →partition 1replica lags 500 mspartition 2replica lags 5 msreaderreads replicas0100200300400500"How long is the wait?"applied on replica"About 10 seconds." (applied at 25 ms)at 100 ms: the answer, with no question✕ The reader sees an answer to a question that does not exist yet
    anomalythe user seesfixcost
    Read your writesTheir own change is missing after a reload.Read the leader after a write; version token; remote_applyLeader load; a short wait; slower commits
    Monotonic readsA comment appears, then disappears.Sticky replica per user; token of the newest position seenUneven replica load; a short wait
    Consistent prefixAn answer before its question.Write related data to one partition, so one log orders itThe partition key choice

    Replica lag breaks three promises: users see their own writes, data never goes backwards, and an answer never appears before its question.

    J

    Try it: lag on a timeline

    recorded from the lab simulation
    scenario
    replica lag
    where reads go

    The user renames their profile, then reloads it 10 ms later. Each read goes to any replica, with no check.

    timeclientleaderreplica 10 msUPDATE name = 'Grace'1 mscommit at 0/5A02 msok11 msSELECT name12 ms'Ada' (old)51 msapply up to 0/5A0
    Stale: the user does not see their own change.

    Same case on a real Postgres standby that applies each commit 300 ms late: a read right after the write returned 'Ada'. A read that first waited for the token returned 'Grace'.

    For read-your-writes I give the client its commit position. The replica waits until it has replayed that far, or the router reads the leader.

    K

    Read your writes with a token

    pseudo code, tested on Postgres
    version tokenpseudo code
    write(name):                  // leader
      UPDATE profiles SET name = ...
      token = WAL position        // after the commit1
      RETURN token
    
    read(token):                  // replica
      WHILE replay position < token2:
        wait 2 ms                 // until a deadline3
      SELECT name ...
    1. 1Read the position in a second statement. A position read inside the transaction comes before the commit record.
    2. 2pg_last_wal_replay_lsn() on the replica, compared as pg_lsn.
    3. 3A deadline caps the wait when a replica is far behind. Then read the leader.
    Tested source Go: write, then read after the token · SQL: token and replay check
    Go: write, then read after the tokengo
    // WriteName changes a name on the leader and returns a version token: the WAL position just past
    // the commit. It reads the position after the commit, in a second statement; a position read
    // inside the transaction comes before the commit record.
    func WriteName(ctx context.Context, leader *pgx.Conn, id int, name string) (string, error) {
      if _, err := leader.Exec(ctx, token["write_name"], id, name); err != nil {
        return "", fmt.Errorf("write name %d: %w", id, err)
      }
      var lsn string
      if err := leader.QueryRow(ctx, token["token"]).Scan(&lsn); err != nil {
        return "", fmt.Errorf("read commit position: %w", err)
      }
      return lsn, nil
    }
    
    // ReadNameAfter reads a name from a replica once the replica has replayed the WAL up to the token,
    // so the caller sees its own write. It gives up at the context deadline; the caller can then read
    // from the leader.
    func ReadNameAfter(ctx context.Context, replica *pgx.Conn, id int, lsn string) (string, error) {
      for {
        var caught bool
        if err := replica.QueryRow(ctx, token["replayed"], lsn).Scan(&caught); err != nil {
          return "", fmt.Errorf("check replay position: %w", err)
        }
        if caught {
          break
        }
        select {
        case <-ctx.Done():
          return "", fmt.Errorf("replica has not replayed %s: %w", lsn, ctx.Err())
        case <-time.After(2 * time.Millisecond):
        }
      }
      return ReadName(ctx, replica, id)
    }
    
    SQL: token and replay checksql
    UPDATE profiles SET name = $2 WHERE id = $1;
    
    -- After the commit: the WAL position just past the commit record is the version token.
    SELECT pg_current_wal_insert_lsn()::text;
    
    -- On the replica: has it replayed the WAL up to the token?
    SELECT pg_last_wal_replay_lsn() >= $1::pg_lsn;
    right after
    'Ada': the standby is 300 ms behind.
    with token
    'Grace', after a wait of about 300 ms.
    remote_apply
    'Hopper': the commit itself waited.

    The commit position is a version token. A replica that has replayed past it has the write.

    L

    Failover and split brain

    promote, then fence
    time →old leaderreplicastorage, clientsterm 1paused (GC, network)promoted, term 2heartbeats missedwakes: write, term 1write, term 2✕✓Every write carries the leader's term. Storage keeps the highest term seenand refuses anything lower, so the old leader cannot write after it wakes.
    1. Detect: the leader stops renewing its key. Patroni gives the key a 30 s time to live by default.
    2. Choose: the standby with the most replayed WAL loses the least.
    3. Promote it and raise the term in the consistent store.
    4. Fence: cut the old leader off, or make storage refuse its term.
    5. Point clients at the new leader: DNS, a proxy, or the router's leader key.

    I promote the standby with the most WAL, raise the term, and fence the old leader so its late writes are refused.

    M

    More than one writer

    multi-leader and leaderless
    conflict rulestatusresult
    One home region per recordApprovedNo conflicts: each record has one writer.
    Merge with a CRDT: counters, sets, textMergeable dataEvery edit survives, merged the same way on every copy.
    Keep both versions, ask the app or userRare conflictsCorrect, but the app must handle siblings.
    Last write wins by timestampLoss acceptedDrops one concurrent edit. Clock skew picks the winner.
    • Multi-leader: each region has a leader and accepts writes locally.
    • Leaderless: any replica takes a write; quorums and read repair keep copies close.
    • Both need a rule for two writes to the same key at the same time.

    I avoid multi-leader writes unless users in several regions must write locally. Then I give each record a home region or merge conflicts with a CRDT.

    N

    Failure cases

    what breaks, and why it stays correct
    eventresultwhy it is safesaved by
    The leader crashes; standbys are asynchronousThe last commits may be missing.Promote the standby with the most WAL. With a synchronous standby, nothing acknowledged is lost.Sync standby
    A network split leaves the old leader runningTwo nodes think they lead.The leader key expires first; the old leader's term is refused.etcd lease
    The only synchronous standby diesCommits that wait for it block.ANY 1 of 2 standbys: the other one acknowledges.Quorum commit
    A standby lags 30 secondsIts reads are old.The router reads replay_lag and skips that standby.Lag check
    A user reads right after their writeA replica may not have it.The token makes the replica wait, or the read goes to the leader.Token
    A dead standby keeps its slotWAL piles up on the leader.max_slot_wal_keep_size caps it; an alert fires on slot lag.Slot limit
    A quorum replica is down during a writeA stand-in keeps a hint.The hint is handed over later. Reads use R + W > N on home replicas.Hinted handoff
    O

    Scale ladder

    start simple; climb only on a signal
    Each step adds one component1Postgres2+ standby3+ read replicas4+ sync standby5+ other regionmore load →
    Read capacity against demand1k10k100k1MDemand, average: 11,574 primary-key reads per secondDemand, average11,574Demand, 5× peak: 57,870 primary-key reads per secondDemand, 5× peak57,870One Postgres: 107,000 primary-key reads per secondOne Postgres107,000+ 2 replicas: 321,000 primary-key reads per second+ 2 replicas321,000+ 5 replicas: 642,000 primary-key reads per second+ 5 replicas642,000primary-key reads per second, log scale
    stepaddit handlesmove up when you see
    1One Postgres, with WAL archived for backups.About 107,000 primary-key reads a second in the lab. Restore from backup after a loss.A restore takes too long for your recovery time.
    2One asynchronous standby and a failover manager.Failover within the leader key time to live. Loses only the commits inside the lag: 0.13 ms in the lab.Reads load the primary.
    3Read replicas behind a router with lag checks and tokens.Reads grow with each replica. Writes do not.No acknowledged commit may be lost.
    4A synchronous standby: ANY 1 of 2.Zero loss while one synchronous standby survives. Each commit waits one round trip.The whole region can fail, or users are far away.
    5A standby in another region, asynchronous.Survives a region loss. A synchronous one there adds about 40 ms per commit at 4,000 km.Writes outgrow one leader: shard, as on the sharding sheet.

    Demand example: 1 billion reads a day is 1,000,000,000 / 86,400 ≈ 11,574 a second. Replica capacity is the lab measurement times the copy count. Light in fibre covers about 200 km per ms, so 4,000 km is 20 ms each way.

    I start with one Postgres and archived WAL. Then I add a standby for failover, replicas for reads, and a synchronous standby when no commit may be lost.

    P

    Drill

    predict, then reveal

    0 of 10 known

    1. A replica lags 50 ms. The user saves a change and reloads 10 ms later. What do they see, and how do you fix it?

    2. Why does a sticky replica fix monotonic reads but not read-your-writes?

    3. N = 3, W = 2, R = 2, and one replica is down. Do reads and writes work? Are reads fresh?

    4. Why can a sloppy quorum return an old value even with R + W > N?

    5. One synchronous standby is configured and it dies. What happens to commits?

    6. You fail over to an asynchronous standby. What is lost?

    7. The old leader wakes after a long pause and accepts writes. How do you stop it?

    8. remote_apply makes reads on the standby fresh. Why not use it for every commit?

    9. Two regions accept writes, and two users change the same field at the same time. What happens?

    10. A standby dies, and its replication slot stays. What breaks next?

    Q

    Numbers to say

    measured, derived or cited
    lag, local
    A commit was visible on a standby on the same machine after 0.13 ms (median).
    sync commit
    0.14 ms local, 0.36 ms with a synchronous standby on localhost: one more round trip.
    far standby
    4,000 km is 20 ms each way in fibre: about 40 ms more per synchronous commit.
    off
    Can lose up to 600 ms of commits: 3 × the default wal_writer_delay.
    quorum
    N = 3, W = 2, R = 2 survives 1 replica down, with fresh reads.
    R = W = 1
    Of N = 3: stale in 6 of 9 cases when the read beats replication.
    reads
    About 107,000 primary-key reads a second on one Postgres.

    Postgres 16 on an 8-core laptop; reads with 32 clients; commit latency over 500 commits. Simulations are seeded and repeat exactly.