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
| Reason | How replicas help |
|---|---|
| Stay up | If one machine dies, another has the data and takes over |
| Scale reads | Spread reads over several copies |
| Be near users | A copy in each region answers local reads quickly |
2 · The common setup
One leader takes writes, followers copy it
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.
| Guarantee | How to get it |
|---|---|
| Read your own writes | Read 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 everything | Read from the leader, or use quorum reads: at the cost of the read scaling you wanted |
5 · Other setups
Several leaders, or none
| Setup | Writes go to | Good at | Watch out for |
|---|---|---|---|
| Single leader | One node | Simple; no write conflicts | Failover; write capacity of one machine |
| Multi-leader | A leader per region | Local writes in every region | Two regions editing the same row: conflicts to resolve |
| Leaderless (Dynamo-style) | Any N replicas; W must acknowledge | No failover; stays writable when nodes fail | Reads of R replicas; W + R > N for overlap |
- Failover, step 1: notice.
The leader stops answering health checks for a few seconds.
- Choose.
The most up-to-date follower is promoted, by consensus or by an orchestrator.
- Redirect.
Clients and the other followers switch to the new leader.
- 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
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.