Replication Strategies

Topics Covered

Why Replicate

Fault Tolerance

Geographic Proximity

Read Scaling

The Tradeoffs

Leader-Follower Replication

How the Replication Log Works

Adding New Followers

Synchronous vs Asynchronous Replication

Failover: When the Leader Dies

Automatic vs Manual Failover

Multi-Leader Replication

When Multi-Leader Makes Sense

Multi-Leader Topologies

The Conflict Problem

Conflict Avoidance

Conflict Resolution Strategies

Leaderless Replication

Quorum Writes and Reads

Read Repair and Anti-Entropy

Sloppy Quorums and Hinted Handoff

When Quorums Are Not Enough

Replication Lag and Its Effects

Read-After-Write Consistency

Monotonic Reads

Consistent Prefix Reads

The Consistency Spectrum

Monitoring Replication Lag

Choosing Consistency Guarantees by Use Case

Every production database is one hardware failure away from total data loss. A single server holds your data, and when that server's disk dies, your data dies with it. Replication solves this by keeping copies of the same data on multiple machines. But fault tolerance is only one of three reasons to replicate, and understanding all three is essential to choosing the right strategy.

Fault Tolerance

If one node goes down, another node has the same data and can serve requests immediately. Without replication, a disk failure means downtime until you restore from backup, and backup restoration takes minutes to hours. With replication, failover can happen in seconds. The key insight is that replication provides continuous availability, not just recovery. Your system keeps serving reads (and possibly writes) during the failure window.

Consider what happens when a single database server fails at 3am. Without replication, the on-call engineer receives an alert, connects to the infrastructure, provisions a new server, restores from the most recent backup (which might be 12 hours old), and then the application comes back online. Total downtime: 30 minutes to several hours. Data loss: everything written since the last backup. With replication, a follower is automatically promoted, clients are redirected, and the system continues serving requests. Total downtime: seconds. Data loss: minimal or zero.

Hardware failures are not rare events. At scale, they are routine. A server with one disk has roughly a 2-4% annual failure rate. A cluster with 100 servers expects 2-4 disk failures per year. With 1,000 servers, you expect a failure every few weeks. Replication is not insurance against an unlikely event. It is a necessary architecture for any system that cannot afford hours of downtime multiple times per year.

Beyond disk failures, there are other failure modes that replication protects against: power supply failures, memory errors, kernel panics, datacenter cooling failures, and even human errors like accidental data deletion (though replication propagates deletes to all replicas instantly, so backups are still needed for that case). The common thread is that any single machine is unreliable over long time horizons. Replication transforms an unreliable single machine into a reliable multi-machine system.

An important distinction: replication protects against hardware and infrastructure failures but does not protect against data corruption bugs. If a bug in your application writes incorrect data to the leader, that incorrect data is faithfully replicated to all followers. Both your primary and all replicas now contain the same wrong data. For protection against application-level errors, you need backups with point-in-time recovery or an event log that allows replaying history.

Geographic Proximity

Users in Tokyo experience 150ms round-trip latency to a server in Virginia. Place a replica in Tokyo and that latency drops to 5ms. For read-heavy applications (product catalogs, news feeds, user profiles), geo-distributed replicas let every user read from a nearby copy. The tradeoff is write complexity: a user in Tokyo who updates their profile must wait for that write to propagate across the ocean.

The physics of the speed of light imposes a hard floor on latency. Light in a fiber optic cable travels at roughly two-thirds the speed of light in a vacuum. A round trip from Tokyo to Virginia covers approximately 20,000 km of cable, giving a minimum latency of roughly 100ms. No amount of software optimization can beat this. The only solution is to move the data closer to the user through replication.

The impact on user experience is significant. Research consistently shows that each additional 100ms of latency reduces conversion rates by roughly 1%. For a user in Tokyo reading a product page served from Virginia, every single page load pays a 150ms penalty that replication would eliminate. Across millions of page views per day, the cumulative impact on revenue is substantial. This is why every major content platform (Netflix, YouTube, Amazon) places read replicas in every region where they have users.

Read Scaling

A single PostgreSQL instance handles roughly 10,000 to 50,000 reads per second depending on query complexity. When your traffic exceeds that, you cannot make reads faster on one machine. But you can spread reads across 5 replicas, each handling its share. Writes still go to one place, but reads (which are typically 90-99% of traffic) scale horizontally. This is why read replicas are the first scaling lever most teams pull.

The math is straightforward. If your application handles 100,000 reads per second and a single node maxes out at 25,000, you need 4 read replicas (plus the leader serving reads) to handle the load. Each replica is a full copy of the data, so any read can go to any replica. A load balancer distributes reads round-robin. When traffic grows to 200,000 reads per second, you add 4 more replicas. The scaling is linear and predictable.

This is also why read replicas are often the first scaling strategy teams adopt before sharding or caching. Adding a read replica requires only a database configuration change and a load balancer update. Sharding requires rearchitecting queries, managing shard keys, and handling cross-shard joins. Caching requires invalidation logic, cache-aside patterns, and reasoning about staleness. Replication delivers immediate read throughput gains with the lowest operational complexity.

The Tradeoffs

Replication is not free. Every replica consumes storage, memory, and network bandwidth. More importantly, replication introduces consistency challenges. When you have one copy of data, every read sees the latest write. With multiple copies, a read might hit a copy that has not received the latest write yet. This gap is called replication lag, and managing its effects is the central challenge of replicated systems.

Replication also complicates writes. In leader-follower setups, all writes must go through one node, which becomes a write bottleneck. In multi-leader setups, concurrent writes to different leaders create conflicts that must be detected and resolved. In leaderless setups, writes must be sent to multiple nodes and coordinated through quorums. Each topology makes a different tradeoff between write simplicity, consistency, and availability. Your choice of replication topology is one of the most consequential architectural decisions you will make. It affects every aspect of your system: write latency, read consistency, failure behavior, operational complexity, and the types of bugs your team will encounter in production.

The three replication topologies covered in this lesson map to different requirements:

Leader-follower is simplest and best for single-region read scaling with strong write ordering. PostgreSQL, MySQL, and MongoDB default to this topology. It is the right choice for applications that value simplicity and strong write ordering.

Multi-leader enables multi-region writes with low latency. CockroachDB and Spanner use consensus-based replication (similar to multi-leader with automatic conflict resolution) for globally distributed databases that need strong consistency. It is the right choice for applications that value write latency across geographic regions.

Leaderless provides maximum write availability with no single point of failure. Cassandra, DynamoDB, and Riak use this for high-availability workloads where eventual consistency is acceptable. It is the right choice for applications that value write availability above all else.

Most applications start with leader-follower because it is the simplest and best supported. Multi-leader and leaderless topologies are only justified when specific requirements demand them. Each additional topology choice adds operational complexity that must be weighed against the benefits.

It is also worth noting that you can combine topologies within a single system. A common pattern is leader-follower replication within a datacenter (for read scaling and fast failover) combined with multi-leader replication between datacenters (for geographic write distribution). Each datacenter has a leader with local followers, and the leaders across datacenters replicate to each other. This layered approach gives you the simplicity of leader-follower for local operations and the latency benefits of multi-leader for cross-region writes.

Interview Tip

In interviews, state all three reasons for replication when the topic comes up. Many candidates only mention fault tolerance. Mentioning geo-proximity and read scaling shows you understand replication as a performance tool, not just a reliability tool.