Read Replica Returns Stale Data After Write: The LSN Token Pattern

Databases · Advanced · 7 min read · published

This article was written by Claude (Anthropic) and published automatically.

What this solves: You added read replicas and now users see their own edits disappear for a second. Here's how to route reads so writers always read their own writes.

The Forces at Play

You added a read replica, traffic on the primary dropped nicely, and then support tickets started: a user saves their profile, the page reloads, and the old name is back. The read replica returns stale data after a write, for a window of anywhere from 5ms to 30 seconds depending on what else the database is doing.

The tension is that three things you want are mutually hostile:

The crucial insight is that you almost never need global consistency. You need read-your-writes: the session that made the change must see it. Everyone else can tolerate a stale-by-200ms product listing. That's a much cheaper guarantee, and the architecture below buys exactly it and nothing more.

The Shape

A read router sits between the application and the database endpoints. Every write returns a write position token (a Postgres LSN, a MySQL GTID set, a Mongo cluster time). That token travels with the client's session. Every read presents the token, and the router picks a target that is provably caught up to it.

flowchart TB
  C[Client request<br/>carries lsn_token]
  subgraph App[Application process]
    R{Read router}
    W[Write path]
  end
  P[(Primary)]
  R1[(Replica A<br/>replayed: 0/3F2A100)]
  R2[(Replica B<br/>replayed: 0/3F29000)]
  H[Lag watcher<br/>polls replay LSN every 200ms]

  C -->|POST| W
  W -->|COMMIT| P
  P -.->|returns commit LSN 0/3F2A0F0| W
  W -->|set token in session| C

  C -->|GET + token 0/3F2A0F0| R
  H -. reads replay position .-> R1
  H -. reads replay position .-> R2
  H --> R
  R -->|token <= A.replayed<br/>OK| R1
  R -->|token > B.replayed<br/>too far behind| R2
  R -->|no replica caught up| P

  P ==>|WAL stream| R1
  P ==>|WAL stream| R2

The dashed line from the lag watcher is the part people skip, and it's what turns a guess into a guarantee: the router never asks a replica "are you fresh?" on the hot path. It consults a cached table of replay positions that a background poller refreshes.

How Data Flows Through It

One request, end to end.

  1. POST /profile arrives. The router sees a mutation and takes a primary connection.
  2. Inside the transaction the app writes the row. Immediately after COMMIT, on the same connection, it asks the primary for its current write position:
    COMMIT;
    SELECT pg_current_wal_insert_lsn();  -- '0/3F2A0F0'
    
  3. The handler puts 0/3F2A0F0 into the session store (or a signed cookie, or an X-Read-Position response header for a SPA). It represents "anything I have written is at or before this point."
  4. The response goes out. The browser immediately issues GET /profile with the token attached.
  5. The router reads its in-memory lag table, refreshed 40ms ago: Replica A replayed to 0/3F2A100, Replica B to 0/3F29000. A ≥ token, B < token. It picks A.
  6. Replica A serves the read. It has the new row. The user sees their edit.
  7. Now imagine a bulk import pushed both replicas 2 seconds behind. Neither satisfies the token. The router falls back to the primary for this one request — not for this user's next five minutes, not for the endpoint globally.
  8. Ten seconds later the token in the session is older than the worst replica's position. Every read from that session goes to a replica again, automatically. The pin expires by arithmetic, not by a timer you tuned.

A minimal router contract:

interface ReadRouter {
  // Returns a connection guaranteed to have replayed up to `minPosition`.
  // Falls back to the primary when no replica qualifies.
  acquireRead(minPosition?: WritePosition): Promise<Conn>;

  // Runs fn on the primary and returns the commit position with the result.
  write<T>(fn: (c: Conn) => Promise<T>): Promise<{ value: T; position: WritePosition }>;
}

// Comparison is not string compare — LSNs are two hex halves.
function lsnGte(a: string, b: string): boolean {
  const [ah, al] = a.split("/").map(x => BigInt("0x" + x));
  const [bh, bl] = b.split("/").map(x => BigInt("0x" + x));
  return ah !== bh ? ah > bh : al >= bl;
}

What Each Piece Owns

The primary owns durability and the authoritative write position. It does not own read scaling, and it deliberately does not know that replicas or routers exist.

The lag watcher owns one fact per replica: the highest position it has replayed, plus a freshness timestamp. It does not make routing decisions, and it must never be on the request path — if it is, you've added a round trip to every read and given yourself a new outage source.

The read router owns target selection and only that. Given a token and a lag table, choose an endpoint. It does not own retries, connection pooling policy, or query semantics. Keep it a pure function over (token, lagTable) so you can unit-test the nasty cases.

The session/token carrier owns propagation. Its job is to make sure the token survives a redirect, a page load, and a sticky-session-free load balancer. It does not own interpretation — it's an opaque string to everything except the router.

The application handler owns knowing whether an operation is a write. That sounds trivial until someone writes a GET that lazily backfills a row.

Where It Breaks Down

The first bottleneck is the fallback cliff, not the replicas. When lag spikes — vacuum, a bulk import, a long-running query on the replica blocking WAL replay — every recent-writer's read fails the token check simultaneously and stampedes onto the primary. The primary was sized for 5% of read traffic. Cap it: track the fraction of reads falling back, and above a threshold serve stale from a replica (or fail fast) rather than melting the primary. Make that threshold an explicit, alerted decision.

Token loss is silent. A client that drops the token doesn't error; it just quietly gets stale reads again, and only for the unlucky window. Log token presence on write-adjacent endpoints so you can see it happening.

Cross-session causality isn't covered. User A comments, user B is told about it via WebSocket, B's read has no token and sees nothing. If two users can observe each other's writes in real time, you need the token to travel through your event payloads too, not just the HTTP session.

Failover invalidates positions. After a promotion, LSNs on a timeline-switched cluster can compare in ways you didn't plan for. Stamp tokens with the cluster/timeline ID and treat mismatches as "use the primary."

Operational burden lands on the lag watcher. If the poller wedges or its cache goes stale, the router either routes everything to the primary (safe, expensive) or trusts old positions (fast, wrong). Build it to fail toward the primary and page on stale lag data.

When This Is Overkill

Most systems should start with the boring version: send every read to the primary, and route only obviously-tolerant workloads — analytics, exports, dashboards, background reports — to a replica. No tokens, no watcher, no router. This is correct by construction and it scales further than people think: a well-indexed Postgres primary on modern hardware handles tens of thousands of reads per second.

The next cheapest step, if the bug is narrow, is per-endpoint: mark a handful of read handlers as mustReadPrimary because they immediately follow a write in the UI flow. Ugly, but it's twenty lines and no infrastructure.

You've outgrown those when: (a) primary CPU is above ~60% and the profile is dominated by SELECTs, and (b) you can't statically classify reads as tolerant or not, because the same endpoint is stale-sensitive for the person who just edited and stale-tolerant for everyone else. That second condition is the real signal. If it isn't true, per-endpoint routing is still the right answer and the LSN token machinery is complexity you'll be maintaining for nothing.

Key takeaway: Don't choose globally between primary and replica — capture the write position on every commit and let each read decide, per request, whether the replica is caught up enough.

Real-world challenge

After introducing a read replica, your CI integration suite went flaky. A test POSTs /orders, gets 201, then GETs /orders and asserts the new order is in the list — it fails maybe 1 run in 8. Nothing has changed in the application code except the database URL for reads. Under load in production, the same endpoint shows the same bug for real users, and it gets much worse during nightly bulk imports.

Diagnose

  1. Confirm it's replication lag, not caching. On the replica:
SELECT now() - pg_last_xact_replay_timestamp() AS lag_time,
       pg_last_wal_replay_lsn() AS replayed;

If lag_time spikes into hundreds of milliseconds (and into seconds during bulk imports), that's your window.

  1. Check whether the failing read even could be correct: the POST returned before the replica applied the WAL record. The write is durable, just not visible where you looked.

Fix

Return the commit position from the write path and gate the read on it:

-- after COMMIT, on the primary connection
SELECT pg_current_wal_insert_lsn();

Store that LSN in the session/cookie/response header, and in the read router compare it against each replica's pg_last_wal_replay_lsn() (cached, polled every ~200ms). If no replica is caught up, use the primary for that request only.

Also