Replica lag is one of those distributed-systems facts everyone accepts in the abstract and nobody plans for in the code. You write to the primary, the write succeeds, a user’s browser immediately reads from a replica that hasn’t caught up — and the data they just saved appears to have vanished. The database is fine. The system simply never told anyone which reads are safe to make where.
These failures have names. They are the read-your-writes anomaly, the monotonic-read anomaly, and their cousins — the session guarantees first formalized in the mid-1990s in the Bayou project and later given a rigorous model by Mahajan, Alvisi, and Dahlin. Every replicated store makes some version of these promises; the engineering work is understanding exactly which ones your database gives you, what they cost, and how to hold your application to them. This post walks through the guarantee hierarchy and the practical patterns for enforcing it in real systems.
The Anomalies, Concretely
Imagine a replicated key-value store with one writable primary and asynchronous read replicas. Four anomalies account for nearly all user-visible consistency bugs:
- Stale read — you read a value and get an older version than what was actually written, with no write of your own involved.
- Read-your-writes violation — you update your profile, the page reloads, and the old value is back. The classic user-report-generating bug.
- Monotonic-read violation — you read a value, refresh, and see an older value than before. Time appears to run backward. This happens when successive reads land on different replicas with different lag.
- Monotonic-write violation — two writes from the same session get applied out of order on some replica, so the effect of the second write is lost.
The strongest guarantee — linearizability — eliminates all of them by making the store behave as if there were exactly one copy of the data and every operation is atomic in real time. It is also the most expensive: reads must coordinate with a quorum or the leader, which adds latency and makes the read path unavailable during partitions. Everything below it is a trade of consistency for speed.
The Session Guarantee Hierarchy
The insight from the session-guarantees literature is that clients don’t usually need global linearizability — they need guarantees scoped to their own session. The standard model orders guarantees by strength:
Read-your-writes (RYW) promises that a session sees its own writes. If you just wrote value X, your next read of that key returns X or newer. Monotonic reads promises that once a session has seen a value, it never sees an older one. Writes-follow-reads and monotonic writes order writes within a session relative to earlier reads and writes. Causal consistency bundles all of these: if operation A might have caused operation B — across sessions, through any chain of reads and writes — every replica applies A before B.
The remarkable result from the formal work is that causal consistency is the strongest guarantee achievable while staying available during partitions — every replica can accept reads and writes without contacting a majority. That is exactly the tier most web applications actually need: a user’s own actions must look immediate, while cross-user convergence can lag by milliseconds without anyone noticing.
How Real Systems Express the Guarantees
The guarantees are implemented with one recurring trick: track a version of the database state with each session, and make replicas wait until they’ve caught up to that version.
MongoDB is the most explicit example. Every operation on a session updates a per-session operation time; with readConcern: "majority" and causally consistent sessions, subsequent reads on any node wait until that node has applied everything the session has seen. Piggybacked metadata in responses keeps the token current.
CouchDB takes a different angle: with eventual consistency as the default, its replication protocol makes conflicts visible and deterministic to resolve, and applications that need stricter behavior implement session checks at the document level.
DynamoDB makes the trade-off a per-read choice: eventually consistent reads are cheaper and faster, strongly consistent reads cost double the capacity but always reflect the latest successful write.
PostgreSQL streaming replication sits at the classic end of the spectrum: replicas are plain lagging copies. The guarantees don’t exist until you build them — which is the exercise for the rest of this post.
Pattern 1: Read-Your-Writes With a Version Token
The workhorse pattern. When a session writes, capture the write’s position — on a primary, a monotonically increasing sequence number; in Postgres, the write-ahead-log location. Hand the position back to the client as an opaque token. On subsequent reads, the replica must have replayed at least that position before answering.
// PostgreSQL: capture the write position after a write on the primary.
func (s *Store) CreateOrder(ctx context.Context, o Order) (string, error) {
var lsnStr string
err := s.primary.QueryRowContext(ctx,
`INSERT INTO orders (user_id, total) VALUES ($1, $2)
RETURNING pg_current_wal_lsn()::text`,
o.UserID, o.Total,
).Scan(&lsnStr)
if err != nil {
return "", err
}
// lsnStr, e.g. "7D9/F31A8B60" — return it to the client as a token
return lsnStr, nil
}
The read side compares the token against the replica’s replay position and either waits or falls back:
// Read from a replica only after it has caught up to the session's LSN.
func (s *Store) GetOrder(ctx context.Context, token, id string) (Order, error) {
var order Order
err := s.waitForLSN(ctx, token)
if err != nil {
// replica didn't catch up in time — fall back to the primary
return s.readPrimary(ctx, id)
}
err = s.replica.QueryRowContext(ctx,
`SELECT user_id, total FROM orders WHERE id = $1`, id,
).Scan(&order.UserID, &order.Total)
return order, err
}
func (s *Store) waitForLSN(ctx context.Context, token string) error {
ctx, cancel := context.WithTimeout(ctx, 50*time.Millisecond)
defer cancel()
for {
var replayed string
err := s.replica.QueryRowContext(ctx,
`SELECT pg_last_wal_replay_lsn()::text`,
).Scan(&replayed)
if err != nil {
return err
}
// pg_wal_lsn_diff returns bytes between two LSNs
var behind int64
err = s.replica.QueryRowContext(ctx,
`SELECT pg_wal_lsn_diff($1::pg_lsn, pg_last_wal_replay_lsn())`,
token,
).Scan(&behind)
if err != nil {
return err
}
if behind <= 0 {
return nil // replica has caught up to the session's write
}
time.Sleep(2 * time.Millisecond)
}
}
Postgres gives you all the primitives here: pg_current_wal_lsn() reads the primary’s current position, pg_last_wal_replay_lsn() reads how far a replica has replayed, and pg_wal_lsn_diff measures the distance between them. Note the comparisons happen via pg_wal_lsn_diff rather than string comparison — LSNs are not lexically ordered across boundaries.
Pattern 2: Sticky Sessions at the Routing Layer
The simpler sibling: after a write, route all of that session’s reads to the primary for a fixed window. No version tokens, no lag checks — just a flag in the session store (a short-TTL Redis key, a signed cookie) that your read path checks before choosing a replica.
It’s coarse but effective. Choose the window as a multiple of your worst-case replication lag: if p99 lag is 200 ms, a 5-second sticky window leaves a comfortable margin. The costs are modest — some primary load from post-write reads, and a rare stale read if lag spikes past the window. For many applications it removes an entire class of bug for one line of routing logic.
Pattern 3: Monotonic Reads Across Replicas
Read-your-writes doesn’t stop time running backward: a session might read from a caught-up replica, then get routed to a lagging one. Monotonic reads require remembering the highest version the session has ever observed and holding reads until every future replica reaches it.
The mechanism composes with Pattern 1: alongside your write LSN, store the maximum LSN seen on any read, and apply the same waitForLSN gate before every replica read. In causally consistent stores like MongoDB, this bookkeeping is what the session token does implicitly — the token advances with every operation, reads and writes alike.
Choosing Your Tier
- Strongly consistent / linearizable — financial balances, inventory decrements, anything with a correctness invariant. Pay the latency; read from the primary or a quorum.
- Read-your-writes + monotonic reads — user-facing profiles, settings, drafts. The version-token or sticky-session pattern above.
- Eventually consistent — feeds, counts, recommendations, analytics. Users don’t hold their own writes against these, so take the speed.
The failure mode to avoid is choosing tier by infrastructure habit rather than by read semantics. Audit your read paths: for each, ask “does the user who caused this read ever observe it?” If yes, that path needs a session guarantee, and now you have three concrete ways to give it one.
Wrapping Up
Session guarantees occupy the sweet spot in the consistency spectrum: strong enough that users never see impossible behavior, cheap enough to keep replication async and the system available. The formal model is three decades old, and modern stores have converged on the same implementation — session-scoped version tokens checked against replica progress. Whether your database does it for you or you build it with pg_wal_lsn_diff and a sticky flag, the contract is the same: writes are visible to their author, and observed values never go backward. For a deeper formal treatment, the Mahajan–Alvisi–Dahlin paper “Consistency, Availability, and Convergence” is the canonical starting point.