How would you choose replication-lag alerting thresholds for a fleet of read replicas, and what should happen automatically when a replica exceeds them?
answer
- Budgets: staleness (reads), recovery point (failover), retention (hard limit)
- Spike that drains vs sustained monotonic rise
- Catch-up time ≈ backlog / (apply rate − write rate)
- Drain from LB with hysteresis; never auto-promote a laggard
- Cap fleet-wide draining — don't evacuate the whole pool onto the primary
basics
~20 sDerive thresholds from what each replica is for: a staleness budget in seconds for read traffic, a recovery-point budget for failover candidates, and a hard byte threshold from log retention. Alert on sustained growth, not single spikes. Above the staleness budget, drain the replica from the read pool automatically, with hysteresis before re-adding.
solid answer
~1 minThresholds should be **derived, not guessed**. Each replica has a role, and each role implies a budget: - **Read-serving:** a staleness budget in seconds set by the product ("search results may be 5s old"). Alert at a fraction of it; act at it. - **Failover candidate:** a recovery-point budget — how much committed data may be lost on promotion. Measured on the flushed position, in bytes and seconds. - **Every replica:** a hard byte threshold derived from **log retention**. Past it, the replica cannot catch up at all and must be rebuilt from a fresh copy; if a replication slot is retaining log on its behalf, the same backlog threatens the primary's disk. This is the highest-severity alert in the set. Alert on **shape**, not instants: sustained growth over N minutes (lag rising and not draining) is a capacity deficit; a spike that drains is normal batch behaviour. Page on the former, ticket the latter. Automatic action: drain the replica from the load balancer above the staleness budget, re-add only after it has been under a lower threshold for a while (hysteresis), and never auto-promote a lagging replica. Guard against draining the whole pool at once — a global lag event must degrade reads to the primary or shed load, not evacuate every replica.
go deeper
Know that thresholds come from how stale the data may be for the feature, and that a lagging replica should stop serving reads.
Distinguish a spike that drains from lag that keeps rising, and connect log retention to the point where catch-up becomes a rebuild.
Give the severity ladder and the automated drain with hysteresis, and separate state alerts (replication stopped) from numeric thresholds.
Derive per-role budgets, reason about catch-up as backlog over apply-minus-write rate, and design the failure modes of the automation itself — fleet-wide drain caps, promotion safety, rebuild rate limits.
## Start from the purpose, not from a number "Alert at 10 seconds" is a number with no derivation, and it is wrong for most replicas in a fleet. A replica exists for a reason, and the reason dictates the budget. **Read-serving replicas** carry a *staleness budget*: the maximum age of data the feature can tolerate. A product-catalogue browse might tolerate 30 seconds; a post-checkout order page tolerates essentially none and should not be served from a replica at all. Different replica pools serving different features legitimately get different thresholds. **Failover candidates** carry a *recovery-point budget*: how much committed work may be lost when this node is promoted. This is measured against the **flushed** position (data durably on the replica's disk), because that — not the replayed position — bounds what a promotion can recover. Express it both in seconds (business language) and bytes (operational language). **Every replica**, whatever its role, is bounded by **log retention**. The primary keeps a finite amount of change log; a replica that falls further behind than that can no longer stream forward and must be rebuilt from a full copy — hours of work and, on a fleet, a capacity event. The mirror risk applies where a replication slot pins log on the primary for a replica's benefit: the backlog then consumes the *primary's* disk, and a full disk on the primary is an outage. Both directions justify a hard, high-severity threshold set well below the actual limit — with enough headroom to intervene. ## Alert on shape, not on instants Lag is spiky by nature: checkpoints, nightly batches, index builds, deploys. Threshold-only alerting on a spiky signal produces noise, and noise produces ignored pages. The informative distinction is between: - **A drainable spike** — lag rises, then falls back. Apply throughput exceeds the average write rate; the replica was momentarily overwhelmed. Only actionable if the peak breaches a budget. - **A sustained deficit** — lag rises monotonically over many minutes. Apply throughput is *below* the incoming write rate, and it will not recover on its own until the write rate falls or the replica's capacity changes. This is the condition that eventually hits log retention. So alert conditions should combine level and duration: "above X for N minutes", or "rising for N minutes with no drain". Recording the *derivative* alongside the level makes catch-up behaviour legible — during recovery lag should fall along a predictable slope, and a stalled recovery (flat, high lag) means the applier is blocked rather than working. A useful mental model: if the primary produces changes at rate *W* and the replica applies at rate *A*, then lag drains at *A − W*. Catch-up time from a backlog *B* is roughly *B / (A − W)*, and when *A ≤ W* it is infinite. That formula is what tells an operator whether waiting is a strategy. ## Severity tiers A workable ladder for a read-serving fleet: 1. **Info / ticket** — lag above a fraction of the staleness budget, sustained. Something changed; investigate during hours. 2. **Automated action** — lag above the staleness budget. The replica is drained from the read pool; no human needed yet. 3. **Page** — a large share of the pool is draining at once, or lag on a failover candidate breaches the recovery-point budget, or lag approaches the retention threshold. 4. **Page, highest severity** — replication stopped or errored, a replica has vanished from the primary's connection list, or a slot's retained log is nearing the primary's disk capacity. These are state alerts, independent of any lag number, and they must also fire when the metric stops arriving at all. ## What automation should do **Drain, don't panic.** Above the staleness budget, remove the replica from the read load balancer. This protects correctness (users stop seeing stale data) without touching replication, and it usually *helps* the replica catch up by freeing I/O and cache that were serving reads. **Use hysteresis.** Re-add only after lag has been below a lower threshold for a sustained period. Without it, a replica flaps in and out of the pool as lag oscillates around the line. **Protect against fleet-wide evacuation.** A cause common to all replicas — a giant batch job on the primary — will breach the threshold everywhere at once. Automation must cap how much of the pool can be drained (for example, never below a quorum), and the fallback must be a deliberate choice: shed load, serve degraded/cached results, or route to the primary if it can absorb it. Silently evacuating every replica onto the primary converts a staleness problem into an availability outage. **Never auto-promote a lagging replica.** Failover selection must consider the flushed position; promoting the most-behind node is data loss by automation. Where the engine supports it, prefer a synchronous or quorum-acknowledged candidate for promotion and let lag disqualify the rest. **Escalate rebuilds explicitly.** Once a replica passes the retention threshold, catch-up is off the table; the runbook is a re-clone. Automating that is reasonable only with rate limits, because re-cloning several replicas simultaneously loads the primary or the backup store precisely when things are already bad. ## Closing the loop Thresholds are only meaningful if the causes are addressed: chunk large write transactions to bound the spikes they create, keep index parity between primary and replicas so apply stays index-driven, and treat a sustained deficit as a capacity or sharding decision rather than an alert to be raised. An alert that fires nightly and is always resolved by waiting is a configuration bug, not monitoring.
- A replica has fallen far enough behind that the primary no longer retains the change log it needs. What are the options?Streaming catch-up is impossible because the required log records are gone, so the replica must be rebuilt from a fresh base copy — a restore or a new snapshot of the primary — and then resume streaming from that point. Prevention is the real answer: alert on a byte threshold well below retention, size retention for the worst realistic backlog, and rate-limit simultaneous rebuilds so recovery does not overload the primary or the backup store.
- Should replication lag alerts be identical across every replica in a fleet?No. Thresholds derive from role: a replica serving a latency-tolerant analytics dashboard can be minutes behind harmlessly, while a failover candidate is governed by the recovery-point objective and a replica behind a read-your-writes-sensitive feature has a budget near zero. One global number either pages constantly on tolerant replicas or fails to protect the strict ones. Only the retention-based hard threshold is genuinely universal.
saying these in an interview costs you the question
- Picking a round threshold with no derivation from a staleness or recovery-point budget
- Alerting on instantaneous lag only, so every nightly batch pages someone
- Auto-promoting whichever replica is available without checking its flushed position
- Automation that can drain every replica simultaneously, dumping all reads onto the primary
- Ignoring log retention and replication-slot disk growth as their own failure mode
- Assuming a lagging replica will always catch up given time