Replication, Partitioning, and Sharding Questions
Scaling and distributing data across nodes: primary-replica and multi-primary replication, read-replica scaling, horizontal partitioning, and sharding strategies with their key-selection and rebalancing challenges. Covers replication lag, failover and split-brain handling, cross-shard operations such as joins, distributed transactions, and global secondary indexes, and the operational cost of a partitioned topology. Key to designing databases that scale horizontally.
Explain how database-level replication lag can affect read-after-write consistency in systems using read replicas. As a data engineer, describe patterns to ensure users see their own recent writes (for example, sticky sessions, session routing, read-your-writes guarantees) and trade-offs involved.
Sample Answer
Direct answer
Replication lag breaks read-after-write consistency (the guarantee that a user who just wrote something sees that write on their very next read) whenever the read happens to land on a follower that has not yet applied it. There are four standard patterns to guarantee a user sees their own recent writes despite this: sticky sessions (route a user's reads to the same node their write went to, for a window of time), session-scoped read-your-writes tracking (remember the position a user's last write reached, such as an LSN, log sequence number: a monotonically increasing marker of a database's write history, and only serve their reads from a replica that has caught up to at least that position), forcing reads-after-writes to the primary for a short window, and, less commonly, bounded-staleness reads that block briefly until a replica catches up rather than either redirecting or serving stale data outright.
The four patterns, and their trade-offs
| Pattern | How it works | Trade-off |
|---|---|---|
| Sticky sessions | Load balancer or client routes a given session/user consistently to the same node (often the primary, or a designated "home" replica) for some time window | Simple to implement, but concentrates that user's reads on one node (losing the read-scaling benefit of replicas for that user) and needs a fallback if the sticky node becomes unavailable |
| Read-your-writes token / session tracking | After a write, the client stores the LSN or a similar log position it reached; on the next read, only route to a replica whose applied position is at or past that token | Precise (only redirects exactly the reads that need it) but requires plumbing the token through the client or session layer and requires each replica to expose its current apply position for comparison |
| Force-read-from-primary flag | Explicitly mark reads that must be fresh (e.g., "show me my own just-submitted order") to bypass replicas and hit the primary directly | Zero staleness for the marked reads, but adds primary load proportional to how many reads get flagged, so it is only safe when a small, well-identified subset of reads actually needs it |
| Bounded-staleness / wait-for-replica reads | The replica delays answering a read until its applied position reaches the requested minimum, instead of redirecting elsewhere | Avoids extra primary load and avoids sticky routing, but turns a routing problem into a latency problem: the read blocks for however long the replica is behind, which is unbounded unless capped with a timeout-and-fallback |
Worked example: a social-network timeline
A user posts a comment, then immediately reloads their own timeline. If the reload hits a replica still 800 milliseconds behind (a fully normal amount of replication lag under moderate load), the comment appears to have vanished, which reads as a bug even though nothing is actually wrong. A concrete implementation: on write, the API server captures the WAL position (or an equivalent monotonic write marker) the write reached and returns it to the client as an opaque read-your-writes token embedded in the response (client-side technique); the client attaches that token to its next read request; the read-routing layer (server-side technique) checks each candidate replica's current applied position and only sends the request to one that has caught up, falling back to the primary if no replica has caught up within a short timeout (for example 50 milliseconds), rather than blocking indefinitely. This combination avoids sticky-session's node-affinity cost and avoids sending every read to the primary, while still guaranteeing the user sees their own comment.
Implementation details worth naming explicitly
A connection pooler alone does not solve this: PgBouncer, the common PostgreSQL pooler, is a pure TCP-level connection pooler with no query inspection and no awareness of a replica's replication position, so it cannot itself implement token-aware read/write splitting, and pointing it at a random or round-robin pool of replicas will pick one that has not caught up. The token-aware routing logic has to live in a layer that actually understands queries and replica state: application-level routing that picks a connection target per request based on the token, or a proxy built for query routing rather than plain pooling (for example Pgpool-II's load-balancing mode, which does inspect queries to split reads from writes) sitting in front of, or instead of, a plain pooler like PgBouncer, which continues to do only what it is built for: pooling connections to whichever destination it is told to use. A retry-on-lag pattern (catch the case where no replica is caught up, retry against the primary once, rather than retrying the same lagging replica) keeps tail latency bounded. Session-affinity state (which token belongs to which in-flight session) needs to live somewhere the load balancer or API layer can cheaply access it on every read, typically in the request context itself (the token travels with the request) rather than in a separate lookup service, to avoid adding a new dependency to every read's critical path.
Trade-offs and pitfalls
Sticky sessions are the easiest pattern to ship first, but teams often leave them in place permanently rather than migrating to a token-based approach, which silently caps how much read-scaling benefit replicas actually provide as traffic grows, since sticky users never spread across the replica fleet. Forcing all writes' subsequent reads to the primary "to be safe" is the opposite mistake: it looks correct but reintroduces exactly the primary-load problem replicas exist to solve, for reads that, in most product flows, do not actually need same-millisecond freshness. The right default is to reserve the primary-forcing pattern for the specific, small set of user actions where staleness is visibly wrong (their own just-written content), and use the token-based approach for everything else.
How would you detect and mitigate silent data corruption or a split-brain scenario in a replicated database? Propose detection mechanisms, automated mitigation steps, and offline repair procedures that preserve data correctness.
Sample Answer
Direct answer
Split-brain (two nodes both believing they are the authoritative primary at the same time, typically after a network partition) and silent data corruption both stem from the same root cause: something wrote data that later turns out to conflict with, or be inconsistent with, the rest of the system, and nobody noticed at write time. Detection has to be active (continuously comparing state) rather than passive (waiting for a user to report wrong data), automated mitigation has to stop the bleeding fast (fence the wrong writer, redirect traffic), and offline repair has to reconcile divergent data without guessing, using whichever side has a verifiable, complete record of what actually happened.
Detection mechanisms
sequenceDiagram
participant N1 as Node A (was primary)
participant N2 as Node B (promoted primary)
participant Detector as Reconciliation job
Note over N1,N2: Network partition isolates Node A
N2->>N2: Promoted after fencing timeout, begins accepting writes
Note over N1: Node A incorrectly believes it is still primary (fencing failed to reach it in time)
N1->>N1: Also accepts writes during the partition window
Note over N1,N2: Partition heals
Detector->>N1: Compare checksums / row versions
Detector->>N2: Compare checksums / row versions
Detector->>Detector: Detect divergent rows written on both sides
Detector->>Detector: Quarantine divergent rows for offline reconciliation
- Consensus-layer detection (best, catches it before it happens). A properly implemented leader-election system (using a distributed configuration store like etcd or ZooKeeper, requiring a majority quorum to grant and renew a leader lease) makes split-brain structurally hard, not just detected after the fact: a node cannot legitimately believe it holds the leader lease unless it can prove it against a majority. This is a design choice made before the incident, not a detection technique applied during one, but it is the most effective mitigation available.
- Fencing-token/generation-number checks. Every write carries a monotonically increasing "generation" or "epoch" number tied to the current leadership term. Any downstream consumer (storage layer, replicas, or an external system) rejects a write carrying an older generation number than one it has already seen, which catches a zombie old-primary's writes even if the primary itself has not yet realized it lost leadership.
- Continuous checksum/version reconciliation. A background job periodically compares row-level checksums or version numbers between what should be identical copies (a shard and its replica, or two nodes that briefly diverged), flagging mismatches for investigation rather than waiting for a user-visible symptom.
- Application-level invariant checks. Domain-specific sanity checks (an account balance that went negative when the business rule says it never should, a row updated with a timestamp older than its own last-modified timestamp) catch corruption that passes structural checks (the row is syntactically valid) but violates business meaning.
Automated mitigation steps
The moment split-brain is suspected (two nodes both claiming leadership, or a generation-number mismatch detected), the automated response should: immediately fence the node with the lower/older generation number (revoke its ability to accept writes, via a network-level block, a credential revocation, or a hardware watchdog reset), stop routing any traffic to it, and freeze (do not auto-merge) any data written during the suspected divergence window rather than guessing which side is "right." Auto-merging without a clear, safe resolution rule is how a detection system turns one incident into a second, worse one.
Offline repair, preserving correctness
Once the systems are stable, reconciliation is a deliberate, evidence-based process, not an automated guess: identify the exact window of potential divergence using the generation-number or timestamp evidence, extract the rows written on each side during that window, and resolve each conflict using domain rules the team defines in advance (for genuinely append-only or idempotent operations, keep both if they do not actually conflict; for a true same-key conflicting write, prefer whichever side has independently verifiable evidence of being the legitimate majority-quorum leader during that window, and treat the other side's writes as candidates for compensating transactions rather than silent overwrites). Anything that cannot be confidently reconciled gets flagged for manual review rather than resolved by a rule that might be wrong, and the incident is not closed until every flagged row has an explicit resolution, not just "most of them look fine."
Trade-offs and pitfalls
The most dangerous mistake is auto-resolving conflicts with a blanket rule like last-writer-wins by wall-clock timestamp, without checking whether clocks were actually synchronized across the two nodes during the incident; clock skew during exactly the kind of network trouble that causes split-brain is common, which means the "last" writer by timestamp may not be the actually-later write. The second mistake is treating detection as a one-time incident-response activity rather than continuous background reconciliation; silent corruption that never triggers an obvious symptom (a slowly-diverging count, not a crash) can persist for a long time before anyone notices without an active, scheduled check.
Design alerting thresholds and anomaly detection for data divergence signals to minimize false positives. Describe statistical baselining, adaptive thresholds, scoring rules for multi-metric signals (e.g., checksum mismatch rate, replication lag, stale-read rate), and escalation policies for SRE teams.
Sample Answer
Direct answer
A single fixed threshold cannot serve a metric with a legitimate daily or weekly cycle: set it loose enough to survive normal peak-hour variation and it misses a real problem that happens during a normally-quiet period, set it tight enough to catch a quiet-period problem and it false-pages on every normal peak. The fix is a baseline that adapts to time-of-day (an exponentially weighted moving average, EWMA, of recent values, which tracks a slowly-changing expected level, plus an adaptive band width derived the same way from recent deviation), scored jointly across the several signals the question names (replication-lag, checksum-mismatch-rate [how often a background comparison of the primary's and a replica's data finds rows whose contents no longer match], stale-read-rate) rather than paging on any single metric alone, because a single adaptive metric is itself noisier than it looks, as the worked example below shows directly.
Statistical baselining and adaptive thresholds
An EWMA baseline weights recent observations more heavily than old ones, with a decay controlled by a span parameter (roughly, "how many recent minutes matter"), so it re-centers on the current time-of-day's normal level rather than a single number for the whole week. The threshold band should be adaptive the same way: track an EWMA of the recent absolute deviation from the baseline (a rolling, exponentially-weighted analog of median absolute deviation, MAD) and set the alert band as a multiple of that, with a small floor so a dead-quiet period does not collapse the band to near zero and start flagging routine noise.
Scoring rules for multi-metric signals
Require correlated evidence across independent metrics before paging, rather than trusting any single adaptive metric on its own. The worked example below shows why this matters concretely: a naive single-metric adaptive detector, tuned sensitively enough to catch a real but moderate incident, also generates a real amount of noise around the legitimate morning ramp-up (where the true baseline is changing fastest and the deviation-tracking EWMA temporarily lags). Scoring replication lag, checksum-mismatch-rate, and stale-read-rate together, and requiring at least two of the three to agree before paging, turns that noisy single-metric detector into a clean one, because unrelated noise on one metric rarely coincides with unrelated noise on a different, independently-measured metric at the same minute, while a real incident moves all of them together.
Escalation policy
- One metric alone flagged: log it, do not page. This is the routine "adaptive threshold caught some noise" case the worked example demonstrates directly.
- Two or more of the three signals agree, sustained past a short confirmation window (a couple of minutes, to avoid paging on a single noisy sample): page on-call with the specific combination of signals that fired, since that combination itself is diagnostic information (lag plus stale-reads without checksum mismatches points toward a capacity problem; checksum mismatches without lag points toward a correctness bug in replication itself, not a speed problem).
- All three sustained and worsening: escalate severity, since this pattern is consistent with a systemic issue rather than an isolated blip.
Worked example
A 7-day, minute-resolution synthetic replication-lag series with a real diurnal cycle (roughly 150ms overnight, up to roughly 900ms at the daytime peak, both legitimate) and one injected, moderate incident: a +500ms lag increase during a normally-quiet overnight window, roughly a 4x increase over that window's own baseline but still well below the LEGITIMATE daytime peak. Seed is pinned for reproducibility.
import numpy as np
rng = np.random.default_rng(seed=42)
MINUTES = 7 * 24 * 60 # one week at 1-minute resolution
t = np.arange(MINUTES)
# diurnal pattern: replication lag naturally rises during business-hours write peaks
# (up to ~900ms) and falls overnight (~150ms baseline). This swing is LEGITIMATE, not
# an incident -- and it is large enough to dominate a whole-week static threshold.
hour_of_day = (t / 60) % 24
diurnal = 150 + 750 * np.clip(np.sin((hour_of_day - 6) / 24 * 2 * np.pi), 0, None) ** 1.5
noise = rng.normal(0, 40, size=MINUTES)
lag_ms = np.clip(diurnal + noise, 5, None)
# inject ONE real but MODERATE incident (e.g. a background compaction stealing IO)
# during a normally-quiet NIGHT window: +500ms on top of a ~150ms night baseline is a
# genuine ~4x lag increase, but it never approaches the big legitimate daytime swings.
incident_start = 3 * 24 * 60 + 2 * 60 # day 3, 02:00 (quiet hours)
incident_len = 30
lag_ms[incident_start:incident_start + incident_len] += np.linspace(80, 500, incident_len)
print(f"series: {MINUTES:,} minutes (7 days), incident at minute {incident_start} "
f"(day {incident_start // 1440}, {(incident_start % 1440) // 60:02d}:00), "
f"duration {incident_len} min, peak injected lag "
f"{lag_ms[incident_start:incident_start + incident_len].max():.0f}ms")
print(f"normal daytime peak lag (no incident) reaches ~{diurnal.max():.0f}ms baseline + noise")
print()
# -- Option A: static threshold. Pick the threshold as 3 standard deviations above the
# GLOBAL mean lag (a common naive baselining choice). --
static_thresh = lag_ms.mean() + 3 * lag_ms.std()
static_flags = lag_ms > static_thresh
print(f"-- static threshold: global mean + 3*std = {static_thresh:.0f}ms --")
print(f" total minutes flagged: {static_flags.sum()}")
false_positive_minutes = static_flags.sum() - min(static_flags[incident_start:incident_start + incident_len].sum(), incident_len)
print(f" of which false positives (outside the real incident window): {false_positive_minutes}")
detected_in_incident = static_flags[incident_start:incident_start + incident_len]
first_detect = np.argmax(detected_in_incident) if detected_in_incident.any() else None
print(f" incident minutes correctly flagged: {detected_in_incident.sum()}/{incident_len}, "
f"first flagged at minute +{first_detect if first_detect is not None else 'NEVER'} of the incident")
print()
# -- Option B: adaptive baseline. EWMA tracks the expected lag (adapts to the diurnal
# cycle), and a robust EWMA of the absolute deviation (like a rolling MAD) sets the
# band width, so the threshold is TIGHT at night and WIDE at a normal daytime peak --
# exactly the multi-metric-safe pattern the question asks for, generalized here to lag
# alone for clarity; checksum-mismatch-rate and stale-read-rate get their own band and
# an incident needs a WEIGHTED SCORE across bands to fire (shown after this block).
def adaptive_flags(series, span=180, k=3.5):
alpha = 2 / (span + 1)
ewma = np.empty_like(series)
mad = np.empty_like(series)
ewma[0] = series[0]
mad[0] = 0.0
for i in range(1, len(series)):
ewma[i] = alpha * series[i - 1] + (1 - alpha) * ewma[i - 1]
dev = abs(series[i - 1] - ewma[i - 1])
mad[i] = alpha * dev + (1 - alpha) * mad[i - 1]
band = np.maximum(k * mad, 25) # floor so a dead-quiet night doesn't get a 0-width band
return series > (ewma + band), ewma, band
adaptive_flagged, ewma, band = adaptive_flags(lag_ms)
print(f"-- adaptive baseline: EWMA(span=180min) + 3.5*EWMA(|deviation|), floor 25ms --")
print(f" total minutes flagged: {adaptive_flagged.sum()}")
fp_adaptive = adaptive_flagged.sum() - min(adaptive_flagged[incident_start:incident_start + incident_len].sum(), incident_len)
print(f" of which false positives (outside the real incident window): {fp_adaptive}")
detected_in_incident_a = adaptive_flagged[incident_start:incident_start + incident_len]
first_detect_a = np.argmax(detected_in_incident_a) if detected_in_incident_a.any() else None
print(f" incident minutes correctly flagged: {detected_in_incident_a.sum()}/{incident_len}, "
f"first flagged at minute +{first_detect_a if first_detect_a is not None else 'NEVER'} of the incident")
print()
print(f"threshold at the moment the incident starts: static={static_thresh:.0f}ms (fixed all week) "
f"vs adaptive={ewma[incident_start] + band[incident_start]:.0f}ms (tracks the time-of-day baseline)")
# -- multi-metric scoring sketch: a real alert should require correlated evidence, not
# one metric alone, to avoid paging on lag noise from an unrelated blip --
checksum_mismatch_flag = np.zeros(MINUTES, dtype=bool)
checksum_mismatch_flag[incident_start + 5:incident_start + incident_len] = True # detector lags lag by 5min
stale_read_flag = np.zeros(MINUTES, dtype=bool)
stale_read_flag[incident_start + 2:incident_start + incident_len] = True
score = adaptive_flagged.astype(int) + checksum_mismatch_flag.astype(int) + stale_read_flag.astype(int)
page_worthy = score >= 2 # require at least 2 of 3 independent signals to agree
print()
print(f"-- combined score (lag + checksum-mismatch-rate + stale-read-rate), page if >=2/3 agree --")
print(f" minutes that would page: {page_worthy.sum()}, all inside the real incident: "
f"{bool(page_worthy[incident_start:incident_start + incident_len].sum() == page_worthy.sum())} "
f"(flagged range: minute {np.argmax(page_worthy)} to "
f"{MINUTES - 1 - np.argmax(page_worthy[::-1])})")
Output:
series: 10,080 minutes (7 days), incident at minute 4440 (day 3, 02:00), duration 30 min, peak injected lag 733ms
normal daytime peak lag (no incident) reaches ~900ms baseline + noise
-- static threshold: global mean + 3*std = 1193ms --
total minutes flagged: 0
of which false positives (outside the real incident window): 0
incident minutes correctly flagged: 0/30, first flagged at minute +NEVER of the incident
-- adaptive baseline: EWMA(span=180min) + 3.5*EWMA(|deviation|), floor 25ms --
total minutes flagged: 57
of which false positives (outside the real incident window): 32
incident minutes correctly flagged: 25/30, first flagged at minute +4 of the incident
threshold at the moment the incident starts: static=1193ms (fixed all week) vs adaptive=270ms (tracks the time-of-day baseline)
-- combined score (lag + checksum-mismatch-rate + stale-read-rate), page if >=2/3 agree --
minutes that would page: 26, all inside the real incident: True (flagged range: minute 4444 to 4469)
The static global threshold (mean plus 3 standard deviations over the whole week) never fires at all: the legitimate daytime swings inflate the whole-week variance enough that a real, moderate overnight incident never crosses it. A single adaptive EWMA-plus-band detector, tuned sensitively enough to actually catch that incident, also produces real false positives, concentrated around the fastest-changing part of the legitimate daily cycle rather than being random noise. Combining it with two independently-injected companion signals (checksum-mismatch-rate and stale-read-rate) and requiring two of three to agree collapses the false positives down to a clean window that lands entirely inside the real incident.
Trade-offs and pitfalls
- A static threshold calibrated on a whole week of data is not conservative, it is blind to any incident smaller than the legitimate daily swing. This is the opposite of what most people assume a "safe, simple" threshold does.
- A single adaptive metric is not automatically a strict improvement over a static one; it trades one failure mode (missing a subtle incident) for a different one (false-positiving on legitimate regime changes) unless it is paired with either seasonal decomposition (modeling and removing the known daily shape before testing the residual) or multi-metric confirmation.
- The confirmation-window length is a real trade-off, not a formality: too short and you page on transient noise that would have self-resolved; too long and you delay detection of a genuine fast-moving incident.
- This entire approach assumes the baseline itself is healthy to begin with. If the system has been quietly degraded for weeks, the EWMA baseline adapts to the degraded state as "normal" and stops flagging it; periodically compare the adaptive baseline against a longer-window reference to catch slow drift that no minute-to-minute detector will ever surface.
Design an automated failover system for primary-replica database pairs that minimizes split-brain risk. Cover health checks, fencing mechanisms, promotion safety checks, and how you would validate and audit automatic promotions. Also explain what data can be lost during the failover and how RPO and RTO differ between manual and automatic failover.
Sample Answer
Direct answer
An automated failover system for primary-replica pairs needs four pieces working together to minimize split-brain risk: a consensus-based health check (a majority, not a single observer, must agree the primary is actually down), fencing that guarantees the old primary cannot keep accepting writes once it loses leadership, a promotion-safety check that refuses to promote a replica too far behind, and an audit trail proving what happened and when. This is the same design pattern Patroni implements for PostgreSQL clusters, and it generalizes to any primary-replica system.
Architecture
flowchart TD
A[Health check: primary missed N heartbeats] --> B{DCS quorum reachable?}
B -->|no quorum| C[Do not promote: minority side, refuse writes]
B -->|quorum reached| D[Elect candidate: most caught-up in-sync replica]
D --> E[Fence old primary via watchdog/STONITH]
E --> F{Fencing confirmed?}
F -->|no| G[Abort promotion, alert on-call: cannot guarantee safety]
F -->|yes| H[Promote candidate to primary]
H --> I[Update DNS/service discovery to new primary]
I --> J[Audit log: who/what promoted, at what LSN, data-loss window]
J --> K[Old primary rejoins later as a follower once fencing is lifted]
Health checks. A single node observing "the primary didn't respond" is not enough evidence, since that node itself could be the one that is network-partitioned. Route the decision through a distributed configuration store (DCS: a small, separately-replicated coordination service such as etcd, Consul, or ZooKeeper) that requires a majority, or quorum, of its own members to agree before any leadership change proceeds; this is the same mechanism Patroni uses for PostgreSQL. If the surviving side of a partition cannot reach a majority of the DCS, it must not promote, even if it genuinely believes the primary is down, because it cannot distinguish "the primary is dead" from "I am the one who is isolated."
Fencing mechanisms. Before any promotion, the old primary must be guaranteed unable to keep accepting writes. Two layers, used together: revoke its ability to renew its leadership lease in the DCS (a soft fence: it should demote itself once it notices the lease expired), and a hardware or OS-level watchdog that forcibly resets the host if the soft fence has not been confirmed within a bounded time (a hard fence, sometimes called STONITH, "shoot the other node in the head"). The hard fence exists specifically for the case where the old primary is unresponsive rather than well-behaved, since a well-behaved node would have already demoted itself.
Promotion safety checks. Never promote a replica that is more than a defined, bounded amount behind the primary's last known position (measured by comparing log sequence numbers, LSNs). If the most caught-up available replica is still too far behind, the safer action is to alert and wait (accepting downtime) rather than promote and silently accept a larger-than-intended data loss; that threshold (how far behind is "too far") is a deliberate, pre-agreed trade-off between availability and data loss, not a default left unexamined.
Validating and auditing automatic promotions. Every automatic promotion writes an audit record: which node was promoted, what LSN it was at, what the old primary's last known LSN was (which bounds the data-loss window), and what triggered the decision. This is what turns "the system did something" into something a human can verify after the fact, and it is what a chaos-testing exercise (deliberately triggering failures in a controlled way, including in production during low-traffic windows, specifically to verify this pipeline works under realistic conditions rather than only in a lab) should be validating end-to-end, including checking that dependent services handling the failover (connection pools reconnecting, in-flight requests retried safely) do not cascade the failure outward.
Manual versus automatic failover: what data can be lost, and the RPO/RTO difference
| Automatic failover | Manual failover | |
|---|---|---|
| Recovery time objective (RTO) | Seconds to low tens of seconds: bounded by health-check interval, DCS quorum round trip, and fencing confirmation time | Minutes: bounded by human detection, human judgment time, and manual execution of the same steps |
| Recovery point objective (RPO) | Determined mechanically by whichever replica was most caught-up at decision time, which may not be the theoretical best choice if the promotion-safety threshold accepted a replica with some lag | Potentially better: a human can inspect multiple replicas' exact positions, attempt to recover unshipped writes from the old primary's log if it is still reachable in a degraded state, and choose the option with the least data loss, at the cost of the extra time that inspection takes |
| Failure mode if wrong | Fast, but a bug in the automation (an incorrect promotion-safety check, a fencing race condition) can promote unsafely at machine speed, before a human notices | Slow, but a human in the loop can catch an obviously wrong situation (e.g., "the primary is actually fine, this is a monitoring false positive") before acting |
| What can be lost | Any write acknowledged by the old primary but not yet replicated to the promoted replica at the moment of the fence; bounded by the promotion-safety threshold | The same category of loss, but potentially smaller, since a human can wait slightly longer to let a lagging replica catch up, or attempt log recovery from the old primary, trading RTO for a better RPO |
Automatic failover is the right default for a well-tested, generic primary-down scenario, because RTO matters more than shaving the last few seconds of RPO for most services. Manual (or manual-confirmed) failover is worth keeping as an option, or even the default, for a scenario where the failure mode is ambiguous (a regional network partition isolating a meaningful fraction of followers, rather than a clean single-node crash), where an automated system cannot reliably distinguish "promote now" from "wait, this might resolve itself" and a wrong automatic decision (over-eager promotion during a transient blip) is worse than the extra minute a human takes to confirm.
Trade-offs and pitfalls
The most common design flaw is building the health check and the fencing mechanism as two independently-tuned systems that can disagree: a health check that decides to promote faster than the fencing mechanism can guarantee the old primary is actually stopped creates exactly the split-brain window the whole design exists to prevent. The second common flaw is never testing the failover path against a real, partial-partition scenario (only testing against a clean node crash), since a partition that isolates a meaningful fraction, not all, of the followers behaves very differently (some followers see a different quorum than others) than a simple full-node failure, and is the scenario most likely to expose a subtle bug in the promotion-safety logic.
Design a monitoring and alerting strategy for a sharded database cluster. List per-shard and cluster-wide metrics to collect (latency percentiles, QPS, replication lag, disk utilization), alert thresholds and severity, and what actions/automation or runbooks should be triggered for common alerts.
Sample Answer
Direct answer
Monitor at two levels that answer two different questions: per-shard metrics tell you which specific shard is unhealthy, and cluster-wide aggregates tell you whether the cluster as a whole is meeting its service-level objective (SLO, the target level of service the team has committed to). Alert on the aggregate for paging (nobody should be paged because one shard out of two hundred has a blip), and keep the per-shard view for diagnosis once a page fires. Thresholds should be set relative to each shard's own healthy baseline, not one global number, since shards can have legitimately different steady-state load.
Metrics to collect
Per-shard:
- Query latency, p50/p95/p99 (95th and 99th percentile response time, the tail behavior that matters more than the average for user-facing impact)
- Queries per second (QPS) and the read/write split
- Replication lag to that shard's replicas (for example Postgres's
pg_stat_replication.replay_lag, the delay between a write committing on the primary and becoming visible on a given replica) - Disk utilization (both space used and IO throughput/IOPS, input/output operations per second, since a shard can be IO-saturated well before it runs out of space)
- Connection count against the configured limit
- Lock wait time and deadlock rate (how often two transactions end up blocking each other in a cycle, forcing the database to kill one so the other can proceed)
Cluster-wide:
- Aggregate QPS and error rate across all shards
- Number of shards currently unhealthy (below their own baseline) as a single number, so a dashboard shows "3 of 200 shards degraded" instead of 200 individual panels
- Rebalance/resharding operations in flight and their progress
- Cross-shard query fan-out latency (the scatter-gather pattern from a query that has to touch multiple shards is only as fast as its slowest shard, so this is a distinct signal from any single shard's own latency)
Thresholds, severity, and what should page
| Signal | Warning | Critical (pages on-call) | Rationale |
|---|---|---|---|
| Per-shard p99 latency | 2x that shard's 7-day rolling baseline | 5x baseline, sustained 5 minutes | A short spike is normal; a sustained multiple of baseline is not, and using the shard's OWN baseline avoids a single global number being wrong for every shard with different traffic |
| Replication lag | 30 seconds | 5 minutes, or growing rather than stable | Below 30s is invisible to most consumers; past a few minutes, stale-read risk and failover data-loss risk (if the primary fails before a lagging replica has caught up, whatever had not yet replicated is gone once that replica gets promoted) both become real |
| Disk utilization | 75% | 90%, or projected to fill within 24 hours at current growth rate | Running out of disk on a primary is a hard outage, not a degradation, so this needs lead time to act, not just a threshold at the edge |
| Connection count | 70% of configured limit | 90% | Hitting the connection limit fails new requests outright; this should never be a surprise |
| Unhealthy shard count (cluster-wide) | 1 percent of shards | 5 percent of shards, or any single shard fully unreachable | A rising unhealthy-shard count often indicates a systemic cause (a bad deploy, a shared dependency failing) rather than 200 independent unlucky shards |
What should happen automatically versus what pages a human
- Automatic, no page: a single shard's connection pool (its fixed set of reusable, already-open database connections) exhausting temporarily under a burst (if the application has retry-with-backoff, this self-heals); routine rebalance operations proceeding on schedule.
- Page, with a runbook attached, not just an alert: replication lag critical (runbook: check for a long-running transaction blocking apply, check network throughput between primary and replica, decide whether to fail over or wait); disk critical (runbook: confirm whether growth is organic or a leak/bug, trigger an emergency resize or archive-and-purge); unhealthy-shard-count critical (runbook: check for a common cause across the affected shards before treating them as independent incidents).
- A runbook is not optional documentation, it is the alert's payoff. An alert that pages someone with no linked, concrete next action just moves the "what do I do now" research into the middle of the incident instead of before it.
Trade-offs and pitfalls
- Alerting on raw per-shard thresholds instead of each shard's own baseline produces both false pages (a shard that's always a bit hotter than average) and missed pages (a shard whose baseline is unusually low, where a real problem never crosses a one-size-fits-all number).
- Too many per-shard alerts wired directly to paging is how on-call burns out. Fan noisy per-shard signals into the aggregate unhealthy-shard-count metric for paging, and keep the per-shard dashboards for triage after the page, not as 200 separate alert rules.
- A metric with no owner decays. Review thresholds periodically against actual incident history: a threshold that never fires might be too loose to catch real problems, and one that fires constantly and gets silenced is worse than no alert at all, since it trains people to ignore it.
Unlock Full Question Bank
Get access to all 14 Replication, Partitioning, and Sharding interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.