What causes replication lag between a leader and its followers, and how would you design a system to guarantee 'read-your-writes' consistency for a user despite that lag?
answer
- lag = time gap leader-to-follower apply
- grows under load, slow follower, network issues
- read-your-writes broken by stale replica read
- fix: route to leader / LSN token / time window
basics
~20 sFollowers apply changes slightly after the leader because copying and applying takes time, especially under load. Read-your-writes means a user should always see their own recent changes - you guarantee it by making sure that right after someone writes, their next read goes to the leader or to a follower confirmed to be caught up, instead of a random possibly-stale replica.
solid answer
~40 sReplication lag is the delay between a write being committed on the leader and that same write being visible on a given follower; it grows under write bursts, slow network links, follower resource contention, or a follower doing expensive background work while catching up. It breaks the natural expectation that after you write something you'll see it if you immediately read it, because a load-balanced read can land on a follower that hasn't caught up yet. Common fixes: route the read that immediately follows a write to the leader for some time window; track a per-user 'last write position' and only route to followers whose position is at or past it; or always read a user's own writes from the leader for a short TTL after writing, falling back to any replica afterward.
go deeper
Should explain in plain terms that copies take a moment to catch up and that's why a fresh write might not show up right away.
Should name concrete causes of lag (load, slow follower, network) and describe at least one fix, like reading own writes from the leader.
Should compare the fixed-window vs log-position-token approaches, know their trade-offs, and connect this to why naive follower routing breaks read-your-writes.
Should design a routing policy across a fleet with mixed workload types and peak-load behavior, and reason about when the fix itself reintroduces the scaling bottleneck it was meant to solve.
## What replication lag is **Replication lag** is the time gap between the moment the leader durably commits a write and the moment a specific follower has fetched, applied, and made that write visible to its own readers. It is not a fixed constant - it fluctuates based on several concrete factors. - **Under normal load**, lag might be single-digit milliseconds, driven purely by network round-trip time between leader and follower. - **Under a write burst** (e.g., a batch job or flash-sale traffic), the leader can produce log entries faster than a follower's apply thread can consume and index them, so the follower's queue backs up and lag grows into seconds or minutes. - **When a follower is resource-starved** - doing a backup, running a heavy analytical query, rebuilding an index, or simply running on slower hardware than the leader - lag also grows, because applying incoming writes competes with that other work for CPU, disk I/O, and lock contention. - **Network partitions or packet loss** between leader and follower can pause replication entirely, in which case lag grows unbounded until the connection recovers. ## Why the lag is unavoidable This lag exists as an unavoidable consequence of physics and engineering trade-offs: making every follower synchronously caught up on every write (eliminating lag entirely) would mean synchronous replication to all followers, which collapses write availability and inflates write latency to the speed of the slowest replica. So virtually every system that reads from followers for scale accepts some amount of lag as the price of that scale-out. ## The read-your-writes problem it forces The trade-off it forces onto application design is exactly the **read-your-writes** problem: a user submits a write to the leader, gets a success response, and then immediately issues a read (e.g., reloading the page after posting a comment) that a load balancer routes to a follower which hasn't yet applied that write. From the user's point of view their own action appears to have silently failed or reverted, which is confusing and erodes trust even though no data was actually lost - it's purely a visibility problem caused by reading a stale replica. ## Techniques that guarantee read-your-writes There are several concrete techniques to guarantee read-your-writes despite lag, each with its own cost. 1. **The simplest** is to route any read that could plausibly be 'my own recent write' back to the leader - for example, always read a user's own profile/posts from the leader, while reading other users' content from replicas; this sacrifices some read scaling on that hot path but is simple to reason about. 2. **A more general version is time-based:** after a user writes, pin their reads to the leader (or force a read from a replica no more than N seconds behind) for some window, say one minute, then fall back to any replica once the risk of hitting stale content is low. 3. **A more precise version tracks log position:** the leader returns a log sequence number (`LSN`) or similar token with the write acknowledgment, the client (or a session store) remembers 'you're entitled to read at least up to position X,' and a router only sends the subsequent read to a follower whose applied position is at or past X, blocking briefly or falling back to the leader if no follower qualifies yet. This is more precise (it doesn't over-penalize users whose replica actually caught up quickly) but requires the follower's replication position to be queryable and the routing layer to track it per-request. ## Failure modes Failure modes to watch for in production: - **Naively pinning 'read from leader after every write'** for all users, all the time, defeats the entire purpose of having followers and just re-funnels all read traffic back to the leader, causing the exact write-throughput/read-scaling bottleneck replication was meant to solve. - **Under-provisioning the 'stickiness window'** (too short a TTL, or too coarse a time-based rule) lets stale reads leak through during traffic spikes when lag is worst - precisely when the naive fix is least effective, because lag is largest exactly when load is highest. - **Systems that use device/session affinity to a single follower** rather than log-position tracking can also produce the related but distinct monotonic-reads anomaly, where a second read from a different, more-lagged follower appears to go backwards in time relative to the first. ## Where it shows up A concrete real-world example: many web applications backed by MySQL or PostgreSQL read replicas implement 'read from primary after write' by writing to a session cookie or a fast key-value store the timestamp or LSN of the user's last write, and having the request router consult that value before choosing a replica; distributed SQL systems like Amazon Aurora expose a read-after-write consistency feature that does exactly this token-based routing under the hood, so application code just calls a consistent-read API without manually plumbing LSNs through every request.
- Why does replication lag tend to get worse exactly when read-your-writes violations matter most, at peak traffic?Peak traffic drives both a burst of writes, which followers must queue and apply, and a burst of reads, which land on those same lagging followers, so the two failure conditions compound at the worst possible time; naive fixed-size buffers or thread pools on the follower's apply path make this nonlinear once queueing begins.
- What's the downside of always routing a user's post-write reads to the leader for a fixed 60-second window, versus using a log-position token?The fixed window either over-serves, still hitting the leader long after the follower actually caught up and wasting the scaling benefit of replicas, or under-serves, since a slow follower can still be behind after 60 seconds under heavy load, so the guarantee silently breaks; a log-position token adapts exactly to real replication state instead of guessing a duration.
- If a mobile app caches a write locally and shows it optimistically before the server round trip completes, does that avoid the read-your-writes problem?It avoids the visible symptom for that one screen, but it's a client-side workaround, not a system-level guarantee - as soon as the user reloads the app or a different screen re-fetches from the server, the same stale-follower read problem can resurface unless the server-side routing is also fixed.
It's like mailing carbon copies of a memo to branch offices: the head office (leader) acts on the memo instantly, but each branch only reflects it once their copy physically arrives. If you call a branch five minutes after HQ approved something, you might be told 'no such memo exists' purely because the courier hasn't gotten there yet - not because the memo was rejected.
saying these in an interview costs you the question
- Thinks replication lag is constant/predictable rather than load-dependent
- Proposes eliminating lag by making everything synchronous with no mention of the availability cost
- Doesn't distinguish read-your-writes from monotonic reads
- Suggests fixing this purely with client-side caching with no server-side guarantee
- Assumes routing all reads to the leader forever is a reasonable general fix