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.

#database#replication#monitoring#high-availability
Cover image for the article: Database Replication Lag Monitoring: Real-Time Alerting for High Availability

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 CaseMax Acceptable LagConsequence of Violation
Financial balances0 (use primary)Incorrect balance display
Inventory counts2 secondsOverselling
User profiles30 secondsStale profile data
Analytics dashboards5 minutesSlightly outdated metrics
Search indexes15 minutesMissing 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:

  1. Heartbeat-based measurement: Write timestamps to primary, read from replicas, measure difference
  2. Native metrics collection: Database-specific lag metrics (pg_stat_replication, SHOW SLAVE STATUS)
  3. 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 Pairp50 Lagp95 Lagp99 Lag
us-east-1 → eu-west-10.8s1.4s2.8s
us-east-1 → ap-southeast-11.2s2.1s4.2s
eu-west-1 → ap-southeast-11.0s1.8s3.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

  1. Heartbeat monitoring is essential. Native metrics tell you lag exists; heartbeats tell you the actual user-visible impact.
  2. Automatic removal saves data. A lagging replica serving stale reads is worse than no replica at all.
  3. Alert on duration, not spikes. Brief lag spikes are normal during checkpoint operations. Alert when lag persists.
  4. Different data needs different thresholds. Financial data needs sub-second lag detection. Content search can tolerate minutes.
  5. 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.

Comments

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