Data Engineering & Real-Time Event StreamsPlaybook3 min readUpdated September 2026

Living With Replication Lag Instead of Pretending It Doesn't Exist

Handling read replica lag means deciding, for each read, how stale the data is allowed to be, because a replica is always at least slightly behind the primary. The classic failure is a user who writes something, reads it back from a replica before replication catches up, and sees stale information.

There's no single fix that's right for every case; the right approach depends on how stale a given read is actually allowed to be, and being explicit about that per use case is what prevents the bug rather than a single blanket policy applied everywhere.

Why lag happens and why it's not a sign something's broken

A replica applies changes from the primary's write ahead log asynchronously, and under normal healthy operation, that application happens within a small fraction of a second. Lag becomes a real problem when it grows well beyond that baseline, usually because the replica is under heavier read load than it can sustain while also keeping up with replication, or because a long running query is blocking the replay of incoming changes.

Monitor replication lag as its own explicit metric, not just replica CPU or query latency, since a replica can look perfectly healthy on every other dashboard while quietly falling further and further behind on replication specifically. A lag metric that's trending upward over days, even if it's still within an acceptable range today, is worth investigating before it becomes a customer facing problem.

Routing reads by how fresh they actually need to be

Not every read has the same freshness requirement, and treating them all identically usually means either overloading the primary unnecessarily or serving stale data somewhere it actually matters. A dashboard showing historical trends can comfortably read from a replica that's a few seconds behind; a confirmation screen showing the write a user just made cannot, and that specific read needs to go to the primary or to a replica you've confirmed has caught up.

Build this routing logic explicitly into your data access layer rather than leaving it to individual engineers to remember case by case. A default that routes to the primary unless a read is explicitly marked as replica safe is a more conservative and generally safer starting point than the reverse.

Read your own writes without hitting the primary for everything

For the specific case of a user reading back their own recent write, a few patterns avoid both the staleness bug and the cost of routing every read to the primary. One is tracking the primary's log sequence number at write time and having the client wait for a replica to reach at least that position before serving the read, which guarantees freshness without a blanket primary only policy. A simpler, less precise option is routing a user's own reads to the primary for a short window right after they write, then falling back to replicas afterward.

Whichever pattern you choose, test it under real replication lag, not just under the low lag conditions a local development environment usually has. A pattern that works fine when lag is a few milliseconds can fail in ways that only show up when lag climbs to a full second or more under real production load.

Common ways to avoid the read your own write bug:

  • Track the primary's log sequence number at write time, and have the client wait for a replica to reach that position before the read is served.
  • Route a user's own reads to the primary for a short window right after they write, then fall back to replicas.
  • Send freshness sensitive reads, such as a confirmation screen, to the primary while dashboards keep reading from replicas.
  • Use synchronous replication for data where staleness is unacceptable, accepting the added write latency that comes with it.

When synchronous replication is worth the latency cost

Synchronous replication, where a write isn't acknowledged until at least one replica has confirmed it, eliminates the staleness problem entirely for whatever replicas are part of that synchronous set, at the direct cost of added write latency since every write now waits on network round trips to the replica. This tradeoff is worth making for data where staleness is genuinely unacceptable, like a financial transaction record, and usually not worth it for the bulk of ordinary application data where a few hundred milliseconds of eventual consistency is a completely reasonable tradeoff.

Mixing synchronous and asynchronous replicas for different tables or workloads is a valid middle ground, rather than treating replication mode as a single all or nothing setting for the whole database. Reserve the latency cost for the specific data where it's actually earning something in return.

Executive Capability Standard

What Good Looks Like

Handling replica lag well means reads are routed by their actual freshness requirement, replication lag is monitored as its own metric, and the read your own write pattern is deliberately chosen and tested under real lag conditions, not assumed to work.

Building The Capability (5-Stage Skill Ladder)

1. Learn:Audit your current read paths and identify which ones read from a replica without any explicit freshness requirement having been considered.
2. Do Manually:Add explicit primary routing for your highest risk read your own write case first, then expand the pattern to other reads.
3. Delegate:Assign a backend engineer to own read routing policy and to monitor replication lag as a standing metric, not just an incident response tool.
4. Automate:Build read routing directly into your data access layer with a conservative default, so new code doesn't have to reason about staleness case by case.
5. Buy:Bring in a fractional CTO or database specialist if a staleness bug has already caused a real customer facing incident, like a user not seeing their own recent action.

How to Get Started

Frequently Asked Questions

How much replication lag is actually acceptable?

It depends entirely on the read. A dashboard showing historical trends can tolerate several seconds of lag without anyone noticing; a screen confirming a user's own write generally can't tolerate any. Rather than picking one acceptable lag number for your whole system, classify reads by their actual freshness requirement and route accordingly.

What's the simplest way to avoid the read your own write staleness bug?

Route a user's own reads to the primary for a short window immediately after they write, then fall back to replicas afterward. It's less precise than tracking log sequence numbers, but it's simple to implement and eliminates the most common version of this bug for most applications.

Is synchronous replication worth the added write latency?

For data where staleness is genuinely unacceptable, like a financial transaction record, generally yes. For most ordinary application data, the latency cost isn't worth it when eventual consistency of a few hundred milliseconds causes no real problem. Mixing modes per table or workload is a reasonable middle ground rather than an all or nothing choice.

About the numbers

This guide doesn't quote a sourced benchmark. Figures in it are estimates or general guidance, so check them against your own numbers.

Related Guides