What does adding read replicas to a relational database actually scale, and what does it not scale?
answer
- replicas replay the same write log
- read capacity × N, write capacity × 1
- total write work = W × (N+1)
- async ⇒ stale reads, lag
- HA target vs throughput are different sizing goals
basics
~20 sReplicas add read capacity: full copies of the data that can serve queries. They do not add write capacity — every replica must apply the same write stream as the primary — and with asynchronous replication their data is slightly behind, so reads can be stale.
solid answer
~60 sA read replica is a full copy of the database kept up to date by streaming the primary's change log and replaying it. Because it holds all the data, it can answer any read query, so you can fan reads out across N replicas and multiply read throughput. What it does **not** do: - **Scale writes.** Every write still executes on the primary *and* is replayed on every replica. Adding replicas increases total write work in the system, and replica apply is often less parallel than the primary's write path — so a write-heavy workload can make replicas fall behind rather than help. - **Give you fresh data.** Asynchronous replication means a replica is behind by anything from microseconds to minutes under load. Read-after-write on a replica can return the old value. - **Reduce storage.** Each replica stores the whole dataset. - **Shrink an oversized working set.** If the problem is that the data doesn't fit in memory, every replica has the same problem. Replicas are also often deployed for high availability (failover target) rather than throughput; the two goals share the mechanism but not the sizing.
go deeper
Say clearly that replicas multiply reads, never writes, and that async replication makes replica data slightly stale.
Add the arithmetic (total write work grows with replica count) and name which reads must still hit the primary.
Discuss lag causes and monitoring, sync vs async tradeoffs, and the fact that moving analytics off the primary can be the real win even when read throughput isn't the constraint.
Frame replicas as a capacity and isolation tool with a consistency cost, and be explicit about which SLOs (read latency, failover RPO/RTO, analytics isolation) each replica is actually bought for.
## What a replica is Primary–replica (leader–follower) replication works from the write-ahead log. Every change the primary makes is recorded in an ordered log — WAL in PostgreSQL, the binlog in MySQL, redo in Oracle. Replicas stream that log and replay it, so they converge to the same state. Reads on a replica see a consistent snapshot of the data *as of some point in the log*, which is at or behind the primary's current position. ## The scaling arithmetic Suppose the system does R reads/s and W writes/s, and one node can handle a mix of both. - Reads distribute: with N replicas, each serves roughly R/N. Read capacity is genuinely multiplied. - Writes do not distribute: the primary does W, and *each replica also does W* as replay. Total write work in the cluster is W × (N+1). So replicas are a lever on the read side only, and they make the write side slightly worse (extra log shipping, extra replay CPU/IO). If you double the replica count and your write load is already near the primary's ceiling, nothing improves and lag gets worse. ## Where replicas genuinely help - Read-heavy application workloads (typical web read:write ratios of 10:1 or higher). - Long-running analytics and reporting queries that would otherwise compete with OLTP on the primary — moving them off protects primary latency even if read throughput was never the constraint. - Geographic read locality: a replica near the users cuts round-trip time. - Backups and heavy exports, taken from a replica. - High availability: a replica is the failover target when the primary dies. ## Where they don't - **Write-bound systems.** If commits, index maintenance, or lock contention are the bottleneck, replicas add load rather than relief. Write scaling requires different tools: batching, reducing index count, queueing, or sharding. - **Consistency-sensitive reads.** Anything that must see its own write, or must be authoritative (checking a balance before a debit, uniqueness checks, `SELECT ... FOR UPDATE`) must go to the primary. - **Working set too big for RAM.** Replicas replicate the problem. - **Connection exhaustion.** If the ceiling is connection count rather than query cost, a pooler helps more than a replica. ## Synchronous vs asynchronous Asynchronous replication (the default almost everywhere) means the primary commits without waiting for replicas; replicas lag. Synchronous replication makes the primary wait for at least one replica to acknowledge before commit returns — this bounds data loss on failover and can bound staleness, but it puts network round-trip time into every commit and makes the primary's write latency depend on the slowest acknowledged replica. Teams generally use one synchronous replica for durability/HA and additional asynchronous replicas for read scaling. ## Replica lag: causes worth naming - Long or large transactions on the primary appear as a single big chunk of replay. - Replay parallelism limits: MySQL historically applied the binlog single-threaded (multi-threaded replication improved this); PostgreSQL's WAL replay is a single startup process. - Long-running read queries on the replica can conflict with replay; PostgreSQL either cancels them or, with `hot_standby_feedback`, delays vacuum on the primary and risks bloat. - Network throughput and replica disk IOPS. Monitoring lag (seconds behind, or bytes/LSN behind) and routing traffic away from a lagging replica is basic hygiene. ## What a good answer sounds like "Replicas multiply read capacity and isolate heavy reads from the primary, and they're the failover target for HA. They don't help writes — every replica replays the whole write stream — and asynchronous replication means their data is stale, so any read that must see its own write, or must be authoritative, still goes to the primary."
- Your application is write-bound and someone proposes adding three more read replicas. What do you say?That it will not help and will likely hurt. Every replica replays the full write stream, so the cluster's total write work goes up while the primary's write ceiling stays the same, and lag typically gets worse. The right moves for a write bottleneck are on the write path itself — batching, removing redundant indexes, shortening transactions, moving non-critical writes to an asynchronous queue, or partitioning/sharding — after measuring which part of the commit path is saturated.
- You add replicas and read latency improves, but the primary's CPU barely drops. What might be happening?Probably that most primary CPU is write work — commit, index maintenance, replication log generation — which replicas cannot take. It can also mean the split is not routing much traffic: only queries explicitly marked read-only are going to replicas, while the bulk still runs on the primary inside read-write transactions. I would measure the read/write split at the connection level before adding more capacity.
Photocopies of a ledger: you can hand out copies so many people read at once, but every entry still has to be written in the master ledger first and then copied into every duplicate — copies never speed up writing.
saying these in an interview costs you the question
- Saying replicas scale the database generally, without separating reads from writes.
- Not realising every replica must apply the entire write stream.
- Assuming replica reads are always up to date.
- Confusing high availability (failover target) with throughput scaling.
- Thinking replicas reduce storage cost or shrink the working set.