Database Replication
Database Replication
Definition: The process of copying and maintaining the same data across multiple nodes to improve availability, fault tolerance, and read scalability.
How It Works
- Leader-follower (primary-replica): one node (the leader) accepts all writes, appends them to a write-ahead log, and streams that log to one or more followers, which apply the same changes in order. Reads can be served by the leader or, for eventually-consistent reads, by followers.
- Synchronous replication: the leader waits for an acknowledgment from one or more replicas before confirming the write to the client. Stronger durability (a follower already has the data if the leader dies) at the cost of added write latency, bounded by the slowest required replica.
- Asynchronous replication: the leader confirms the write immediately and streams it to replicas afterward. Lower, more predictable write latency, but a leader crash before replication completes means the acknowledged write can be lost.
- Semi-synchronous replication: a middle ground, the leader waits for acknowledgment from at least one replica (not all) before confirming, trading some latency for a bounded worst-case data loss window.
- Multi-leader replication: more than one node accepts writes, and each propagates its writes to the others. Useful across multiple data centers so each region gets local write latency, at the cost of needing conflict resolution when the same record is written in two places concurrently.
- Leaderless replication (Dynamo-style): any replica can accept a write; the client (or a coordinator) writes to W replicas and reads from R replicas, using quorum overlap (
W + R > N) to make reads likely to see the latest write, with conflicts resolved via vector clocks, version vectors, or last-write-wins. - Replication topology varies too: single-leader is simplest to reason about; multi-leader and leaderless scale writes further but push more complexity into conflict handling.
- Failover: when a leader fails, a follower is promoted (automatically via a consensus-based control plane, or manually). Correct failover must ensure the old leader can’t keep accepting writes after a new leader is chosen, or the cluster ends up with two leaders (split brain).
- Under the hood, log shipping: most leader-follower systems (Postgres, MySQL) replicate by streaming the write-ahead log (WAL), not by re-running SQL. Every change is assigned a monotonically increasing log sequence number (LSN); a follower’s replication lag is literally
leader's current LSN - follower's last applied LSN, which is why it’s usually reported in bytes rather than a wall-clock estimate. - A follower connects to the leader, requests the WAL starting from its last applied LSN, and the leader streams every subsequent record. On restart after a disconnect, the follower resumes from its last known LSN rather than re-copying the whole dataset, as long as the leader has retained WAL back that far (older WAL segments get recycled, which is why a follower disconnected too long needs a full resync instead).
Trade-offs
- Synchronous replication buys durability guarantees (an acknowledged write survives a leader failure) at the direct cost of write latency, since every write now waits on a network round trip to at least one other node.
- Asynchronous replication buys low, predictable write latency at the cost of a real, non-zero data-loss window if the leader fails before followers catch up.
- Read replicas buy read throughput and geographic read latency improvements at the cost of replication lag: a replica’s data is only as fresh as its last applied log entry, which is always at least slightly behind the leader.
- Multi-leader and leaderless designs buy write availability and low write latency in every region at the cost of needing an explicit conflict resolution strategy, something single-leader systems never have to build.
- More replicas generally mean more fault tolerance and more read capacity, but also more storage cost, more replication network traffic, and a larger quorum to coordinate for synchronous or leaderless writes.
Worked Example: Quantifying Replication Lag Risk
A leader accepts writes at 2,000/second, each roughly 1KB, producing about 2MB/second of WAL. An asynchronous replica on a network path with 100ms latency and enough bandwidth to keep up under normal load falls behind by roughly the network delay, around 200KB, or 0.1 seconds worth of writes, at steady state. If that replica’s disk I/O briefly can’t keep up (say, during a backup job), the lag can grow to seconds or minutes instead, at 2MB/second, a 30-second lag represents roughly 60MB, or about 60,000 unreplicated writes sitting only on the leader. If the leader fails during that window, promoting this replica loses all of them, a concrete way to reason about the data-loss window an asynchronous setup is actually exposed to, rather than treating “asynchronous replication risks data loss” as an abstract warning.
Why It Matters
- Read replicas offload read traffic from the primary, which is essential for read-heavy workloads (dashboards, search, content feeds) where reads outnumber writes by orders of magnitude.
- It provides fault tolerance: if the primary fails, a healthy replica can be promoted, avoiding a full outage from a single machine failure.
- It’s a prerequisite for multi-region deployments, since serving users from a nearby replica (rather than round-tripping to a single primary on another continent) is what makes global-latency SLAs achievable.
- It underlies most database backup and disaster-recovery strategies: a replica in a separate availability zone or region is a live, continuously updated backup, not just a periodic snapshot.
Common Pitfalls
- Replication lag causing the classic read-your-writes problem: a user writes data via the leader, immediately reads from a lagging replica, and doesn’t see their own change.
- Assuming synchronous replication is free. It directly adds round-trip network time to every write, and if a required replica is slow or unreachable, it can stall writes entirely unless the system degrades gracefully.
- Failing to handle split-brain during failover: if the old leader doesn’t step down cleanly (e.g., it’s just slow, not dead) and a new leader is promoted, both can accept writes simultaneously, corrupting data.
- Using multi-leader replication without a real conflict resolution plan, so concurrent writes to the same record in different regions silently produce undefined or arbitrary results.
- Treating a replica as a substitute for a backup. Replication propagates mistakes (a bad
DELETE, corrupted data) just as fast as it propagates good writes; only point-in-time backups protect against that. - Not monitoring replication lag as a first-class metric, so a slowly degrading replica goes unnoticed until a failover promotes a replica that’s minutes behind.
- Promoting whichever replica responds first during failover instead of the most caught-up one, risking a promotion that loses more acknowledged writes than necessary.
- Assuming a “healthy” replication connection (no errors in logs) means low lag. A replica can be connected and error-free while still falling behind purely because it can’t keep up with the leader’s write rate.
Comparison
| Synchronous Replication | Asynchronous Replication | Multi-Leader / Leaderless | |
|---|---|---|---|
| Write latency | Higher, waits on replica ACK | Lower, confirms immediately | Low, local write accepted |
| Data loss on leader failure | None (for acknowledged writes) | Possible, unreplicated writes lost | Depends on quorum settings |
| Conflict resolution needed | No, single writer | No, single writer | Yes, concurrent writers |
| Best fit | Strong durability requirements | High write throughput, latency-sensitive | Multi-region write availability |
| Topology | Write scalability | Complexity | Typical use |
|---|---|---|---|
| Single-leader | Limited to one node’s capacity | Low, easiest to reason about | Most OLTP relational databases |
| Multi-leader | Scales with number of leaders | Medium, needs conflict resolution | Multi-datacenter deployments needing local write latency |
| Leaderless (quorum) | Scales with number of nodes | High, needs read-repair/anti-entropy | High write-availability systems (Dynamo-style) |
Real-World Scenario
An e-commerce platform runs a Postgres primary with two asynchronous read replicas serving product catalog and dashboard queries. During a flash sale, a customer places an order, and the order confirmation page immediately queries a read replica that hasn’t caught up yet, briefly showing “order not found.” The team fixes this by routing any read immediately following a write from the same session to the primary (or a replica confirmed caught up via log position), while leaving general browsing traffic on the replicas, keeping the primary free to handle the actual write load.
Debugging Walkthrough: Split-Brain During a Failover
- The primary Postgres node in one availability zone starts experiencing high GC-style pause times under load (or, in a managed setup, the health-check agent has a bug), causing it to miss heartbeats to the failover controller without actually being down.
- The failover controller, seeing missed heartbeats, promotes a replica in another availability zone to primary and updates DNS/connection routing to point application traffic there.
- The original primary’s pauses end a few seconds later. It never received the “you are demoted” signal, because the network hiccup that caused missed heartbeats also delayed that message. It resumes accepting writes from any client still holding a connection to its old address, believing itself to still be primary.
- For a window of time, both the old and new primary accept writes independently. Application servers that already re-resolved DNS write to the new primary; a handful of connections still pinned to the old primary’s IP (from before the failover, via a connection pool that hadn’t refreshed) keep writing there.
- On-call notices the problem not from an availability alert, since both nodes are technically “up”, but from a data-integrity alert: a uniqueness constraint violation once the two primaries’ replication streams are later reconciled, or a customer report of an order that “disappeared.”
- Root cause identified: the failover process lacked a fencing step, no mechanism forcibly cut off the old primary’s ability to accept writes (revoking its database credentials, blocking it at the network level, or using a STONITH-style “shoot the other node in the head” action) before promoting the new one.
- Fix: add explicit fencing to the failover procedure, so a demoted primary is made unable to accept writes (not just asked nicely to stop) before a replica is promoted, closing the exact race that produced two simultaneous writers.
FAQ
Does more replicas always mean more availability? Only up to a point. More replicas add fault tolerance, but each one also adds coordination overhead for synchronous or quorum-based writes, and eventually the marginal replica adds more operational cost than reliability.
What’s the difference between replication and sharding? Replication makes full copies of the same data for redundancy and read scale; sharding splits the data into disjoint pieces across nodes for write and storage scale. They’re usually combined.
Can a replica serve writes? Only in multi-leader or leaderless topologies. In single-leader replication, followers are read-only, and any write sent to a follower is rejected or forwarded to the leader.
Does replication help with disaster recovery on its own? Only partially. A replica in another region protects against a regional outage, but it replicates mistakes too (bad deletes, corruption), so DR strategy still needs point-in-time backups alongside replication, not instead of it.
How is replication lag usually measured? As the difference between the leader’s current log position (or timestamp) and the position the replica has applied, often reported in bytes of unapplied log or in seconds behind.
History
- Early relational databases (1990s MySQL, Postgres) implemented single-leader, mostly asynchronous, statement-based replication primarily for read scaling and basic disaster recovery, not for automated failover.
- The 2000s web-scale era pushed replication toward log-based (row/WAL-based) replication for correctness, since statement-based replication could produce different results on different replicas for non-deterministic statements (e.g., ones using
NOW()orRAND()). - Amazon’s 2007 Dynamo paper popularized leaderless, quorum-based replication (
W + R > N) for high write-availability systems, directly influencing Cassandra and Riak. - Consensus-based replication (Raft, Paxos) later became the standard for systems requiring strict correctness on failover, replacing manual or heuristic-based leader promotion with a provably safe protocol.
Common Interview Questions
- Why does asynchronous replication risk data loss but synchronous doesn’t? Because an async leader confirms a write before any replica has it, so a crash immediately after acknowledgment loses that write; sync replication waits for the replica first.
- What causes replication lag, concretely? Network delay, a replica applying changes slower than the leader produces them (e.g., due to single-threaded log application), or a replica busy serving heavy read traffic.
- How do you avoid split-brain during a leader failover? Use a consensus-based control plane (or fencing tokens) to guarantee only one node believes it’s the leader at a time, and make the old leader step down or get fenced off before a new one is promoted.
- What’s the difference between statement-based, row-based, and log-based (WAL) replication? Statement-based replicates the SQL statement itself (can behave non-deterministically); row-based replicates the actual changed rows; log-based streams the underlying write-ahead log, which is what most modern systems use for reliability.
- When would you choose multi-leader replication despite the conflict-handling cost? When write latency in every region matters more than avoiding conflicts, and conflicts are rare or mechanically resolvable (e.g., CRDTs, per-field last-write-wins).
Example
A Postgres primary handles all writes; two asynchronous read replicas serve SELECT-heavy dashboard queries, reducing load on the primary. MySQL Group Replication and Postgres synchronous replication are used where durability on failover matters more than write latency, such as financial transaction logs.
Design Checklist
- Can the application tolerate replication lag for this specific read, or does it need to read its own writes?
- What’s the acceptable data-loss window if the leader dies mid-write: zero (synchronous), or a few seconds of the most recent writes (asynchronous)?
- If using multi-leader or leaderless replication, is there an actual, tested conflict-resolution strategy, not just an assumption that conflicts “won’t happen in practice”?
- Is replication lag monitored and alerted on as its own metric, separate from general database health?
Related Terms
Referenced by