Fault Tolerance, High Availability, and Disaster Recovery Questions
Keeping a system serving despite failure, from code-level resilience to infrastructure-level recovery: circuit breakers, retries with backoff and jitter, timeouts, bulkheads, graceful degradation, and preventing cascading failures, alongside redundancy, failover (active-active versus active-passive), RPO and RTO objectives, backup and restore, and multi-region failover. Covers dependency-failure isolation, chaos engineering to validate resilience, failure-mode analysis, designing to nines of availability, cost-versus-availability tradeoffs, and recovery runbooks. Spans both the patterns that isolate partial failure and the disaster-recovery planning that restores a business-critical system after a major outage.
Design a chaos engineering program that progressively increases risk across service, database, and network layers for a critical system, starting with the safest experiments and working up. For each layer, what's your hypothesis, your blast-radius control, and your rollback criteria?
Sample Answer
Direct answer
Structure the program as a pyramid of increasing blast radius: start with single-instance, single-connection experiments in a canary slice of traffic, and only widen scope once the previous step stayed green. Every experiment, regardless of layer, needs the same three things stated up front: a falsifiable hypothesis (what should happen if the system is as resilient as you believe), a blast-radius control (the mechanism that caps how much traffic or infrastructure the experiment can touch), and a rollback criterion (the automated trigger that aborts the experiment before it becomes an incident).
Program structure by layer
| Layer | Hypothesis | Blast-radius control | Rollback criteria | Safest to riskiest experiments |
|---|---|---|---|---|
| Service | Losing a fraction of worker instances doesn't breach the SLO, because retries and circuit breakers absorb it | Canary AZ, capped at 1 to 5 percent of real traffic, feature-flag kill switch | Error rate more than 2x baseline AND p95 latency over SLO for 5 minutes, or success rate drops more than 1 percentage point absolute | (1) kill one non-primary worker process, (2) terminate 5 percent of workers in one AZ, (3) inject added latency into a canary slice, (4) disable retries on a canary path to check the fallback actually engages |
| Database | Read replicas and connection pooling keep reads available; write failures retry or queue without data loss | Target one replica or one connection pool at a time, throttle at the connection level, never touch the primary directly in early stages | Replication lag over a fixed threshold (for example 30 seconds), write failure rate spikes more than 1 percentage point absolute, or any detected data divergence | (1) throttle one read replica's I/O by 10 percent, (2) pause replication on one replica briefly, (3) close 5 percent of connections from a non-critical pool, (4) simulate primary failover, first in staging, then in a production canary |
| Network | Timeouts, retries, and the service mesh absorb transient network faults without payment (or equivalent critical-path) loss | Confine faults to one AZ and a capped traffic percentage using mesh-level fault injection, not a real router or switch | End-to-end success rate on the critical path drops more than 1 percentage point, or a circuit breaker stays open across more than 2 dependent services simultaneously | (1) add 50ms latency to one client-to-service hop, (2) inject 1 percent packet loss in one AZ for 5 minutes, (3) simulate a route flap between two internal services, (4) blackhole a non-critical downstream dependency and confirm graceful degradation, not failure |
The pyramid runs left to right within a layer, and layer to layer (service before database before network) because a service-level failure is the easiest to reason about and the cheapest to roll back; database and network faults touch more of the system at once and take longer to reverse cleanly.
Execution discipline
- Pre-flight: a written runbook, on-call and stakeholders notified, an automated abort mechanism wired to the rollback criteria (not a human watching a dashboard and deciding), and the experiment coded as a reproducible, version-controlled script rather than an ad hoc manual action.
- During: watch the rollback-criteria metrics in real time; the abort has to be automatic and fast, because by the time a human notices a metric crossing threshold and manually intervenes, the blast radius has often already grown past what the control was meant to cap.
- After: a lightweight postmortem regardless of outcome (a clean pass is still evidence worth recording), and only widen scope for the next run once the current one is unambiguously green, not "green with an asterisk."
Trade-offs & pitfalls
The single most common mistake is skipping straight to a production-wide experiment because a staging environment "doesn't reproduce the failure mode," which is often true but doesn't change the fact that the first production run of any new fault type belongs in the smallest blast radius you can construct, even if that means accepting a less realistic signal initially. A second pitfall is defining rollback criteria in terms of the fault itself (for example, "abort if packet loss exceeds 2 percent") instead of user-facing impact (error rate, latency, success rate); the fault is the input you're controlling, the rollback trigger has to watch the output, or you can hit exactly your intended fault level while still causing an unacceptable customer-facing outage. Not every fault type generalizes across domains the same way either: a GPU training job's most dangerous failure mode isn't a crashed worker (checkpointing handles that cheaply) but silent numerical divergence (the training job keeps running, but silently starts computing mathematically wrong updates to the model, with no crash or error to announce it), where the job keeps running and producing wrong gradients (a gradient is the per-step adjustment the training process makes to the model's internal numbers; a wrong one nudges the model in a bad direction instead of a good one) with no immediate error signal, so the safe blast-radius control there is different in kind from an HTTP service's traffic-percentage cap: it's about capping how long a divergence can run undetected before an automated metric check, watching a loss curve (a plot of the model's error over time, which should trend down) or a gradient norm (a single number summarizing how large the model's updates are; a sudden spike signals training has gone unstable), kills the job, not about capping how many requests are affected. A resilience program that only ever tests one fault type at a time also under-tests: real incidents are frequently two failures at once (a network blip during a deploy, a slow dependency during a traffic spike), so a mature program's later stages deliberately combine fault types once single-fault experiments across all three layers are consistently passing.
Walk through the common replication topologies, single-leader, multi-leader, and quorum-based, and how each affects consistency, latency, and availability.
Sample Answer
Direct answer: Single-leader replication routes all writes through one node and copies them out to followers, giving strong consistency on the leader but a failover gap if it dies. Multi-leader replication lets several nodes accept writes independently and merge them later, trading consistency for local write availability. Quorum-based replication has no fixed leader; reads and writes each require acknowledgment from a configurable subset of replicas, and the overlap between those subsets is what determines the consistency guarantee.
Structured elaboration
| Topology | Consistency | Write latency | Availability under partition |
|---|---|---|---|
| Single-leader | Strong on the leader; followers can lag (eventual, unless reads are forced to the leader) | Low (single write path, no coordination) | Writes unavailable if leader partitioned away until failover completes; reads can continue from followers |
| Multi-leader | Eventual; requires conflict resolution (last-write-wins, CRDTs, app-level merge) | Low locally at each leader | High: each site keeps accepting local writes during a partition, at the cost of divergence to reconcile later |
| Quorum-based | Tunable, from eventual to strong, depending on read/write quorum sizes | Higher (must wait for multiple acknowledgments, not just one) | Survives a minority of node failures without going unavailable; a true majority-losing partition halts progress |
The quorum math that determines consistency: for N replicas, a write quorum of W nodes and a read quorum of R nodes, the system guarantees a read overlaps with the most recent write whenever
W+R>NThis is a direct pigeonhole argument: if W and R are subsets of the same N-replica set and ∣W∣+∣R∣>N, they cannot be disjoint (two disjoint subsets can sum to at most N elements total), so they must share at least one replica, and that shared replica has both the latest write and is included in the read.
Worked example with N=3:
- W=2,R=2: W+R=4>3, so every read quorum is guaranteed to overlap every write quorum by at least one replica. This gives strong (read-your-writes) consistency, at the cost of needing acknowledgment from 2 of 3 replicas on both reads and writes.
- W=1,R=1: W+R=2≤3, no overlap is guaranteed. A write can land on replica A while a read is served entirely from replica B, missing it. This is fast (single-replica round trip) but only eventually consistent.
Trade-offs & pitfalls
- Single-leader is the simplest to reason about and the default choice unless you have a specific reason not to use it; its main weakness is the failover window (detecting the leader is gone and safely promoting a replacement), not steady-state operation.
- Multi-leader avoids that failover gap for writes but pushes complexity into conflict resolution; it's the right choice specifically when you need low-latency local writes at multiple sites and can tolerate (or algorithmically resolve) concurrent edits, not as a general-purpose upgrade over single-leader.
- Quorum systems let you dial the W/R trade-off per workload (e.g., W=1 for a write-heavy, tolerant-of-staleness workload; W=N for a read-heavy workload that wants every read to be a single, fast, guaranteed-fresh replica read), but that tunability is also a footgun: teams often ship with W+R≤N by default (e.g., both set to 1 for speed) without realizing they've silently given up the consistency guarantee they assumed they had.
- A common wrong turn: treating "quorum-based" as automatically stronger than single-leader. With W+R≤N it is weaker, not stronger, than a single-leader system with synchronous replication to at least one follower.
Design the failure detection that decides when to trigger an automated failover for a critical service. What health signals would you check, how would you set thresholds and windows to avoid mistaking a blip for a real failure, and when would you still want a human in the loop instead of a fully automatic failover?
Sample Answer
Direct answer: I'd layer health signals from cheap-and-fast (process liveness) to expensive-and-meaningful (dependency and business-metric checks), require multiple consecutive failures before declaring a node unhealthy to avoid reacting to a single blip, and keep a human in the loop specifically for the step that's hardest to reverse: promoting a new primary for stateful services, where an incorrect automated failover can cause data loss or a split-brain (two nodes each believing they're the one true primary and accepting conflicting writes at the same time), versus something like removing an unhealthy instance from a load-balancer pool, which is cheap to reverse and safe to fully automate.
Structured elaboration
| Signal | What it catches | Suggested check interval | Failures needed before acting |
|---|---|---|---|
| Process liveness | Process crashed or hung | Every 5s | 3 consecutive (15s) |
| Readiness / dependency connectivity | Process is up but can't reach its DB, cache, or queue | Every 10s | 2 consecutive (20s) |
| Latency / error-rate threshold | Process is up and connected, but degraded (slow, erroring) | Every 10-15s, rolling window | Sustained breach over a window (e.g., p95 > threshold for 3 samples), not a single sample |
| Business-metric sanity | Everything upstream looks healthy but the service is doing something wrong (e.g., checkout success rate collapsed) | Every 30-60s | Requires a real threshold breach, not a spike; slower-moving signal used as a final gate |
Why "N consecutive failures" instead of one: a single failed check can be a genuine blip (a GC pause, a brief network hiccup) rather than a real failure. If each check independently has some baseline flakiness probability p (a transient failure unrelated to a real outage), then requiring N consecutive failures before declaring unhealthy makes the false-trigger probability:
P(false trigger)=pNWith, say, p=0.05 (5% chance any single check fails transiently) and N=1 (react on the first failure), the false-trigger probability is just p=5%, meaning roughly 1 in 20 blips would incorrectly trigger action. Requiring N=3 consecutive failures drops that to:
P(false trigger)=0.053=0.000125=0.0125%a 400x reduction, at the cost of a real failure now taking 3 check intervals (here, up to 15 seconds at a 5s interval) longer to detect instead of one. That's the actual dial being turned: detection speed versus false-positive rate, and it should be set from an observed flakiness rate for your specific checks, not copied from another team's runbook.
Decision flow, including where automation stops and a human is required:
flowchart TD
A[Health check runs] --> B{N consecutive<br/>failures?}
B -->|No| A
B -->|Yes| C{What kind of<br/>action?}
C -->|Remove from LB pool| D[Fully automated:<br/>cheap, instantly reversible]
C -->|Open circuit breaker| D
C -->|Promote new primary<br/>for stateful service| E{Safety checks pass?<br/>replica caught up,<br/>quorum reachable}
E -->|No| F[Page human,<br/>do not auto-promote]
E -->|Yes, and blast radius<br/>is well-understood| G[Auto-promote,<br/>but page for review]
E -->|Yes, but ambiguous<br/>e.g. partition, not clear failure| F
"Quorum reachable" in that safety check means enough replicas are online and able to vote that a new primary can be safely elected, a majority agreeing on who's in charge, without risking two nodes each believing they're the primary at once.
Trade-offs & pitfalls
- Fully automated failover is right when the action is cheap and reversible (removing an unhealthy node from rotation); it's risky when the action is expensive or irreversible (promoting a database replica, since promoting the wrong one, or promoting during a network partition rather than an actual failure, can cause split-brain or data loss). The line isn't "how critical is the service," it's "how reversible is this specific action."
- Requiring consecutive failures trades detection speed for false-positive protection; too aggressive a requirement (e.g., N=10) means a genuine failure runs uncaught far longer than the blast radius justifies.
- Checks that themselves depend on a shared resource (e.g., every health check queries the same central database) can produce correlated, simultaneous "failures" across an entire fleet when that shared resource degrades, which looks like a mass outage but is really a single point of failure in the monitoring path itself.
- A common wrong turn: only checking process liveness and assuming that's sufficient. A process can be alive, passing liveness checks, and still be completely unable to serve real traffic because its only database connection pool is exhausted; readiness and dependency checks catch what liveness checks structurally cannot.
A multi-region failover led to a million records being double-processed, because the failover wasn't gated on the sinks actually being idempotent, and there was a race during leader election. Walk through the root cause and the concrete architecture changes that would prevent this from recurring.
Sample Answer
Direct answer
Two separate failures compounded here: the failover wasn't gated on the old leader actually being provably stopped (a real split-brain window during leader election, a period where two nodes can each believe they're the current leader and both may accept writes at the same time), and even a clean failover doesn't save you if the sinks aren't idempotent, because at-least-once delivery will still double-write on any redelivery. The durable fix addresses both: close the leader-election race with fencing, and make sinks idempotent as defense in depth, not either one alone.
Two-track architecture fix
flowchart LR
P[Producer] --> S[Partitioned stream]
S --> OL[Old leader epoch 5]
S --> NL[New leader epoch 6]
OL -->|late write epoch 5| GATE{Fencing check}
NL -->|write epoch 6| GATE
GATE -->|epoch >= committed| COMMIT[Upsert by event_id]
GATE -->|epoch stale| DROP[Discard write]
Fencing tokens close the leader-election race. Every leader term gets a monotonically increasing epoch number. The sink (or a shared coordination point in front of it) only accepts writes carrying an epoch at or above the last committed epoch:
epochold leader=5<epochnew leader=6⇒writes at epoch 5 arriving after commit are rejectedConcretely: the old leader holds epoch 5, becomes unreachable during a network partition, but is still alive and later resumes sending writes at epoch 5. A new leader was elected at epoch 6 after the lease expired and starts committing. Once epoch 6 has committed anything, the sink rejects any further epoch-5 writes outright, so the old leader's late writes are discarded rather than double-applied. This requires the storage layer to support an atomic conditional write (compare-and-swap on epoch, or a native fencing primitive); without that, the check can itself race.
Idempotent sinks are the second, independent layer. Every event carries a stable, globally unique event_id assigned at ingestion. Sinks write with an upsert keyed on event_id (INSERT ... ON CONFLICT (event_id) DO NOTHING, or a warehouse MERGE), so even a legitimate at-least-once redelivery, with no leader-election bug involved at all, lands exactly once in the sink's state.
Worked example
Both mechanisms are demonstrated in the fencing example above: the epoch comparison (5 versus 6) is the concrete, reproducible logic that decides accept versus reject, and it's independent of the idempotency layer. If fencing alone failed for some reason (a coordination bug elsewhere), the event_id upsert would still catch the duplicate, because the old leader would be resending events that already have IDs the sink has seen. Layering the two means either one failing doesn't reintroduce the bug, only both failing together would.
Trade-offs and pitfalls
The pitfall that caused this exact incident is fixing only the sink you noticed the duplicates in. Idempotency has to be audited across the full fan-out, a warehouse MERGE keyed on event_id doesn't help if a downstream webhook or notification consumer isn't also deduplicated, and duplicate user-visible side effects (a second email, a duplicate charge) will still occur even after the "main" sink is fixed. Fencing tokens add real requirements on the storage layer: many systems don't support atomic conditional writes natively, which pushes you toward optimistic-concurrency primitives (row versioning, conditional writes) as a design constraint on the whole pipeline, not just the failover logic. Idempotent upserts also add write overhead (a uniqueness check or MERGE is more expensive than a blind insert), which is a fair trade for correctness here but should be sized against actual throughput, not assumed free.
Your write region goes down for a couple of hours but your read regions are healthy. Design a graceful degradation plan: what stays available in read-only mode, what fails outright, and how do you communicate the degraded state to users?
Sample Answer
Direct answer
Keep serving reads from the healthy read regions immediately and unconditionally, since they don't depend on the write region at all; reject writes that require strong consistency or exactly-once guarantees outright with a clear error rather than pretending to accept them; and for writes that tolerate eventual consistency, queue them locally and durably so nothing is lost, then replay the queue once the write region recovers. The one rule that should never be violated is silently accepting a write and losing it: every write path is either "accepted and durably queued for later replay" or "rejected immediately," never "accepted and quietly dropped."
Architecture
Writes accepted by the anchor leader flow out to every read region as a CDC stream (change data capture: a continuous feed of every row-level change, read off the leader's own write log and replayed onto each replica), which is what keeps the read replicas caught up during normal operation, shown as the CDC node in the diagram below.
flowchart TD
CLIENT[Clients worldwide] --> RA[Region A read replica]
CLIENT --> RB[Region B read replica]
CLIENT --> RC[Region C read replica]
RA --> FWD[Write forwarder]
RB --> FWD
RC --> FWD
FWD --> LEADER[Anchor write leader, sharded by key range]
LEADER --> CDC[CDC stream]
CDC --> RA
CDC --> RB
CDC --> RC
When the write region (the anchor leader) goes down, the forwarder can no longer reach it. That's the trigger for the degradation workflow below; everything downstream of "detect the write region is unreachable" is about what each request type does next.
What stays available, what fails, what queues
| Request type | Behavior during the outage | Why |
|---|---|---|
| Reads (any kind) | Served normally from the nearest healthy read replica | Reads never depended on the write region; replication lag from before the outage started is the only staleness, and it stops growing once the write region is down (no new writes are landing to replicate) |
| Critical writes (billing, auth changes, anything needing exactly-once) | Rejected immediately with a 503 and a Retry-After header, structured as {mode: "read-only", eta: ..., writeAllowed: false} | These require strong consistency; queuing and replaying them risks double-charging or auth-state corruption, so refusing cleanly is safer than accepting and reconciling later |
| Non-critical writes (comments, drafts, "last viewed" timestamps) | Accepted at the edge, appended to a durable, region-local, idempotency-keyed queue, and the client gets a 202 Accepted with a queue ID | These tolerate eventual consistency, so queuing preserves the user's action without requiring the write region to be up right now |
| Complex cross-region features (distributed transactions, long multi-step jobs) | Disabled via feature flag for the duration of the outage | These need coordination that the write region's absence makes impossible to do correctly; a degraded system telling the user "this feature is temporarily unavailable" is better than a feature that silently does the wrong thing |
Communicating the degraded state
- API-level: every write endpoint that's disabled returns a structured error body, not a bare 503, so client applications can distinguish "the write region is down" from a generic server error and render an appropriate message instead of a raw failure.
- UX-level: a persistent banner ("Changes are temporarily disabled; you can still view everything. We'll apply queued updates once service resumes.") rather than silent failures scattered across individual actions, and per-action status for anything queued ("Saved locally, will publish when service resumes") so users aren't left wondering whether their action actually did anything.
- Ops-level: a status page update and, for major customers, proactive notification with an ETA, since 2 hours of read-only mode is the kind of thing enterprise customers expect to hear about before they notice it themselves.
Reconciliation when the writer returns
Queued writes are replayed in the order they were queued, using the idempotency key each write was tagged with to detect and skip duplicates (a client that retried a queued write multiple times during the outage shouldn't apply multiple times on replay). Conflict handling depends on the data type: commutative operations (counters, like counts) merge deterministically with no ambiguity; simple last-write-wins fields use the write's original logical timestamp, not the replay time, so a user's genuinely earlier edit doesn't overwrite a genuinely later one just because it replayed second; and content edits with a real risk of conflicting concurrent changes (two people editing the same document) get surfaced to the user as a merge conflict rather than silently resolved, because a silent wrong resolution is worse than asking.
Trade-offs & pitfalls
The scope of what gets classified "critical" versus "queueable" is the actual design decision here, and it's easy to get wrong in both directions: classifying too much as critical means the read-only window feels far more restrictive than it needs to be, while classifying too much as queueable risks a replay-time conflict resolution mess for data where "eventually consistent" was never actually an acceptable property. A related pitfall is bounding the queue: an unbounded local write queue during a multi-hour outage can grow large enough that the replay itself becomes a second incident (a burst of stale updates hitting the recovered write region all at once), so the queue needs both a size cap (reject new non-critical writes past a threshold, same as critical ones) and a TTL, with the two-hour outage in this scenario sized against realistic write volume before committing to "queue everything non-critical" as the policy. The same read-only degradation shape applies whether the thing behind the write region is a recommendation service, an image-preview pipeline, a personalization engine, or an analytics dashboard. In every case the same triage question decides the design: does this write need to be right immediately, or can it be right eventually, and only the answer to that question, not the specific feature, determines whether it queues or rejects.
Unlock Full Question Bank
Get access to all Fault Tolerance, High Availability, and Disaster Recovery interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.