replication lag
Replication Lag
Learn why replicas can temporarily remain behind the latest accepted writes, how lag differs across leader-follower, multi-leader, and leaderless replication, how stale reads affect users, and how read-your-writes, monotonic reads, consistent-prefix reads, freshness-aware routing, replication positions, repair, monitoring, failover, and application design control the resulting risks.
Introduction
Replication maintains copies of data on multiple machines. These copies can improve availability, read capacity, fault tolerance, and geographic access.
However, a data change cannot appear on every replica instantaneously. The change must be transferred across a network, received by each replica, persisted according to the database policy, and applied to the replica's local state.
Client writes new data
|
v
Leader commits the write
|
v
Replication record is transmitted
|
v
Follower receives the record
|
v
Follower persists and applies it
|
v
Follower can return the new value
The delay between the authoritative write and its visibility on another replica is called replication lag.
Core idea: Replication lag means that replicas can temporarily represent different points in data history. An application must decide which reads can tolerate an older version and which reads must wait for or route to a sufficiently current replica.
Prerequisites
| # | Prerequisite | Why It Is Needed |
|---|---|---|
| 1 | Leader-follower replication | Writes commonly reach one leader before being copied to followers. |
| 2 | Multi-leader and leaderless replication | Lag can appear between several writing locations or among independently updated replicas. |
| 3 | Synchronous and asynchronous replication | Write acknowledgment timing determines which replicas are current when the client receives success. |
| 4 | Consistency models | Applications need clear guarantees about the versions clients are allowed to observe. |
| 5 | Quorums and versioning | Leaderless reads can compare versions returned by several replicas. |
| 6 | Failover | Promoting a lagging replica can omit recently accepted writes. |
What Is Replication Lag?
Replication lag is the difference between the latest accepted state of a dataset and the state currently available on another replica.
Leader:
Latest replication position = 15,800
Follower A:
Applied position = 15,800
Follower B:
Applied position = 15,720
Follower C:
Applied position = 15,100
Follower A is current at the measured position. Followers B and C remain behind and can return older values for changes they have not yet applied.
Lag can be represented using:
- Time behind the source
- Bytes of unapplied replication data
- Log-sequence position difference
- Number of pending records or operations
- Transaction or commit position difference
- Oldest unapplied event age
Replication-lag Stages
| Stage | Meaning | Possible Bottleneck |
|---|---|---|
| Generation | The source creates replication records from committed changes | High write volume or large transactions |
| Transmission | Replication records travel to the replica | Network latency, congestion, interruption, or bandwidth |
| Receipt | The replica receives replication data | Receiver CPU, buffers, network processing, or connection health |
| Persistence | The replica records the received change durably according to policy | Storage latency or throughput |
| Application | The replica applies the change to its queryable state | CPU, locks, I/O, transaction ordering, or apply-worker capacity |
A replica can have received a change without having made it visible to ordinary reads. Monitor the stage that matches the database's actual replication architecture.
Synchronous vs Asynchronous Replication
Synchronous Replication
Client sends write
|
v
Leader commits or prepares write
|
v
Required follower acknowledges
|
v
Leader returns success
Synchronous replication places a configured replica acknowledgment in the client write path. It can provide stronger durability or visibility guarantees, depending on what the acknowledgment means.
The trade-offs can include:
- Higher write latency
- Dependence on replica and network response time
- Reduced write availability when required replicas are unavailable
Asynchronous Replication
Client sends write
|
v
Leader commits write
|
v
Leader returns success
|
v
Followers receive and apply
the change later
Asynchronous replication allows the leader to acknowledge a write before followers confirm it. This can reduce write latency and preserve leader availability during follower problems, but followers can return stale data.
Acknowledgment rule: Verify what a successful write means in the selected database. A follower receiving a change, persisting it, and applying it are different stages and can provide different guarantees.
Lag across Replication Topologies
Leader-Follower
Writes
|
v
Leader
|
+-- Replication stream -> Follower A
+-- Replication stream -> Follower B
+-- Replication stream -> Follower C
Followers can be at different replication positions. A read from a lagging follower can omit recent leader writes.
Multi-Leader
Region A Leader
|
| cross-region replication
v
Region B Leader
Region B Leader
|
| cross-region replication
v
Region A Leader
Each leader can accept local writes before remote leaders receive them. Cross-location lag can therefore combine stale reads with concurrent-write conflicts.
Leaderless
Coordinator
|
+-- Replica A: Version 12
+-- Replica B: Version 12
+-- Replica C: Version 10
Leaderless systems can observe stale or concurrent versions across replicas. Quorum reads, version comparison, read repair, hinted handoff, and anti-entropy are used according to the database's design.
User-visible Stale Read
1. User updates course progress.
2. Write is committed on the leader.
3. Application returns success.
4. User immediately refreshes the page.
5. Read is sent to a lagging follower.
6. Follower returns the previous progress value.
The database can be operating according to its asynchronous replication contract while the application still appears incorrect to the user.
Read-Your-Writes Consistency
Read-your-writes consistency means that after a client successfully writes a value, that client will not subsequently read a version older than the accepted write.
Client writes:
Progress = 72%
Required subsequent observation:
Progress >= Version containing 72%
Possible implementation patterns include:
- Read from the leader after a client write
- Read from the leader for data owned or modified by the current user
- Track the write's replication position and select a replica that has reached it
- Wait for an eligible replica to catch up within a bounded deadline
- Use a database-provided consistency level that meets the requirement
Position-based Read Routing
The application can capture a replication position associated with a successful write and require subsequent reads to use a replica that has applied at least that position.
Write commits at:
Position 8,500
Subsequent read requires:
Applied position >= 8,500
Replica A:
Position 8,420
Not eligible
Replica B:
Position 8,510
Eligible
Conceptual Client Context
{
"resourceId": "course-progress-1042-42",
"minimumRequiredVersion": 8500
}
The version or position can remain opaque to the public client. The server should validate any client-carried consistency context rather than trusting an arbitrary client-selected database position.
Freshness-aware Read Router
Read request
|
v
Does request require
a minimum version?
|
+-- No:
| route to an eligible replica
|
+-- Yes:
compare required version
with replica positions
|
+-- Current replica exists:
| route to replica
|
+-- No replica caught up:
route to leader,
wait within deadline,
or fail according
to the contract
Freshness rule: Do not require every read to use the freshest possible copy when only selected operations need that guarantee. Classify read paths according to their actual freshness requirements.
Monotonic Reads
Monotonic reads ensure that after a client has observed one version, later reads do not return an older version.
First read from Follower A:
Course progress = Version 14
Second read from Follower B:
Course progress = Version 11
Problem:
The client observes data
moving backward.
Possible directions include:
- Route one client to a replica that is at least as current as the previous read
- Carry a minimum observed version in trusted session context
- Use affinity to an eligible replica where appropriate
- Route to the leader when no follower satisfies the minimum version
Consistent-prefix Reads
Consistent-prefix reads preserve a meaningful order between causally related writes.
Write 1:
Create discussion question
Write 2:
Create reply to that question
Invalid observation:
Reply appears before
the original question
Related data can be routed through the same partition, replication stream, or ordering mechanism when causal ordering must remain visible.
Three Important Guarantees
| Guarantee | Client Expectation | Example |
|---|---|---|
| Read-your-writes | A client sees that client's own accepted update | A changed profile remains changed after refresh |
| Monotonic reads | Later reads do not return older versions than previously observed | Progress does not move from 72% back to 65% |
| Consistent prefix | Causally related writes appear in a sensible order | A reply does not appear before its question |
Replication Lag and Failover
If the leader fails, the system can promote a follower. A lagging follower can be missing recent writes acknowledged by the old leader.
Old leader:
Committed through Position 10,000
Follower selected for promotion:
Applied through Position 9,970
Potential missing range:
Positions 9,971 to 10,000
The exact outcome depends on replication acknowledgment, failover, durable logging, election, and recovery rules.
Failover design should define:
- Which followers are eligible for promotion
- How replication freshness is compared
- Whether acknowledged writes can be lost
- How the old leader is fenced
- How divergent or missing writes are handled
- When traffic can resume safely
Common Causes of Replication Lag
| Cause | How Lag Develops |
|---|---|
| High source write volume | The source generates replication records faster than the replica applies them |
| Large transaction | A large volume of changes must be transferred and applied |
| Network interruption | The replica temporarily stops receiving changes |
| Network congestion | Replication transfer cannot keep pace with generated traffic |
| Replica storage bottleneck | The replica cannot persist or apply changes quickly enough |
| Replica CPU pressure | Replication processing competes for compute resources |
| Heavy replica reads | Analytical or long-running queries compete with replication work |
| Lock contention | Replication apply waits behind another operation |
| Schema or index change | Data-definition work or additional maintenance affects apply throughput |
| Inactive subscriber or slot | Replication records accumulate because a consumer is not advancing |
Write Rate vs Apply Rate
Let \(G\) be the rate at which the source generates replication work and \(A\) be the rate at which the replica applies it.
The replica can catch up when:
\[ A > G \]
Lag continues growing when:
\[ G > A \]
When generation stops or decreases, a simplified catch-up estimate is:
\[ CatchUpTime \approx \frac{ ReplicationBacklog }{ ApplyRate - GenerationRate } \]
This is a planning approximation. Transaction dependencies, parallelism, retries, network variation, and database implementation affect actual recovery.
Replica Read Workload
Read replicas are commonly used for reporting, search indexing, exports, and analytical queries. These workloads can compete with replication apply.
Replica resources
|
+-- Replication receiver
+-- Replication apply
+-- Application reads
+-- Reporting queries
+-- Backup operations
+-- Maintenance tasks
Separate workload pools, query limits, dedicated analytical systems, or replica-specific routing can protect replication capacity.
Long-running Queries
A long-running replica query can hold resources or maintain an older database view. Depending on the database, this can interfere with cleanup, replication apply, or storage processing.
Review:
- Query duration
- Locks and transaction state
- Temporary storage
- CPU and I/O consumption
- Cancellation and statement timeouts
- Whether reporting should use another data path
Replication Slots and Log Retention
Some replication systems retain source log data until a subscriber confirms that the data is no longer required.
Source generates replication logs
|
v
Subscriber stops consuming
|
v
Required logs cannot be removed
|
v
Retained log volume grows
|
v
Source storage pressure increases
An inactive replication slot or subscriber can therefore create both lag and source-storage growth. Monitor subscriber state, retained log volume, and recovery progress.
Conceptual PostgreSQL Lag Inspection
SELECT
application_name,
client_addr,
state,
sent_lsn,
write_lsn,
flush_lsn,
replay_lsn,
write_lag,
flush_lag,
replay_lag
FROM pg_stat_replication;
The availability and meaning of fields depend on the PostgreSQL version, replication mode, and current activity. Interpret the status using the selected platform's official documentation.
Conceptual Position Difference
SELECT
application_name,
pg_wal_lsn_diff(
sent_lsn,
replay_lsn
) AS pending_replay_bytes
FROM pg_stat_replication;
A byte difference shows outstanding replication volume. It does not directly equal a fixed catch-up duration because apply throughput varies.
Lag Measured in Time vs Bytes
| Measurement | Useful For | Limitation |
|---|---|---|
| Time lag | Understanding data freshness in user-facing terms | Can be unavailable or misleading during idle periods |
| Byte lag | Estimating outstanding replication volume | The same byte count can require different apply times |
| Position lag | Comparing source and replica progress | Requires position-aware interpretation |
| Transaction lag | Identifying unapplied transaction count | Transactions can vary greatly in size and cost |
Read-routing Strategies
| Strategy | Benefit | Trade-off |
|---|---|---|
| Read everything from leader | Simple freshness direction | Reduces read-scaling benefit and increases leader load |
| Read everything from replicas | Maximum read distribution | Freshness-sensitive paths can return stale data |
| Route freshness-sensitive reads to leader | Protects selected critical paths | Requires operation classification |
| Position-aware replica routing | Uses replicas without violating a minimum version | Requires replica progress tracking |
| Wait for replica catch-up | Can preserve replica reads | Adds latency and requires a bounded deadline |
| Session or client affinity | Can support monotonic reads from one replica | The chosen replica can become unavailable or remain behind |
Learning-platform Examples
| Data | Lag Risk | Possible Direction |
|---|---|---|
| Public course description | A recent editorial update appears later on a replica | Document acceptable freshness and use replica reads |
| Learner progress | Progress appears to move backward after a refresh | Provide read-your-writes and monotonic-read behaviour |
| Enrollment confirmation | A confirmed enrollment temporarily appears missing | Read from an authoritative or sufficiently current source |
| Search index | A newly published course does not immediately appear | Expose or document indexing freshness when relevant |
| Recommendation model | Recommendations use slightly older activity data | Accept measured staleness when it does not affect correctness |
| Certificate issuance | A replica omits a newly issued or revoked certificate | Use a freshness policy appropriate for authoritative status |
Do Not Use Stale Reads for Strict Invariants
Some decisions must use sufficiently current authoritative data.
Examples include:
- Checking whether an enrollment already exists
- Validating an active session after revocation
- Checking available inventory before reservation
- Confirming a payment state
- Enforcing a unique business reference
- Authorizing access after a permission removal
Reading a stale replica and then writing to the leader can produce a time-of-check to time-of-use problem. Enforce critical invariants through transactions, conditional writes, uniqueness constraints, locking, or another approved coordination mechanism at the authoritative write path.
Retry and Duplicate-operation Risk
Client submits enrollment
|
v
Leader commits enrollment
|
v
Client reads lagging replica
|
v
Enrollment appears missing
|
v
Client submits again
Protect state-changing operations with idempotency rather than relying on an immediate replica read to determine whether the original request succeeded.
Idempotent Request
POST /api/enrollments HTTP/1.1
Host: api.example.com
Authorization: Bearer access-token
Idempotency-Key: unique-operation-key
Content-Type: application/json
{
"courseId": 42
}
Read Repair and Anti-Entropy
Leaderless and selected replicated systems can use repair mechanisms to bring stale replicas up to date.
Read Repair
Read several replicas
|
v
Detect an older version
|
v
Return selected current value
|
v
Update stale replica
Anti-Entropy
Background replica comparison
|
v
Identify missing or divergent ranges
|
v
Transfer required versions
|
v
Verify convergence
Read repair helps frequently accessed data. Background repair is needed for data that might not be read often enough to trigger repair.
Conceptual Freshness Policy
replicaReads:
default:
source: eligible-replica
freshness: documented-eventual
readYourWrites:
minimumRequiredVersion: trusted-session-context
fallback:
- current-replica
- leader
- bounded-wait
- controlled-failure
monotonicReads:
trackMinimumObservedVersion: true
criticalOperations:
source: authoritative-path
allowStaleRead: false
routing:
maximumApprovedLag: workload-specific
monitoring:
replicaPositions: enabled
lagByReplica: enabled
staleReadEvents: enabled
This is conceptual configuration. Use the actual consistency, routing, and replication features provided by the selected database.
Security Impact
Replication lag affects security-sensitive state as well as business data.
Examples include:
- A revoked session remaining visible as active on a lagging replica
- A removed role remaining visible temporarily
- A disabled account appearing enabled
- A revoked API credential appearing valid
- An outdated access-control list permitting previous access
Security decisions should use a consistency level appropriate to the risk. Do not direct revocation-sensitive authorization checks to an arbitrarily stale replica.
Observability
Useful replication-lag metrics include:
- Lag time by replica
- Outstanding replication bytes
- Source and replica positions
- Replication receive rate
- Replication apply rate
- Pending replication operations
- Oldest unapplied transaction age
- Replication connection state
- Write, flush, and replay progress where available
- Stale-read count
- Read reroutes to leader
- Reads waiting for replica catch-up
- Read-repair and anti-entropy backlog
- Retained log or replication-slot volume
- Replica CPU, memory, I/O, and query load
Structured Lag Event
{
"replica": "course-read-replica-2",
"sourcePosition": "15800",
"appliedPosition": "15720",
"positionDifference": 80,
"replicationState": "catching-up",
"routingEligibility": "restricted"
}
Do not include credentials, connection secrets, or sensitive record values in replication-monitoring events.
Alert Conditions
Alert when:
- Lag exceeds the approved freshness objective
- Lag continues increasing rather than stabilizing
- Replication receive or apply stops
- A replica remains in catch-up state unexpectedly
- Retained replication logs threaten source storage capacity
- A replica repeatedly serves stale data beyond policy
- No replica satisfies freshness-sensitive reads
- Leader traffic increases because replicas remain ineligible
- Read repair or anti-entropy backlog grows
- Failover candidates remain substantially behind
- Replica CPU, I/O, or storage becomes saturated
- Lag increases after a schema or deployment change
Troubleshooting Workflow
- Identify the replication topology and affected replicas.
- Determine whether lag is measured in time, bytes, positions, or transactions.
- Compare source generation rate with replica receive and apply rates.
- Check replication connection and subscriber state.
- Check network latency, interruption, congestion, and bandwidth.
- Check replica CPU, memory, disk latency, and storage throughput.
- Check long-running queries and lock contention.
- Check large transactions and bulk changes.
- Check recent schema, index, deployment, and configuration changes.
- Check replication slots and retained logs.
- Reduce or isolate nonessential replica workloads where safe.
- Confirm that freshness-sensitive reads avoid the lagging replica.
- Estimate whether the replica can catch up at its current apply rate.
- Repair or rebuild the replica through the approved database procedure if necessary.
Common Replication-lag Mistakes
Assuming Every Replica Is Current
Asynchronous replication allows different replicas to represent different points in data history.
Reading from a Replica Immediately after Writing
The client can fail to observe the write that the leader already accepted.
Measuring Lag Using One Metric Only
Time, bytes, positions, and apply state reveal different parts of the replication process.
Using a Lagging Replica for Authorization
A revoked session, role, or permission can temporarily appear active.
Promoting the Most Reachable Replica
The selected replica can omit newer writes available on another candidate.
Sending Heavy Reporting Queries to Every Replica
Read workload can consume resources required by replication apply.
Ignoring Inactive Replication Slots
Required source logs can accumulate and consume growing storage.
Retrying after a Stale Read without Idempotency
The client can duplicate a write that already succeeded on the leader.
Routing All Reads to the Leader Permanently
This removes much of the read-scaling benefit instead of classifying freshness-sensitive operations.
Waiting Indefinitely for Catch-up
A freshness-sensitive read can consume resources without completing before the caller's deadline.
Assuming Lag Will Always Recover Automatically
A replica cannot catch up while its apply capacity remains below the source generation rate.
Skipping Failover Testing under Lag
The acknowledged-write loss and client-consistency behaviour remains unverified.
Recommended Test Cases
| Test | Expected Evidence |
|---|---|
| Normal asynchronous replication | Replicas converge within the documented operating objective |
| Read after write | The writer observes the accepted update |
| Monotonic reads | A client does not move from a newer version to an older version |
| Consistent-prefix read | Causally related records appear in the expected order |
| Replica read overload | Replication apply remains protected from uncontrolled query load |
| Large transaction | Lag and catch-up behaviour are measured and bounded |
| Network interruption | The replica reconnects and resumes from the correct position |
| Inactive subscriber | Retained-log growth is detected before source storage is threatened |
| Leader failure during lag | Failover follows the documented durability and promotion policy |
| Duplicate client retry | Idempotency prevents a second business effect |
| Security revocation | Authorization does not rely on an unacceptably stale replica |
| Replica recovery | The replica catches up or is rebuilt through a controlled process |
Replication-lag Best Practices
Recommended Practices
- Define acceptable freshness for every important read path.
- Measure lag by replica rather than only at cluster level.
- Track replication positions as well as time-based lag.
- Distinguish receive, persistence, and apply lag.
- Use read-your-writes for interactive user updates.
- Preserve monotonic reads where data moving backward would confuse users.
- Preserve causal ordering for related records.
- Route critical reads to an authoritative or sufficiently current source.
- Use position-aware replica selection where supported.
- Bound any wait for replica catch-up.
- Protect replication resources from uncontrolled analytical queries.
- Monitor replication slots, retained logs, and inactive subscribers.
- Use idempotency when stale reads can cause client retries.
- Do not enforce strict business invariants using stale reads.
- Protect security revocation paths from excessive staleness.
- Select failover candidates using freshness and safety criteria.
- Monitor read repair and anti-entropy where applicable.
- Test lag during bursts, large transactions, and network failures.
- Test failover while replicas are behind.
- Verify product-specific lag and acknowledgment semantics.
Practice Exercise
Design replication-lag handling for your online learning platform.
Requirements
- Identify which data is read from replicas.
- Define acceptable freshness for course catalog data.
- Provide read-your-writes after a learner updates progress.
- Prevent progress from moving backward across replica reads.
- Preserve ordering between discussion questions and replies.
- Keep enrollment validation on an authoritative path.
- Protect session and permission revocation from stale reads.
- Capture the write position required by a subsequent read.
- Route reads only to replicas meeting the required position.
- Define a bounded fallback when no replica has caught up.
- Add idempotency to enrollment and progress-submission requests.
- Monitor lag time, bytes, positions, and apply state.
- Test a large bulk course update.
- Test a disconnected replica.
- Test leader failover while followers are behind.
Freshness-policy Template
| Read Path | Freshness Requirement | Routing Direction | Failure Behaviour |
|---|---|---|---|
| Public course catalog | Documented bounded or eventual freshness | Eligible read replica | Serve approved data or report temporary unavailability |
| Own learner progress after update | Read-your-writes | Replica at required position or authoritative source | Bounded wait or controlled fallback |
| Enrollment confirmation | Must include accepted enrollment | Authoritative or sufficiently current source | Do not report enrollment missing from an arbitrary stale replica |
| Recommendations | Measured staleness can be acceptable | Read replica or projection | Omit optional recommendations when unavailable |
| Authorization revocation | Security-policy freshness | Approved authoritative consistency path | Deny access when current authorization cannot be validated |
Frequently Asked Questions
What is replication lag?
Replication lag is the difference between the latest accepted data and the data currently received, persisted, or applied by another replica.
Why does replication lag occur?
Replication requires network transfer, storage, and change application. Lag grows when a replica cannot receive or apply changes as quickly as the source produces them.
Does asynchronous replication cause stale reads?
It can. A leader can acknowledge a write before followers receive and apply it, allowing a follower to return an older version.
What is read-your-writes consistency?
It ensures that a client does not read a version older than that client's own previously accepted write.
What are monotonic reads?
Monotonic reads ensure that after a client observes one version, later reads do not return an older version.
What are consistent-prefix reads?
They ensure that causally ordered writes appear in a sensible order, such as showing a question before its reply.
Should every read go to the leader?
Not necessarily. Route only freshness-sensitive reads to the leader or a sufficiently current replica, while stale-tolerant reads can use ordinary replicas.
How can an application select a current replica?
The application or router can track the minimum required replication position and select a replica that has applied at least that position.
Can replication lag cause data loss?
Asynchronous lag can create a window in which acknowledged writes are not available on a promoted follower. The exact durability outcome depends on the database's acknowledgment and failover rules.
Can heavy reads increase replication lag?
Yes. Reporting or analytical queries can compete with replication for CPU, memory, storage, locks, and I/O on the replica.
Why are inactive replication slots dangerous?
A source can retain replication logs required by an inactive subscriber, allowing retained log storage to continue growing.
How should replication lag be monitored?
Monitor time, bytes, source and replica positions, receive and apply state, subscriber health, retained logs, and client-visible stale-read events.
Key Takeaway
Replication lag is the delay between an accepted write and that write becoming visible on another replica. It is especially important with asynchronous replication, where the source can acknowledge a write before followers receive and apply it. Lag can produce stale reads, read-your-writes violations, non-monotonic observations, incorrectly ordered related records, and unsafe failover to an outdated replica. Classify read paths by freshness requirement instead of treating every read equally. Use an authoritative source or position-aware replica for freshness-sensitive operations, bound any wait for catch-up, and use idempotency when stale reads can trigger duplicate writes. Monitor replication time, bytes, positions, receive and apply rates, replica resources, retained logs, and subscriber state. Finally, protect authorization and strict business invariants from stale data, and test large transactions, heavy replica reads, network interruptions, and leader failover while replicas are behind.