6.2 Problems with Replication Lag
And on where to solve it: dealing with these issues in application code is complex and easy to get wrong.
Replication lag
Read-after-write. The user submits a write, then immediately reads it back and it is missing — the read hit a follower that has not caught up. This is the anomaly users notice fastest, because they know what they just typed. With this lag the read lands inside the window and observes the anomaly.
Read-scaling architecture: many followers, reads distributed across them. This removes load from the leader and lets reads be served by nearby replicas.
This realistically works only with ASYNCHRONOUS replication. If you synchronously replicated to all followers, a single node failure or network outage would make the entire system unavailable for writing. And the more nodes you have, the likelier one is down — so a fully synchronous configuration would be very unreliable.
Eventual consistency: reading from an async follower may return outdated information. Run the same query on the leader and a follower at the same time and you may get different results. Stop writing and wait, and the followers eventually catch up.
The term "eventually" is deliberately vague; in general, THERE IS NO LIMIT to how far a replica can fall behind. Normally a fraction of a second — but near capacity or with network problems, easily seconds or minutes.
(Note: it's not only NoSQL databases that are eventually consistent — followers in an asynchronously replicated RELATIONAL database have the same characteristics.)
The three anomalies
“To the user, it looks as though the data they submitted was lost.” The guarantee is that a user always sees updates they submitted themselves; it promises nothing about other users.
Mr. Poons: “How far into the future can you see, Mrs. Cake?”
Mrs. Cake: “About 10 seconds usually, Mr. Poons.”
Mrs. Cake: “About 10 seconds usually, Mr. Poons.” ← answer
Mr. Poons: “How far into the future can you see?” ← question
Poons's shard replicated slower than Cake's. “Such psychic powers are impressive but very confusing.”
2.1 Implementing read-after-write consistency
| Technique | How | Limitation |
|---|---|---|
| Read modifiable things from the leader | E.g. always read the user's OWN profile from the leader, other users' profiles from a follower | Requires knowing whether something MIGHT have been modified without querying it. Fails if most things are user-editable — most reads would hit the leader, negating read scaling |
| Time-based | Track the time of the last update; for one minute after it, make all reads from the leader. Also: monitor replication lag and prevent queries on any follower more than one minute behind | Coarse; wastes leader capacity |
| Client remembers a timestamp | The client remembers the timestamp of its most recent write; the system ensures the serving replica reflects updates at least until that timestamp. If not sufficiently up to date, either route to another replica or WAIT until it catches up | Timestamp is either logical (a log sequence number) or the actual system clock — in which case CLOCK SYNCHRONIZATION becomes critical (Ch 9) |
| Cross-region | Any request that must be served by the leader must be routed to the region containing the leader | Added latency and complexity |
Cross-device read-after-write adds two more problems:
- Remembering the user's last-update timestamp breaks, because code on one device doesn't know what happened on the other. The metadata must be CENTRALIZED.
- Different devices may route to different regions — desktop on home broadband vs mobile on cellular data take completely different network routes. If your approach requires reading from the leader, you may first need to route all of a user's devices to the same region.
2.2 Implementing monotonic reads
Make sure each user always reads from THE SAME REPLICA (different users can use different replicas) — e.g. choose the replica by a HASH OF THE USER ID rather than randomly.
However, if that replica fails, the user's queries must be rerouted to another replica.
2.3 Implementing consistent prefix reads
This is particularly a problem in SHARDED databases. If the database always applies writes in the same order, reads always see a consistent prefix and the anomaly can't happen. But in many distributed databases, different shards operate independently, so there is no global ordering of writes — a reader may see some parts of the database in an older state and some in a newer state.
Solutions:
- Make sure any writes that are causally related to each other go to the same shard — but in some applications that can't be done efficiently
- Explicitly track causal dependencies (→ §4.5, happens-before)
2.4 The pragmatic advice
Think about how the application behaves if replication lag increases to several minutes or even hours. If the answer is "no problem," that's great. If it's a bad experience, design for a stronger guarantee.
Pretending that replication is synchronous when in fact it is asynchronous is a recipe for problems down the line.
And on where to solve it: dealing with these issues in application code is complex and easy to get wrong.
The simplest programming model is to choose a database providing STRONG CONSISTENCY (linearizability, Ch 10) and ACID transactions (Ch 8), so you can mostly ignore replication challenges and treat the database as if it had a single node.
The historical arc: in the early 2010s the NoSQL movement argued these features limited scalability and that large-scale systems must embrace eventual consistency. Since then, a number of databases provide strong consistency and transactions WHILE ALSO offering fault tolerance, high availability, and scalability — the NewSQL trend.
But weaker consistency still has legitimate reasons: stronger resilience in the face of network interruptions, and lower overheads compared to transactional systems.
6.1 Single-Leader Replication
The leader logs every write statement it executes and sends the statement log to followers; each follower parses and executes that SQL as if received from a client.
6.3 Multi-Leader Replication
Both perform automatic merges for all the types above; different design philosophies and performance characteristics.