The short answer
Quick answer: Replication means keeping copies of the same data on several database servers. In the most common setup, one server (the leader or primary) accepts all writes and sends a stream of changes to the others (followers or replicas), which apply them in the same order. Replicas provide high availability (if the leader fails, a replica takes over), read scaling (queries can be spread across replicas) and lower latency for distant users. The main complication is replication lag: replicas are usually slightly behind the leader, so a read from a replica can return stale data.
Why replicate
- Survive failures. Hardware dies. With a copy on another machine, the data and the service survive.
- Scale reads. Most applications read far more than they write. Replicas share the read load.
- Be closer to users. A replica in another region answers local reads quickly.
- Offload heavy work. Run backups and analytics on a replica without slowing the primary.
Replication does not scale writes. Every write still goes through the leader and must be applied on every replica. For write scaling you need sharding.
Leader-follower replication
This is the standard design for PostgreSQL, MySQL, SQL Server, MongoDB and many others.
- The leader executes a write and records it in its log.
- It streams that log to each follower.
- Each follower applies the changes in order.
- Clients send writes to the leader, and reads to either.
What gets shipped
| Method | What is sent | Notes |
|---|---|---|
| Statement-based | The SQL statement itself | Breaks on non-deterministic functions such as NOW() or RANDOM() |
| Physical (WAL shipping) | The low-level write-ahead log | Exact byte-for-byte copy; replicas must run a compatible version |
| Logical (row-based) | Row-level changes: "row 42 now has these values" | Works across versions and for feeding other systems (change data capture) |
The PostgreSQL high availability documentation compares these approaches.
Synchronous vs asynchronous
When does the leader tell the client "committed"?
| Asynchronous | Synchronous | |
|---|---|---|
| Leader confirms | As soon as it has written locally | After a replica confirms it has the data |
| Write latency | Low | Higher (adds a network round trip) |
| If the leader dies | Recently committed writes may be lost | No committed data is lost |
| If a replica is slow or down | No effect on writes | Writes stall |
A common compromise is semi-synchronous: require confirmation from at least one replica (or a majority), and let the rest follow asynchronously. This is durable without letting one slow replica stop all writes.
Replication lag and its surprises
With asynchronous replication, followers trail the leader by milliseconds, or by much more under load. This creates visible anomalies:
- Reading your own writes. A user saves their profile, the page reloads from a replica that has not caught up, and the old value appears.
- Going back in time. Two successive reads hit different replicas, and the second shows older data than the first.
- Out-of-order views. A reply appears before the message it answers.
Common fixes:
- Read from the leader for a short time after a user writes, or for data the user just changed.
- Pin a user's session to one replica.
- Track the log position of the user's last write and only read from replicas that have reached it.
- Monitor lag and remove badly lagging replicas from rotation.
This behaviour is a form of eventual consistency.
Failover
When the leader fails, a follower must take over. The steps:
- Detect the failure, usually by missed heartbeats over a timeout.
- Choose a new leader, ideally the replica with the most recent data.
- Promote it and reconfigure the other replicas to follow it.
- Redirect clients, via a proxy, DNS or a service registry.
Each step has pitfalls:
- Data loss. With asynchronous replication, the new leader may be missing the old leader's last writes.
- Split brain. If the old leader was only cut off by a network problem and still believes it is in charge, two leaders accept writes and the data diverges. Systems prevent this by requiring a majority of nodes to agree on a leader, and by fencing the old one so it cannot write.
- False alarms. Too short a timeout causes needless failovers during a brief slowdown; too long extends real outages.
Leader election is typically delegated to a consensus system. See how consensus algorithms like Raft work.
Multi-leader replication
Several nodes accept writes and replicate to each other. This is used for multi-region deployments and for apps that work offline and sync later.
The catch is write conflicts: two leaders change the same row at the same time. Strategies for resolving them:
- Last write wins, by timestamp. Simple, but silently discards data, and clocks are unreliable.
- Merge the values, where the data type allows.
- Conflict-free replicated data types (CRDTs), designed so concurrent changes always merge predictably.
- Ask the application or user to resolve it.
Multi-leader setups are powerful and notoriously tricky. Avoid them unless you need them.
Leaderless replication
In systems inspired by Amazon's Dynamo, such as Cassandra and Riak, any replica accepts writes. The client (or a coordinator) sends each write to all N replicas and waits for W acknowledgements. Reads query several replicas and wait for R responses.
If W + R > N, the sets of nodes written to and read from must overlap, so at least one response contains the latest value. For example, with N = 3, choosing W = 2 and R = 2 tolerates one node being down.
Replicas that missed writes are repaired in the background and during reads. These systems stay available during failures, in exchange for weaker consistency guarantees. The trade-off is the subject of the CAP theorem.
Replication vs backups
A replica is not a backup. If someone runs DELETE FROM users without a WHERE clause, that mistake replicates to every copy within milliseconds. You need both:
- Replication for hardware failure and availability.
- Backups (and point-in-time recovery) for human error, bugs and corruption.
Frequently asked questions
What is a read replica?
A copy of the database that receives changes from the primary and serves read-only queries, taking load off the primary.
What is replication lag?
The delay between a write being committed on the leader and being applied on a replica. During that time, the replica returns older data.
What is the difference between replication and sharding?
Replication keeps full copies of the same data on several servers. Sharding splits the data so each server holds a different part. They are often combined: each shard is replicated.
Does replication slow down writes?
Asynchronous replication adds almost no delay. Synchronous replication adds the time to reach a replica and get an acknowledgement.
Conclusion
Replication keeps a system running when machines fail and lets reads spread across many servers. The price is living with copies that are not always identical. Decide how much data you can afford to lose and how stale a read can be, and those two answers will choose your replication mode for you.
Related articles
- How Write-Ahead Logs Prevent Data Loss During Crashes
- Sharding Explained: How Databases Scale Beyond One Machine
- The CAP Theorem Explained With Real Examples
- What Is Eventual Consistency and When Is It Acceptable?
