Database Replication
Keeping multiple copies of the same database in sync, so no single database is a single point of failure.
What is it?
A single database server is a risk: if it goes down, the whole application loses access to its data. Replication means continuously copying data from one database (the primary) to one or more additional databases (replicas), so there's always more than one copy of the data available.
Replicas are also useful even when nothing has failed — since reads (fetching data) usually vastly outnumber writes (changing data), spreading reads across several replicas can handle far more traffic than one database ever could alone.
Explain like I'm 10
It's like a company keeping backup copies of an important physical ledger in multiple offices, updated continuously as changes come in. If the main office burns down, another office already has an up-to-date copy — and in the meantime, staff in every office can read from their local copy instead of everyone calling the main office.
Examples
Reads from a replica, writes to the primary
// Simplified pattern
async function getUser(id) {
return replicaDb.query("SELECT * FROM users WHERE id = ?", [id]);
}
async function createUser(data) {
return primaryDb.query("INSERT INTO users ...", [data]);
// this change then gets replicated to replicaDb automatically
}How it works
The primary database records every change it makes (often in a log), and continuously streams that log of changes to each replica, which applies the same changes in the same order to stay in sync. Because this streaming takes a small amount of time, a replica's data is typically slightly behind the primary's — this delay is called replication lag.
Why does it exist?
Replication protects against losing a single database entirely (a replica can be promoted to take over), and lets read-heavy applications scale far past what one database server could handle alone by spreading reads across many replicas.
When to use it
Add replication once a single database is either a reliability risk you can't accept, or a read bottleneck — most real production databases run with at least one replica for exactly these reasons.
When not to use it
For a small, early-stage application with low traffic and where brief downtime is acceptable, a single database is simpler to operate and reason about — replication adds real operational complexity (replication lag, failover logic) that isn't worth it until the risk or the read load actually justifies it.
Common mistakes
Reading data immediately after writing it from a replica, and being surprised the just-written data isn't there yet — that's replication lag.
Forgetting that replicas are (usually) read-only, and mistakenly sending writes to one.
Assuming replication alone is a backup strategy — a mistaken delete on the primary replicates to every replica too, just as fast as a legitimate change.
Practice exercises
- Easy:
Explain, in your own words, why spreading reads across replicas helps an application handle more traffic.
- Medium:
Describe what 'replication lag' is and a real scenario where it could cause a confusing bug.
- Hard:
Explain why replication is not a substitute for backups, even though both involve keeping copies of your data.
Interview questions
What is database replication?
Continuously copying data from a primary database to one or more replica databases, so there's more than one up-to-date copy available.
Why do replicas help with read scalability?
Because read traffic can be spread across many replicas instead of all hitting a single database, while writes still go to the primary.
What is replication lag?
The small delay between a change being made on the primary and that same change appearing on a replica, since replication isn't instantaneous.
What's the difference between synchronous and asynchronous replication?
In synchronous replication, the primary waits for at least one replica to confirm it received the write before telling the client the write succeeded — guaranteeing that confirmed data exists in more than one place. In asynchronous replication, the primary confirms the write immediately and streams it to replicas afterward, which is faster but means a just-confirmed write could still exist only on the primary for a brief window.
What's the tradeoff of synchronous replication — why doesn't every system just always use it, since it avoids lag?
Every write now has to wait on a round trip to at least one replica before it's acknowledged, adding real latency to every write, and if that replica is slow or unreachable, writes can stall or fail entirely — trading raw write speed and availability for the stronger guarantee that a confirmed write already has a second copy.
What is "semi-synchronous" replication, and what problem does it try to solve?
The primary waits for acknowledgment from at least one replica (not all of them) before confirming a write, then replicates to the rest asynchronously — a middle ground that guarantees a confirmed write survives the primary failing alone, without paying the latency and availability cost of waiting on every replica.
Scenario: a primary crashes with writes that had been confirmed to the client but hadn't yet reached any replica. What happens to that data?
Under asynchronous replication, that data is simply lost — the client was told the write succeeded, but no replica has it, and there's nothing left to promote with it intact. This is precisely the failure mode synchronous or semi-synchronous replication is designed to prevent.
What is failover, and what's the difference between automatic and manual failover?
Failover is promoting a replica to take over as the new primary after the original primary fails. Manual failover requires a human to detect the failure and trigger the promotion — slower but deliberate; automatic failover has the system detect and promote on its own — faster recovery, but riskier if the failure detection itself is wrong (e.g. a network blip mistaken for a real outage).
What is "split-brain" in a replicated system, and how does it happen?
Two nodes each believe they are the legitimate primary at the same time — often because a network partition made the old primary unreachable, a replica was promoted, and then the original primary came back online still thinking it's in charge. Both now accept writes independently, and those writes can conflict or be lost when the partition heals.
How do quorum-based writes help prevent split-brain or ambiguous leadership?
By requiring a write (or a leadership claim) to be acknowledged by a strict majority of nodes before it's considered valid, it becomes mathematically impossible for two disjoint groups to both reach a majority at the same time — so at most one side of a network partition can ever have a legitimate primary.
What is the "read-after-write" (read-your-writes) consistency problem in replicated systems, and how might a system solve it?
A user writes data, then immediately reads it back — but if that read is routed to a replica that hasn't caught up yet, they see their own change as if it never happened. Common fixes: route a user's reads to the primary for a short window after their own write, or route their reads consistently to the same replica they wrote through.
Scenario: a user updates their profile picture, then immediately reloads the page and sees the old one. What's likely happening, and how would you fix it?
The reload's read almost certainly hit a replica that hasn't yet received the update — classic replication lag. Fixing it means giving that user a read-your-writes guarantee, e.g. briefly reading from the primary right after their own write, rather than routing every read to whichever replica happens to be available.
What's the difference between single-leader, multi-leader, and leaderless replication topologies?
Single-leader: one primary accepts all writes, replicas are read-only — simple, no write conflicts, but the primary is a bottleneck and single point of failure for writes. Multi-leader: several nodes each accept writes and replicate to each other — better write availability across regions, but now concurrent writes to the same data on different leaders can conflict. Leaderless: any node can accept a write, and the system uses read/write quorums to stay consistent enough — no single leader bottleneck, but conflict resolution and consistency reasoning become the client/system's job.
What kind of conflict can arise in multi-leader replication that doesn't happen in single-leader?
Two different leaders can each accept a write to the same record at nearly the same time, each believing theirs is the only update — when the two leaders later replicate to each other, there are now two conflicting versions of the same record and no single, obviously correct order between them.
How is a write conflict typically resolved in a multi-leader or leaderless system?
Common strategies: last-write-wins (pick whichever write has the later timestamp, accepting that the earlier one is silently discarded), merging the conflicting versions with application-specific logic, or surfacing the conflict for a human or the application to resolve explicitly — there's no universally correct answer, only tradeoffs.
What is "cascading replication," and why might you use it?
Instead of every replica pulling directly from the primary, some replicas replicate from other replicas instead, forming a chain or tree. This reduces the network and CPU load the primary must spend serving every replica directly, at the cost of extra replication hops (and so extra lag) for replicas further down the chain.
What is log-based (WAL-shipping) replication versus statement-based replication, and what's the tradeoff?
Log-based replication ships the primary's low-level write-ahead log entries (the exact bytes changed) to replicas, which apply them directly — deterministic and exact, but tied to the specific storage engine/version. Statement-based replication ships the actual SQL statements to re-execute on each replica — more portable and compact, but any statement that isn't perfectly deterministic can produce different results on each replica.
What could go wrong with statement-based replication that log-shipping avoids?
A statement using something non-deterministic — like NOW(), RAND(), or an auto-incrementing value computed independently — can execute at a slightly different moment or produce a different value on each replica, silently causing replicas to diverge from the primary even though every replica "correctly" ran the same statement.
How would you monitor replication lag in production, and what would you do if it grows unbounded?
Track lag as a first-class metric (seconds or bytes behind the primary) per replica, and alert past a threshold your application's tolerance allows. Unbounded, growing lag usually means a replica can't keep up with write volume — options include investigating what's slowing that specific replica down, temporarily routing more reads away from it, or adding capacity/replicas to spread the read load it was meant to absorb.
Scenario: replication lag balloons from 200ms to 30 seconds during a traffic spike. What's likely the cause, and what are your options?
Likely either a spike in write volume that the replica's apply process can't keep up with, or the replica itself being starved of resources (CPU/IO) by a surge of read queries also being routed to it. Options: temporarily reduce read traffic to that replica, scale it up, or investigate whether a batch job or write-heavy event is producing an unusual write burst.
What happens during promotion of a replica to primary, and what risk exists if the old primary comes back online afterward?
The chosen replica stops applying incoming replication and starts accepting writes directly as the new primary, and other replicas (and the application) are pointed at it. If the old primary comes back and isn't explicitly told it's been demoted, it may still think it's the primary and keep accepting writes independently — the split-brain scenario, with two divergent sets of writes that now need to be reconciled.
Why is automatic failover risky without careful design?
It has to correctly distinguish a genuine primary failure from a transient network blip — a false positive can trigger an unnecessary promotion and create split-brain once the 'failed' primary reappears — and it has to account for the fact that a promoted replica might not have every write the old primary had confirmed, risking silent data loss right at the moment of failover.
How does replication lag relate to the choice of consistency model?
Replication lag is the mechanical reason eventual consistency exists at all — an eventually-consistent system is essentially one that accepts replication lag as a normal, expected condition rather than something to eliminate. A system that instead requires strong consistency is, underneath, choosing to wait out (or avoid) that lag on the read path rather than tolerate it.
Scenario: an application always reads from a replica, but a user who just wrote data expects to see it immediately in the same session. What pattern fixes this without going fully synchronous everywhere?
Read-your-writes routing — send just that user's reads to the primary (or to a replica confirmed caught-up past their write) for a short window after their own write, while every other user's reads continue to go to any replica as normal. This targets the guarantee only where it's actually needed instead of paying its cost globally.
Why does adding more read replicas not help with write scalability?
Every replica still has to apply every single write to stay in sync — adding replicas multiplies how many places a write must eventually land, not how fast the primary can accept new writes. Only the primary (or a different scaling strategy like sharding) determines write throughput; replicas only help spread out reads.
Why can heavy read traffic against a replica itself worsen replication lag?
The replica has a fixed amount of CPU, memory, and IO capacity, shared between serving read queries and applying the incoming replication stream. If read load consumes enough of that capacity, the replica falls behind on applying replicated writes — meaning the very read-scaling replicas exist for can, under enough load, undermine how fresh their own data stays.
What's a realistic strategy for deciding how many replicas to run, and why isn't "more is always better"?
Size it to the read throughput you actually need to absorb beyond what the primary alone (or your caching layer) can handle, with some headroom for one replica being down. Every additional replica is another target for the primary to replicate to (more replication overhead) and another moving part to monitor and fail over — replicas you don't need trade operational complexity for capacity you'll never use.
Why can it be risky to promote a replica that was itself lagging at the time the primary failed?
That replica is missing whatever writes hadn't reached it yet, so promoting it doesn't just risk losing the writes that never reached any replica — it can also silently lose writes that other, more caught-up replicas already had, if the lagging one becomes the new source of truth for the whole system.
What is chain replication, briefly, and what property does it optimize for?
Replicas are arranged in an ordered chain; writes enter at the head and are passed down the chain, with a write only considered complete once it's reached the tail. It optimizes for strong consistency with clearly defined read/write rules (e.g. reads served from the tail always see fully-replicated data), at the cost of write latency scaling with the chain's length.
Scenario: you need strong, immediate consistency for one critical page (an order confirmation) right after a write, but want general read scaling everywhere else. How would you architect this with replication?
Route reads for that specific page (or that user's data, for a short window after their write) directly to the primary, while every other, less consistency-sensitive read continues to hit replicas as usual — applying the stronger guarantee narrowly, only where staleness would actually cause a problem.
Follow-up: what's the cost of routing reads to the primary for that one page?
It adds load to the primary that read scaling via replicas was specifically meant to avoid, and that page's read latency is now tied to the primary's load rather than benefiting from spreading across many replicas — an acceptable, narrow tradeoff only because it's limited to the one page that actually needs it.
Why is monitoring replication lag considered a first-class operational metric rather than an afterthought?
Lag directly determines how stale the data behind every read replica actually is at any given moment — without visibility into it, a team has no real way to know whether their eventual consistency guarantees are holding at milliseconds or have silently degraded to minutes, which changes what the rest of the system can safely assume about read freshness.