Replication
Keeping copies of data on several machines for durability, read capacity and locality, and living with copies that briefly disagree.
Distribution
Learn it
Replication keeps copies of a database on several machines. Usually one primary accepts all writes and streams each change to replicas, which apply the changes in the same order. Replicas can then serve reads.
Check
Which of these does adding read replicas not give you?The primary can send changes in two ways:
- Synchronous: it waits for a replica to confirm a change before telling the client the write succeeded.
- Asynchronous: it tells the client right away and sends the change afterwards. This is the usual default, because writes stay fast.
Check
With asynchronous replication, the primary confirms a write to a client and then crashes before sending it to any replica. You promote a replica. What happened to the write?With asynchronous replication, replicas are always slightly behind the primary. This is replication lag: usually milliseconds, but it can grow to seconds or more under heavy load.
Think first
A user edits their profile. The write goes to the primary, then the page reloads and reads from a replica. What might they see?Check
A user refreshes a page twice. The first request goes to a replica 10 ms behind; the second to a replica 2 seconds behind. A comment was posted one second ago. What do they see?
Quick reference
The same ideas, condensed for revision.
How it goes wrong
- Stale read after write
- The user's own change is missing on the next page.
- Data loss on failover
- Asynchronously replicated writes committed just before the crash are missing on the new primary.
- Lag spiral
- A replica falls behind under load, serves increasingly stale data, and may need to be removed from rotation.
- Caching a miss
- A 404 served from a lagging replica is cached and outlives the lag by hours.
Instead, consider
- Caching
- Reads repeat the same keys; a cache gives read capacity without a full copy of the database.
- Partitioning
- Writes or data size exceed one machine; replication alone does not scale writes.
In practice
- Postgres streaming replication
- Async by default; synchronous_commit for synchronous replicas.
- MySQL replication / Aurora replicas
- Leader-follower with replica lag metrics.
- DynamoDB global tables, Cassandra
- Multi-region, multi-writer with last-writer-wins conflict handling.
It assumes
- Reads dominate, or are latency-sensitive in places far from the primary.
- The application can tolerate, or explicitly works around, replication lag.
Explain it in your own words
Where you practise it
Further reading
Engineers describing it in systems they run.
- What we learned from a 22-Day storage bug (and how we fixed it)
Mux · Drew Rodman and Constantin Britcov · Post, Mar 2026
A candid incident report. Three small races in a segment store, exposed when a scaling change slowed object storage, and why it took weeks to connect the symptoms to the cause.
- The Great Re-shard: adding Postgres capacity (again) with zero downtime
Notion · Arka Ganguli and others · Post, Jul 2023
The payoff of fixed logical shards: growing to more machines by moving shards whole.
- Scaling Memcache at Facebook
Meta (Facebook) · Rajesh Nishtala and others · Paper, Apr 2013
The reference on running a look-aside cache hard: leases, invalidation from the commit log, failover without hammering the database, and consistency across regions.
Related concepts
- Durability
What has to have happened before a system may say "saved": which failures the data must survive, and where that guarantee is actually made.
- Caching
Keeping a copy of data closer to where it is used, trading freshness and complexity for speed and reduced load on the source.
- Partitioning
Splitting data or work by key so each part is handled independently: scaling out, and giving each key a single owner.
- Conflict resolution and convergence
When replicas accept concurrent changes, a deterministic rule must merge them so every replica ends in the same state without losing intent.
- Online data migrations
Moving live data to a new schema or store without downtime: write to both, backfill the past, verify, switch reads, then switch writes, with a way back at every step.
- Consistent hashing
Mapping keys to nodes so that adding or removing a node moves only a small share of keys, instead of reshuffling almost all of them.