Database Replication Lag Monitoring: Real-Time Alerting for High Availability
Real-time replication lag monitoring and alerting systems that prevent stale reads and maintain consistency in distributed databases.

Replication lag is the silent availability killer. Your database looks healthy — primary is responding, replicas are online, connections are stable. But a replica that's 30 seconds behind the primary is serving stale data to every user routed to it. In a financial system, that means showing incorrect balances. In an e-commerce system, that means overselling inventory. We learned this the hard way after a replication lag spike caused $34K in inventory discrepancies. If you are using RDS Proxy in production, lag monitoring becomes even more critical since the proxy transparently routes queries to replicas.
This article covers the monitoring and alerting infrastructure we built to detect, measure, and respond to replication lag across PostgreSQL, MySQL, and DynamoDB Global Tables.
Why Replication Lag Matters
In a read-replica architecture, you route reads to replicas to reduce primary load. This only works if replicas are reasonably current. The acceptable lag depends on your use case:
| Use Case | Max Acceptable Lag | Consequence of Violation |
|---|---|---|
| Financial balances | 0 (use primary) | Incorrect balance display |
| Inventory counts | 2 seconds | Overselling |
| User profiles | 30 seconds | Stale profile data |
| Analytics dashboards | 5 minutes | Slightly outdated metrics |
| Search indexes | 15 minutes | Missing recent content |
The problem: most teams only discover lag when users report inconsistencies. By then, damage is done.
Architecture: Lag Measurement System
Our monitoring system measures lag from three angles:
- Heartbeat-based measurement: Write timestamps to primary, read from replicas, measure difference
- Native metrics collection: Database-specific lag metrics (pg_stat_replication, SHOW SLAVE STATUS)
- Application-level verification: Write-then-read consistency checks from the application layer
import asyncio
import time
from dataclasses import dataclass, field
from datetime import datetime, timezone
from typing import Optional
import asyncpg
import aiomysql
@dataclass
class LagMeasurement:
replica_id: str
lag_seconds: float
measurement_method: str
timestamp: datetime
is_healthy: bool
threshold_seconds: float
class HeartbeatLagMonitor:
"""Measure replication lag via heartbeat writes and reads."""
def __init__(
self,
primary_dsn: str,
replica_dsns: dict[str, str],
heartbeat_interval: float = 1.0,
table_name: str = "_replication_heartbeat"
):
self.primary_dsn = primary_dsn
self.replica_dsns = replica_dsns
self.heartbeat_interval = heartbeat_interval
self.table_name = table_name
self._primary_pool: Optional[asyncpg.Pool] = None
self._replica_pools: dict[str, asyncpg.Pool] = {}
async def initialize(self) -> None:
"""Create connection pools and heartbeat table."""
self._primary_pool = await asyncpg.create_pool(self.primary_dsn, min_size=1, max_size=2)
for replica_id, dsn in self.replica_dsns.items():
self._replica_pools[replica_id] = await asyncpg.create_pool(dsn, min_size=1, max_size=2)
# Create heartbeat table on primary
async with self._primary_pool.acquire() as conn:
await conn.execute(f"""
CREATE TABLE IF NOT EXISTS {self.table_name} (
id INTEGER PRIMARY KEY DEFAULT 1,
heartbeat_ts TIMESTAMP WITH TIME ZONE NOT NULL,
sequence_num BIGINT NOT NULL DEFAULT 0
)
""")
await conn.execute(f"""
INSERT INTO {self.table_name} (id, heartbeat_ts, sequence_num)
VALUES (1, NOW(), 0)
ON CONFLICT (id) DO NOTHING
""")
async def write_heartbeat(self) -> float:
"""Write current timestamp to primary. Returns write timestamp."""
ts = datetime.now(timezone.utc)
async with self._primary_pool.acquire() as conn:
await conn.execute(f"""
UPDATE {self.table_name}
SET heartbeat_ts = $1, sequence_num = sequence_num + 1
WHERE id = 1
""", ts)
return ts.timestamp()
async def read_replica_lag(self, replica_id: str) -> LagMeasurement:
"""Read heartbeat from replica and calculate lag."""
pool = self._replica_pools[replica_id]
async with pool.acquire() as conn:
row = await conn.fetchrow(f"""
SELECT heartbeat_ts, sequence_num FROM {self.table_name} WHERE id = 1
""")
if row is None:
return LagMeasurement(
replica_id=replica_id,
lag_seconds=float('inf'),
measurement_method='heartbeat',
timestamp=datetime.now(timezone.utc),
is_healthy=False,
threshold_seconds=5.0
)
replica_ts = row['heartbeat_ts'].timestamp()
current_ts = time.time()
lag = current_ts - replica_ts
return LagMeasurement(
replica_id=replica_id,
lag_seconds=max(0, lag),
measurement_method='heartbeat',
timestamp=datetime.now(timezone.utc),
is_healthy=lag < 5.0,
threshold_seconds=5.0
)
async def measure_all_replicas(self) -> list[LagMeasurement]:
"""Measure lag on all replicas concurrently."""
await self.write_heartbeat()
await asyncio.sleep(0.1) # Brief delay to allow propagation
tasks = [
self.read_replica_lag(replica_id)
for replica_id in self._replica_pools
]
return await asyncio.gather(*tasks)
Alerting Pipeline
Raw measurements feed into a multi-tier alerting system:
from enum import Enum
from typing import Callable
class AlertSeverity(Enum):
INFO = "info"
WARNING = "warning"
CRITICAL = "critical"
EMERGENCY = "emergency"
@dataclass
class AlertRule:
name: str
threshold_seconds: float
duration_seconds: float # Must exceed threshold for this long
severity: AlertSeverity
action: str
class ReplicationAlertManager:
"""Multi-tier alerting for replication lag."""
RULES = [
AlertRule("lag_warning", 5.0, 30, AlertSeverity.WARNING, "slack_notify"),
AlertRule("lag_critical", 15.0, 60, AlertSeverity.CRITICAL, "page_oncall"),
AlertRule("lag_emergency", 60.0, 10, AlertSeverity.EMERGENCY, "remove_replica"),
]
def __init__(self):
self._violation_start: dict[str, float] = {}
self._active_alerts: dict[str, AlertRule] = {}
def evaluate(self, measurement: LagMeasurement) -> Optional[AlertRule]:
"""Evaluate measurement against alert rules."""
key = f"{measurement.replica_id}"
# Find highest severity rule that's violated
triggered_rule = None
for rule in sorted(self.RULES, key=lambda r: r.threshold_seconds, reverse=True):
if measurement.lag_seconds >= rule.threshold_seconds:
triggered_rule = rule
break
if triggered_rule is None:
# Clear any existing violation tracking
self._violation_start.pop(key, None)
self._active_alerts.pop(key, None)
return None
# Track duration of violation
if key not in self._violation_start:
self._violation_start[key] = time.time()
elapsed = time.time() - self._violation_start[key]
if elapsed >= triggered_rule.duration_seconds:
if key not in self._active_alerts or self._active_alerts[key] != triggered_rule:
self._active_alerts[key] = triggered_rule
return triggered_rule
return None
def should_remove_replica(self, replica_id: str) -> bool:
"""Check if replica should be removed from read pool."""
key = replica_id
alert = self._active_alerts.get(key)
return alert is not None and alert.severity == AlertSeverity.EMERGENCY
Automatic Replica Removal
When lag exceeds emergency thresholds, we automatically remove the replica from the read pool:
class ReadPoolManager:
"""Manages read replica pool with automatic removal on lag."""
def __init__(self, replicas: list[str], alert_manager: ReplicationAlertManager):
self.all_replicas = set(replicas)
self.healthy_replicas = set(replicas)
self.alert_manager = alert_manager
def get_read_replica(self) -> Optional[str]:
"""Get a healthy replica for read routing."""
if not self.healthy_replicas:
return None # Fallback to primary
# Round-robin among healthy replicas
replica = min(self.healthy_replicas)
return replica
def update_health(self, measurements: list[LagMeasurement]) -> dict:
"""Update replica health based on lag measurements."""
changes = {'removed': [], 'restored': []}
for measurement in measurements:
if self.alert_manager.should_remove_replica(measurement.replica_id):
if measurement.replica_id in self.healthy_replicas:
self.healthy_replicas.discard(measurement.replica_id)
changes['removed'].append(measurement.replica_id)
else:
if measurement.replica_id not in self.healthy_replicas:
if measurement.replica_id in self.all_replicas:
self.healthy_replicas.add(measurement.replica_id)
changes['restored'].append(measurement.replica_id)
return changes
DynamoDB Global Tables Lag
DynamoDB Global Tables don't expose lag metrics natively. We measure it using the same heartbeat approach, writing to the primary region and reading from replicas. When lag grows unacceptable, database sharding may be the next step to distribute write load:
| Region Pair | p50 Lag | p95 Lag | p99 Lag |
|---|---|---|---|
| us-east-1 → eu-west-1 | 0.8s | 1.4s | 2.8s |
| us-east-1 → ap-southeast-1 | 1.2s | 2.1s | 4.2s |
| eu-west-1 → ap-southeast-1 | 1.0s | 1.8s | 3.5s |
Production Results
After deploying comprehensive lag monitoring:
- Mean time to detect lag spikes: from 12 minutes (user reports) to 5 seconds (automated)
- Replica removal time on lag emergency: 10 seconds (automated) vs 3 minutes (manual)
- Stale read incidents per month: from 4-6 to 0 in the last 8 months
- Inventory discrepancy cost: from $34K one-time to $0
Key Takeaways
- Heartbeat monitoring is essential. Native metrics tell you lag exists; heartbeats tell you the actual user-visible impact.
- Automatic removal saves data. A lagging replica serving stale reads is worse than no replica at all.
- Alert on duration, not spikes. Brief lag spikes are normal during checkpoint operations. Alert when lag persists.
- Different data needs different thresholds. Financial data needs sub-second lag detection. Content search can tolerate minutes.
- Test your failover path. Removing a replica under lag must be tested regularly. An untested automatic removal is a liability.
Replication lag monitoring is cheap insurance against expensive consistency violations. Build it before you need it. For related database infrastructure patterns, see RDS Proxy in production for connection pooling that prevents replica overload, and Kubernetes cost optimization for right-sizing the infrastructure that runs your monitoring systems.
Recommended reading

Per-Team Cost Allocation in Shared Kubernetes Clusters: From Chaos to Clarity
Implementing accurate per-namespace cost allocation in multi-tenant Kubernetes clusters, covering request vs. usage attribution, shared resource amortization, and building showback dashboards that drive accountability.

Measuring and Eliminating Toil: From 40% to 12% of Engineering Time
A systematic approach to identifying, measuring, and automating toil—the repetitive operational work that scales linearly with service growth and prevents engineers from doing creative work.

Serverless Postgres in Production: Branching, Scale-to-Zero, and the End of Database Provisioning
Running Neon serverless Postgres in production for 8 months — covering database branching workflows, scale-to-zero economics, connection pooling, and migration from RDS.

Comments
No comments yet. Be the first to share your thoughts.