System designCore3 min

Database replication

Keep copies of the data on several machines, so you survive losing one and can serve more reads. The hard part is that the copies are never quite in step.

1 · Why copy

Three reasons to keep replicas

ReasonHow replicas help
Stay upIf one machine dies, another has the data and takes over
Scale readsSpread reads over several copies
Be near usersA copy in each region answers local reads quickly

2 · The common setup

One leader takes writes, followers copy it

ApplicationLeaderFollower 1Follower 2writeschange logchange logreadsApplicationLeaderFollower 1Follower 2writeschange logchange logreads

The leader records every change in a log and streams it to the followers, which apply it in the same order. Writes go to the leader; reads can go anywhere.

3 · The key choice

Wait for followers, or don't

Synchronous: the leader waits for the follower, so an acknowledged write is on two machines. Asynchronous: the leader answers straight after its own write, which is faster, but a write can be lost if the leader dies before sending it.

Many systems wait for one follower and let the rest catch up (semi-synchronous): no write is on a single machine, and one slow follower can't stall everything.

4 · The catch

Replication lag: reading your own old data

Asynchronous followers run behind, usually by milliseconds, sometimes by seconds under load. A user who updates their profile and reloads can be served by a follower that hasn't seen the change yet.

GuaranteeHow to get it
Read your own writesRead from the leader for a short while after a user writes, or until the follower has reached the write's position in the log
Monotonic reads (never go back in time)Pin each user to one follower
Up-to-date reads for everythingRead from the leader, or use quorum reads: at the cost of the read scaling you wanted

5 · Other setups

Several leaders, or none

SetupWrites go toGood atWatch out for
Single leaderOne nodeSimple; no write conflictsFailover; write capacity of one machine
Multi-leaderA leader per regionLocal writes in every regionTwo regions editing the same row: conflicts to resolve
Leaderless (Dynamo-style)Any N replicas; W must acknowledgeNo failover; stays writable when nodes failReads of R replicas; W + R > N for overlap
  1. Failover, step 1: notice.

    The leader stops answering health checks for a few seconds.

  2. Choose.

    The most up-to-date follower is promoted, by consensus or by an orchestrator.

  3. Redirect.

    Clients and the other followers switch to the new leader.

  4. Fence the old one.

    If the old leader comes back believing it's still in charge, two leaders accept writes (split brain). It must be shut out.

6 · In a real system

Three copies of every short link

URL Shortener System Design

Leaderless replication with quorums

The URL shortener keeps every link on three Cassandra replicas. A write is acknowledged by a majority in the local datacenter, and a read waits for the same majority, so a link that was just created never reads as missing, even with one replica down.

Check yourself

3 questions

1. With asynchronous replication, what can happen if the leader crashes?
2. A user edits their profile and immediately sees the old version. What's the cause?
3. Leaderless, N = 3 replicas. Which W and R make every read see the latest acknowledged write?

Takeaways

Remember this

  • Replicas buy availability, read capacity and nearness to users.
  • Synchronous is safer, asynchronous is faster; waiting for one follower is a common middle.
  • Followers lag: route read-your-writes traffic to the leader.
  • Failover must fence the old leader, or you get split brain.