Suppose your whole application depends on a single database server. The moment that one machine dies, loses a disk, or needs maintenance, your entire site goes down with it. And even while it's healthy, every read and every write in the system has to funnel through that one box — so as traffic grows, it becomes the bottleneck everyone is waiting on.
Replication is the answer to both problems. Instead of trusting one copy of your data, you keep several identical copies on different machines and keep them in sync. If one machine fails, another can take over. If reads pile up, you can spread them across the copies. This is one of the foundational tricks behind almost every large-scale system.
Why replicate at all
There are three big reasons to keep multiple copies of the same data. The first is availability and failover: if the primary database crashes, a replica that already holds the data can be promoted to take its place, so the system keeps running instead of going dark. The second is scaling reads: most applications read far more often than they write, and by letting many replicas answer read queries you multiply your read capacity without overloading a single machine.
The third reason is putting data closer to users. If your customers are spread across the globe, a single database in one region means everyone outside that region waits on a long network round-trip. By placing replicas in different regions, a reader in Tokyo can hit a nearby copy instead of crossing an ocean to reach a server in Virginia, cutting latency dramatically.
Leader and follower
The most common arrangement is leader/follower replication, also called primary–replica. One node is designated the leader (the primary), and it is the only one allowed to accept writes. Whenever a client changes data, the change is applied on the leader and then sent out to every follower (replica), which applies the same change to its own copy.
Reads, on the other hand, can be served by any node — the leader or any of the followers. This split is what makes the pattern so useful: writes stay simple because there's a single authoritative copy to coordinate them, while reads scale out across all the replicas.
How it works
Picture a client sending a write to the primary database. The primary records the change in its own storage and appends it to an ordered log of changes, numbered one after another. It then forwards each change to its replicas over the network. Every replica replays the log in order and remembers how far it has got, so its copy ends up matching. Meanwhile, read requests from clients are spread across all the copies — some hit the primary, but many are answered by the replicas, taking load off the leader.
The crucial detail is timing. This propagation is usually asynchronous: the primary confirms the write to the client immediately, before all the replicas have caught up. For a brief window, a replica may still be holding the old value while the primary already has the new one, and different replicas can be at different points in the log.
That timing matters most when the leader dies. Failover means promoting a follower to be the new leader, and the follower you promote decides what survives. Step through one write reaching the followers at different times; just before the end the leader crashes, and you predict which follower should take over.
The gap has a name: replication lag. It's the short delay between a write landing on the primary and that same change appearing on a replica. Usually it's milliseconds, but under heavy load or network trouble it can stretch to seconds — long enough for a user to notice stale data.
Eventual consistency
Because replicas lag behind the primary for a moment, a read served from a replica might return data that's slightly out of date. This is the heart of eventual consistency: reads may be stale for a short time, but if writes stop, all copies will converge to the same value given enough time. Nothing is permanently wrong — the system just isn't instantaneously identical everywhere.
This is a deliberate trade-off against strong consistency, where every read is guaranteed to see the most recent write no matter which copy answers it. Strong consistency is what the strict transactional model gives you on a single server — see ACID & transactions. Distributed systems often relax that guarantee in favor of the BASE philosophy (Basically Available, Soft state, Eventually consistent), accepting brief staleness so the system can stay fast and available.
Why accept staleness at all? The CAP theorem explains the bind: when the network between nodes breaks (a partition), a distributed system has to choose between staying consistent and staying available. Eventually-consistent systems lean toward availability, answering with possibly-stale data rather than refusing to respond — a fine choice for a like count or a feed, but a dangerous one for a bank balance.
Here's what a stale read looks like in time. A shop owner changes a price and a shopper's read lands on a follower that's busy catching up. Predict what the shopper sees, then switch the follower's lag from a busy day to a quiet one.
Your store reads prices from replicas. In testing, a price change always shows up immediately. During a big sale, some shoppers see the old price for a second or two after a change. What's the most likely explanation?
Pitfalls and variations
Replication lag creates some subtle bugs. The classic one is read-your-own-writes: a user updates their profile, the write goes to the primary, but their next read is served by a replica that hasn't caught up yet — so they see their old profile and assume the save failed. There are two common fixes. The simplest: for a short time after a user writes (say a minute), send that user's reads to the primary, which always has their change. The more scalable one: the write's reply carries its position in the log, the app keeps it in the user's session, and reads only go to a replica that has reached that position (or wait until it has). Simply pinning a user to one replica for their whole session isn't enough on its own: that replica can still be behind their write.
Step through it below. Predict what you see after the reload, then flip between the fixes and notice what each one costs.
Leader/follower isn't the only model. In multi-leader replication, more than one node can accept writes — useful when you have datacenters in several regions and want each to take local writes. The catch is conflict resolution: if two leaders accept conflicting changes to the same record at the same time, the system has to decide which one wins, which is genuinely hard to get right. For most applications, a single leader with several read replicas is the simplest setup that delivers the big wins.
After editing their bio, some users see the old bio when the page reloads, and they edit it again. Reads go to any replica. Which fix makes sure people see their own edits without sending everyone's reads to the leader?
Don't blindly read from a replica after a write. If correctness depends on seeing the freshest data — a payment confirmation, an inventory check at checkout — route that read to the primary or wait for the replica to catch up. Eventual consistency is great for tolerant data and quietly wrong for the data that really matters.