Database Replication
Leaders, followers, failover and the lag that bites: what copies of a database buy you, and what each kind of copy costs.
An interactive System Design lesson: 22 steps, about 32 minutes, on a live simulation in your browser.
A social site stores one row per user: user:1 is Ada, user:2 is Grace, and so on up to user:12. The client opens a profile page, the app server asks the database for the row, and the answer comes back in 29 ms: 10 ms each way to the app, 1 ms each way to the database, 2 ms of app work and 5 ms for the query.
Two more database machines stand on the stage, follower-1 and follower-2. For now the app ignores them: every read and every write goes to the machine called leader. As far as the application is concerned there is one copy of the data.
What you will learn
One copy of the data
- A profile page and one database: With one database, the whole site's reads, writes and data live and die with one machine.
- The database dies
Leader, followers and the log
- Every write goes into a log: Replication is one machine's ordered list of changes, replayed by every other machine in the same order.
- When is the write confirmed?: An asynchronous leader says yes as soon as it has the write. Followers get it afterwards, and the client is never told how long afterwards.
- Streaming replication in Postgres
Reading from followers
- Send the reads to the followers: Followers multiply read capacity. They do nothing for writes: every write still lands on the one leader and on every follower.
- Save, then reload twice: Asynchronous replication means a follower can return data older than a write the same user has already been told succeeded.
- Eventually consistent
- Read your own writes: Read-your-writes is a per-user promise. Keep it by reading that user's recent data from the leader, or from a follower that has provably caught up.
- Drill: how far behind is it?
Buying durability with latency
- Wait for one follower, or for all: Synchronous replication moves the cost of durability onto every write: the leader cannot answer until the follower it waits for has answered.
- A dead follower stops every write
- Drill: wait for any one standby
Failover
- Three saves, all confirmed
- Promote a follower: Failover with asynchronous replication loses every write the old leader confirmed but had not yet shipped: at most the replication lag's worth.
- The old leader comes back
- Two leaders: A deposed leader that does not know it was deposed keeps accepting writes. Two leaders means two diverging histories.
- Fence the old leader: Fencing makes the old leader unable to write before the new leader starts. Without it, every failover risks split brain.
- Drill: promote a standby
Across regions
- A follower an ocean away: A follower in another region buys survival of a regional outage and local reads there. It cannot be waited for on every write without paying the round trip on every write.
Recap & playground
- Cheat sheet
- Playground