Learn Labs
6. Replication

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

anomaly
mitigation
leader
write v2
follower
v1v2 applied
client read
reads v1 — stale
Problem

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.

Eventual consistency is a statement about the limit, not about any particular read. The anomalies below are what 'eventually' feels like from inside an application.

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

① Read-after-write (read-your-writes)
userleaderfollowerwriteasync — not yet arrivedreadold value

“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.

② Monotonic reads — time appears to go backward
read #1follower Alittle lagsees commentread #2follower Bgreater lagcomment goneVery likely if the user refreshes and each request goes to a random server.
③ Consistent prefix reads — violation of causality
actual

Mr. Poons: “How far into the future can you see, Mrs. Cake?”
Mrs. Cake: “About 10 seconds usually, Mr. Poons.”

observed

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.”

Figure 6.2.1The three anomalies

2.1 Implementing read-after-write consistency

TechniqueHowLimitation
Read modifiable things from the leaderE.g. always read the user's OWN profile from the leader, other users' profiles from a followerRequires 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-basedTrack 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 behindCoarse; wastes leader capacity
Client remembers a timestampThe 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 upTimestamp is either logical (a log sequence number) or the actual system clock — in which case CLOCK SYNCHRONIZATION becomes critical (Ch 9)
Cross-regionAny request that must be served by the leader must be routed to the region containing the leaderAdded latency and complexity

Cross-device read-after-write adds two more problems:

  1. 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.
  2. 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.


On this page