In a read/write-split architecture where writes go to a primary database and reads are load-balanced across asynchronous read replicas, what concrete strategies can a team use to avoid serving a user stale data right after their own write, and what does each strategy cost?
answer
- async replication → lag is inherent, not a bug
- primary pinning after write = simplest fix, adds primary load
- lag-aware routing via heartbeat/LSN/GTID
- session stickiness = monotonic reads per user, not global
- classify reads by staleness tolerance first
basics
~20 sRight after saving something, reading it back from a delayed copy can show the old version. Teams fix this by briefly reading from the same up-to-date database that took the write, or by tracking how far behind each copy is and steering around the ones that are behind, or by giving up strict freshness only for reads where it doesn't matter.
solid answer
~50 sThe core problem is asynchronous replication lag: a replica may not yet reflect a write that just committed on the primary. Common mitigations: (1) route reads for a short window after a write to the primary (or the same connection/session) - simple, but adds load back to the primary and only covers that user's own follow-up read; (2) track per-replica replication lag (e.g., via a heartbeat table or LSN/GTID comparison) and route reads only to replicas under a lag threshold, falling back to the primary if all replicas are too far behind - adds monitoring/routing complexity; (3) session/sticky consistency, pinning a user's session to one replica or the primary for a period so they see a monotonically advancing view - simpler than global lag tracking but doesn't fix cross-user staleness; (4) accept eventual consistency for reads that don't need freshness (activity feeds, recommendations) and reserve stronger guarantees only for reads that need them, minimizing how much traffic pays the routing/monitoring cost at all.
go deeper
Should understand that replicas can be behind the primary and that reading from one right after a write can show old data.
Should be able to describe primary-pinning-after-write as a basic fix and recognize it as the direct cause of the 'my change disappeared' bug.
Should be able to propose and compare multiple mitigation strategies (pinning, lag-aware routing, session stickiness, staleness-tolerant classification) with their respective costs, and pick appropriately per read type.
Should design the routing/consistency policy as a system-wide budget - deciding which subsystems get which guarantee, how lag monitoring itself is kept reliable, and how the policy evolves as replica count and load grow.
## Why lag exists at all **Asynchronous replication** - the default mode for read replicas in essentially every mainstream relational database (MySQL, PostgreSQL streaming replication, etc.) - means the primary commits a write and returns success to the client without waiting for any replica to apply it. The replica receives the change afterward, over the network, via a replication stream (binlog events in MySQL, WAL segments in PostgreSQL), and applies it on its own schedule. - Under light load this lag is typically single-digit milliseconds and invisible in practice. - Under load spikes, long-running queries competing for the replica's resources, network partitions, or a replica falling behind because it simply can't apply changes as fast as they arrive, lag can grow to seconds or, in bad cases, minutes. A read/write-split architecture that blindly load-balances all `SELECT`s across replicas will, during any nonzero lag window, sometimes return data that doesn't yet reflect a very recent write - the canonical symptom being a user who just saved a change and, on the very next page load, sees the old value, which reads to them as data loss even though nothing was actually lost. ## Primary pinning after a write The most direct mitigation is **read-after-write consistency via primary pinning**: after a client performs a write, route that same client's subsequent reads to the primary (or to the connection that performed the write) for some bounded window - a few seconds, or for the remainder of the request/session - before letting them fall back to replica-routed reads. This is simple to reason about and guarantees the specific user who made the change sees it immediately, but it doesn't help a different user who reads the same data through a replica moments later, and it adds load back onto the primary exactly for the traffic that read/write splitting was meant to offload, so it has to be scoped narrowly (only the affected rows/user, only a short window) rather than applied broadly. ## Lag-aware routing A more general approach is **lag-aware routing**: continuously measure how far behind each replica is - commonly - via a small **heartbeat table** that the primary updates with a timestamp every second or so, which each replica then reflects with its own delay, letting the router compute 'this replica is currently 340ms behind' by diffing timestamps; - or via native replication position markers (PostgreSQL LSNs, MySQL GTIDs) compared between primary and replica. The router then excludes replicas whose lag exceeds a threshold from receiving read traffic, falling back to the primary (or to a less-lagged replica) if all replicas are over budget. This generalizes read-after-write consistency to any client, not just the one who wrote, but it costs real engineering: a monitoring/heartbeat mechanism, a routing layer that consults it on every read (or on a cached recent measurement), and a failure mode of its own - if lag monitoring itself lags or breaks, the system can either wrongly route to a stale replica or wrongly overload the primary by excluding replicas that are actually fine. ## Session or sticky consistency **Session or sticky consistency** takes a middle path: rather than tracking lag globally, pin a given user's session to a single replica (or the primary) for its duration, so that user always sees a monotonically advancing view of the data - they'll never see something 'go backward' between two of their own reads, even if their view is a few hundred milliseconds behind absolute real time. This is cheaper to implement (a hash of session/user ID to a replica, no lag telemetry needed) but only solves consistency within one user's session; two different users hitting two different replicas can still momentarily disagree about the current state, which is fine for most consumer-facing reads but not for anything requiring a single global source of truth at read time. ## Paying only where freshness matters The strategy that costs the least is simply not paying for strong consistency where it isn't needed: classify reads by how much staleness they can tolerate, and only route the freshness-sensitive minority (e.g., 'show me the order I just placed,' 'show my current account balance before I authorize a debit') through primary-pinning or lag-aware routing, while the freshness-tolerant majority (product listings, an activity feed, aggregate counts, search results) flow through plain replica load balancing with no special handling at all. A concrete real-world shape: a social app pins reads to the primary for the 2 seconds immediately following a user's own post (so they see their post appear instantly in their own feed), routes everyone else's view of that same feed through ordinary lag-tolerant replica reads (a few hundred milliseconds of staleness in someone else's feed is imperceptible), and reserves lag-aware routing specifically for a payments subsystem where even other users' reads of balance data must not be stale.
- How would you measure replication lag in practice without relying only on the database's built-in lag metric?A common technique is a heartbeat table on the primary that gets updated with the current timestamp every second or so; each replica eventually applies that update, and comparing the replica's heartbeat timestamp to the current wall-clock time gives an application-observable lag figure, independent of whatever internal metric the database reports. This is especially useful because it measures lag as actually observed by a client, including any queuing effects.
- Why doesn't session stickiness alone solve staleness for a public-facing read like a product's current stock count shown to many different users?Session stickiness only guarantees a monotonically advancing view for one pinned session/user - it says nothing about whether two different users, pinned to two different replicas with different lag, see the same value at the same moment. For data where all readers need to agree (like stock count near zero), you need either primary reads, lag-aware routing with a tight threshold, or a different consistency mechanism altogether.
- What's the downside of setting the lag threshold in lag-aware routing too conservatively (too low)?A very low threshold means replicas get excluded from serving reads more often, pushing more traffic back onto the primary during exactly the load spikes that tend to cause lag in the first place - potentially defeating the purpose of having replicas at all. It's a tuning trade-off between freshness guarantees and how much read offloading you actually get.
It's like a group chat where your own message appears instantly on your phone but takes a beat to sync to everyone else's - to make sure you never see your own message vanish, your app always shows you your latest local copy first, while other people's phones catch up a moment later.
saying these in an interview costs you the question
- Thinks replicas are always perfectly synchronous
- Proposes routing all reads to the primary as the fix (defeats the purpose of replicas)
- Doesn't distinguish read-your-own-writes from global consistency across users
- Can't name any concrete lag-measurement mechanism
- Assumes one consistency strategy fits every kind of read in the system