skip to content

In an asynchronous primary-replica setup, what causes replication lag, and what techniques can an application use to avoid serving a user stale data right after their own write?

level: middleimportance: must knowfreq 85%

answer

  1. async trades freshness for throughput
  2. read-your-writes
  3. monotonic reads
  4. sticky sessions
  5. LSN/GTID position tracking

basics

~20 s

The replica copy of the data takes a little time to catch up after the main database changes, so reading from a replica right after writing can show old data. Apps fix this by reading recent writes from the primary or by waiting for the replica to catch up.

solid answer

~40 s

Replication lag is the delay between a write committing on the primary and that write becoming visible on a replica; it grows when the primary's write throughput outpaces a replica's apply rate, under long-running transactions, network delay, or a replica doing heavy read load or maintenance work. It can be reduced with semi-synchronous replication for critical paths, but most systems stay async for throughput and instead solve staleness at the read layer: read-your-writes consistency (route the user's own subsequent reads to the primary, or wait for a replica to reach a tracked log position), monotonic reads (pin a session to one replica so it never sees data 'rewind'), or bounded-staleness reads that exclude a replica once its lag exceeds a threshold.

go deeper

for a junior

Knows replicas can be behind and that reading right after writing can show stale data.

for a middle

Can explain what causes lag and describe read-your-writes/sticky-session mitigations.

for a senior

Can choose sync vs async per use case, design lag-aware read routing, and reason about monotonic/causal consistency.

for a principal

Sets org-wide staleness SLAs per data class, weighs RPO/RTO trade-offs, and decides when to escalate specific critical paths to synchronous replication or quorum reads.

## What lag actually is Replication lag is simply the gap in time between a write being durably committed on the primary and that same write becoming visible on a given replica. In an asynchronous replication setup - the default for most relational and NoSQL systems because it keeps write latency low - the primary acknowledges a write to the client as soon as it is durable locally, then ships the corresponding log entry (a WAL segment in Postgres, a binlog event in MySQL) to replicas in the background. The replica reads that stream and re-applies each change in order. ## What widens the gap Under normal conditions this happens in tens of milliseconds, but several things can widen the gap: - the primary accepting writes faster than a replica can apply them (a burst of traffic, a bulk import, or a single very large transaction the replica must apply as one unit); - network latency or bandwidth constraints, especially across regions; - the replica itself being busy doing something else, like serving a flood of read queries or running index/vacuum maintenance; - or the replica simply running on weaker hardware. Lag can range from milliseconds to, in pathological cases, minutes. ## Why it matters to a reader This matters operationally because many architectures route reads to replicas specifically to spread load off the primary. If lag is nonzero, a client that just wrote data and then reads it back from a replica can see the pre-write state - the classic **'read-your-writes' violation**, one of the most common bugs reported in systems that naively load-balance reads across replicas. A subtler problem is **monotonic-read violation**: if a client's successive reads are load-balanced across replicas with different lag, it can see data 'go backward in time' - read a value, then on the next request read an older value from a less-caught-up replica. A third problem is **causal consistency** across related writes: if write A causally precedes write B (e.g. a comment posted after a post), and B propagates to a replica faster than A, a reader can see B without A, which looks broken to a user. ## You cannot simply delete the lag There is no free way to eliminate lag entirely without giving up the throughput/latency benefit of async replication. The alternative, synchronous replication, has the primary wait for the replica to apply the write before acknowledging - this trades write latency and availability, since a slow or down replica now blocks writes, for zero lag. Most production systems instead solve staleness at the application or routing layer rather than eliminating it at the source. ## Concrete mitigation techniques Concrete mitigation techniques - (1) read-your-writes routing, (2) sticky sessions, (3) bounded-staleness reads, (4) client-side consistency tokens, (5) return the data the client just sent: 1. **Read-your-writes routing** - after a write, route that user's subsequent reads to the primary for some window, or track the write's log sequence number/commit timestamp and have the read wait until a replica reports it has applied at least that position (Postgres exposes this via `pg_last_wal_replay_lsn()`; MySQL via GTID position checks). 2. **Sticky sessions / monotonic reads** - pin a user's session to a single replica (e.g. via consistent hashing on session ID) so they only ever move forward in time, never backward. 3. **Bounded-staleness reads** - have the load balancer or client library check a replica's lag metric and exclude any replica lagging beyond an acceptable threshold, e.g. one second, from the read pool. 4. **Client-side consistency tokens** - some managed databases (e.g. DynamoDB's `ConsistentRead` flag, or Spanner's read timestamps) let the client explicitly request a strongly consistent read when freshness matters, and fall back to cheaper eventually-consistent reads elsewhere. 5. **Return the data the client just sent** instead of re-reading it from the database - often the simplest fix, since the client already knows what it just wrote and doesn't need a round-trip to confirm it. ## How it shows up in production In production, replication lag shows up as a support ticket, not an error message: a user updates their profile picture, refreshes, and still sees the old one; a user places an order and the order-confirmation page, read from a replica, briefly says 'order not found.' Teams that don't design for it end up bolting on ad hoc fixes - e.g. 'always read orders from the primary for the first five seconds after creation' - which is exactly the read-your-writes pattern above, just discovered the hard way. This exact bug pattern (a user's own post briefly disappearing after posting because the detail page read from a lagging replica) is common enough in social and e-commerce apps that routing post-write reads to the primary, or to a confirmed-caught-up replica, is now considered standard practice rather than a special case.

  • Why doesn't a database just always use synchronous replication to avoid lag entirely?
    Because it ties write latency and availability to the slowest or least reachable replica - every write must wait for that replica to acknowledge, so a slow network link or a stalled replica directly stalls or blocks writes. Teams reserve synchronous replication for a small critical subset of data, or accept the latency cost only when durability guarantees demand it.
  • How would you detect that a replica has fallen dangerously behind before routing traffic to it?
    Expose and monitor a lag metric - the difference between the primary's and replica's log position or timestamp - and have the load balancer or connection pool exclude replicas whose lag exceeds a threshold. Postgres exposes this via pg_stat_replication; MySQL exposes Seconds_Behind_Master.
  • What's the difference between read-your-writes consistency and monotonic-read consistency, concretely?
    Read-your-writes guarantees a user sees their own prior writes on their next read. Monotonic reads guarantees a user's reads never go backward in time, even for data they didn't write themselves. You can have one without the other - sticky-session routing alone gives monotonic reads but not read-your-writes if the sticky replica hasn't caught up to the write yet.

Replication lag is like a group chat where one person types a message (primary) and it takes a moment to appear on everyone else's phone (replicas); if you immediately ask a friend 'did you see what I just sent?' before it has synced to their phone, they'll say no even though you know you sent it.

saying these in an interview costs you the question

  • says replicas are always perfectly in sync
  • proposes eliminating lag by 'just making replication faster' with no architecture change
  • doesn't distinguish read-your-writes from general staleness
  • assumes synchronous replication has no cost
  • can't name a concrete mitigation beyond 'add caching'

context