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.
Using the CAP theorem, walk through the trade-off you'd make for a financial ledger service versus an analytics event aggregator. Which guarantee does each give up during a network partition, and why?
Sample Answer
CAP theorem says that during a network partition, a distributed system must choose between consistency (every read sees the latest acknowledged write) and availability (every request gets a response), because it can't guarantee both while the partitioned halves can't talk to each other. Partition tolerance itself isn't optional for any system that spans more than one node, so the real choice interview questions are testing is CP versus AP, and the right answer depends entirely on what happens if you get it wrong.
Financial ledger: choose CP
A ledger getting a stale or divergent balance is a correctness bug with real financial consequences (double-spend, incorrect balance shown, a transaction accepted twice), so consistency has to win. The concrete mechanism is a consensus protocol requiring a write quorum.
Deriving what "sacrifice availability" actually means, with a pinned example: take a 5-node replica set (a typical Raft cluster size) where a majority quorum is required to commit a write:
quorum=⌊2N⌋+1=⌊25⌋+1=3That cluster tolerates up to N−quorum=2 node failures while still committing writes. Now suppose a network partition splits the 5 nodes into a 3-node side and a 2-node side. The 3-node side still has a majority (3 ≥ quorum of 3), so it keeps accepting writes and stays both consistent and available. The 2-node side does not have a majority (2 < 3), so by design it must refuse writes, becoming unavailable, specifically to prevent both sides from independently committing conflicting transactions. That refusal on the minority side, not a global shutdown, is the literal, computable meaning of "CP sacrifices availability during a partition": only the minority partition goes unavailable, and only for writes.
Analytics event aggregator: choose AP
An analytics pipeline getting an event a few seconds late, or briefly double-counted before deduplication catches up, is a rounding error on a dashboard, not a financial loss, so availability wins: every node keeps ingesting even when it can't see the others.
The mechanism looks different from the ledger: instead of a write quorum gatekeeping every write, use leaderless or partitioned ingestion (each partition or node accepts writes for its own shard independently, as in Kafka-style partitioned logs or a Dynamo-style leaderless store), so there's no majority to lose and no write path that can be blocked by a partition. The cost shifts from "unavailable during a partition" to "eventually reconciled after one": deduplication on ingest, idempotent consumers, and background reconciliation jobs to merge whatever diverged while the partition was open.
Comparing the two designs
| Dimension | Ledger (CP) | Analytics aggregator (AP) |
|---|---|---|
| Write path | Quorum-gated (majority must ack) | Leaderless / partitioned, no quorum gate |
| During a partition | Minority side refuses writes | Both sides keep accepting writes independently |
| Consistency mechanism | Consensus (Raft/Paxos-style) | Eventual consistency + reconciliation |
| Failure mode if you pick the wrong side | Double-spend, incorrect balances | Stale dashboard, temporary undercounts |
| Client-side pattern needed | Idempotency keys so a client retry after a rejected write is safe | Deduplication keys so late/duplicate events don't double-count on reconciliation |
Trade-offs and pitfalls
The common misreading of CAP is treating CP as "the whole system goes down during any partition," when the quorum math above shows it's specifically the minority side, and only for writes, reads can often still be served (possibly stale, depending on the read-consistency level chosen). The common mistake on the AP side is stopping at "it's eventually consistent" without actually building the reconciliation path: leaderless ingestion without idempotent consumers and deduplication just relocates the correctness problem downstream instead of solving it. And even a CP ledger still needs idempotency keys on the client side: a client that times out waiting for a quorum ack and retries the same transaction must not have it applied twice, which is a consistency concern CAP itself doesn't cover but that any real ledger design has to handle regardless of which side of CAP it lands on.
For availability targets of 99.9%, 99.99%, and 99.999%, calculate the allowed downtime per year and per month for each. Then walk through what architectural changes actually get you from one tier to the next.
Sample Answer
Direct answer: Availability is the fraction of time a system is usable, and "N nines" is shorthand for how close that fraction is to 100%. Going from 99.9% to 99.99% to 99.999% shrinks allowed downtime by roughly 10x at each step, and each step also costs roughly an order of magnitude more in engineering and infrastructure, because you're eliminating an entire category of failure (single-host, then single-zone, then single-region) rather than just adding more of the same redundancy.
Structured elaboration
Allowed downtime per year is derived from the availability target directly:
downtimeyear=(1−A)×8760 hourswhere 8760 is the number of hours in a 365-day year (24 x 365), and per-month downtime uses 730 hours (8760 / 12):
downtimemonth=(1−A)×730×60 minutesPlugging in each target:
| Availability | Allowed downtime / year | Allowed downtime / month |
|---|---|---|
| 99.9% ("three nines") | (1−0.999)×8760=8.76 hours -> 8h 45m 36s | (1−0.999)×730×60=43.8 min |
| 99.99% ("four nines") | (1−0.9999)×8760=0.876 hours -> 52m 34s | (1−0.9999)×730×60=4.38 min |
| 99.999% ("five nines") | (1−0.99999)×8760=0.0876 hours -> 5m 15s | (1−0.99999)×730×60=0.438 min -> 26s |
Each jump divides allowed downtime by exactly 10, because each availability target divides (1−A) by 10.
What actually changes architecturally between tiers
- 99.9% -> 99.99%: eliminate single points of failure inside one facility. Multi-AZ deployment, N+1 redundancy (one extra standby unit beyond what's strictly needed to handle normal load, so a single failure doesn't drop capacity below what's required) on stateful components (load balancers, databases with a standby replica), automated health-check-driven failover, and a real on-call rotation with paging. Most of the gain here comes from removing manual recovery steps: a human restarting a service takes minutes and that alone can burn the entire four-nines monthly budget.
- 99.99% -> 99.999%: eliminate the facility (zone or region) itself as a single point of failure. Multi-region active-active or hot standby (a fully-running backup kept ready to take over instantly, unlike a cold standby that would first need to be started up and warmed), automated cross-region failover (not human-triggered), synchronous or tightly-bounded-lag replication for the data that must survive a region loss, and rigorous testing of the failover path itself (chaos drills: deliberately triggering the failover in a controlled test so a broken failover path is discovered on a Tuesday afternoon, not during a real outage), because at this tier the failover mechanism is now a bigger risk to availability than the failures it's protecting against.
- Beyond 99.999%, the limiting factor usually isn't infrastructure, it's deployment risk (bad releases) and dependency risk (a vendor or DNS provider you don't control), so the remaining budget goes to progressive rollouts, fast automated rollback, and reducing the number of hard external dependencies on the critical path.
Worked example: composing a dependency chain
A request that serially depends on a load balancer (99.99%), an app tier (99.95%), and a database (99.99%) has a combined availability equal to the product of the individual availabilities, because all three must be up simultaneously:
Aserial=0.9999×0.9995×0.9999=0.9993That's roughly 99.93%, worse than any single component, which is why a system built entirely from 99.99%-rated pieces chained together does not automatically deliver 99.99% end to end. Adding a redundant standby database (parallel, either one being up is sufficient) with independent 99.99% availability changes only that term:
Adb,pair=1−(1−0.9999)2=1−0.00012=0.99999999so the pair is effectively always up, and the chain's availability is then bounded by the weakest remaining serial link (the app tier at 99.95%), not the database.
Trade-offs & pitfalls
- Availability composes multiplicatively across a serial chain and the weakest link dominates: chasing five nines on your database while your app tier sits at three nines is wasted spend.
- Each nine costs disproportionately more: 99.9% to 99.99% is mostly process and automation (cheap-ish); 99.99% to 99.999% usually means paying for a second region and the operational overhead of keeping it truly independent (expensive, and dangerous if the failover path itself is untested).
- Downtime budgets don't distinguish planned from unplanned; a team that spends its whole error budget on deploy-related outages hasn't actually built a more resilient system, just a riskier release process.
- A very common interview trap: treating "99.99% uptime" as a promise about any single request rather than a time-integrated average. A system can meet 99.99% for the year while having a full 50-minute outage in one bad afternoon.
A downstream service you depend on starts responding slowly, and requests to it start backing up on your side, growing queues and increasing latency. Walk through your immediate mitigations and your longer-term architectural fix, and explain the trade-off each one introduces.
Sample Answer
Direct answer: The immediate priority is to stop the slowdown from consuming your own resources: set aggressive timeouts, open a circuit breaker so you stop calling the failing dependency, and isolate the connection/thread pool used for that call so it can't starve everything else. The longer-term fix is architectural: decouple the caller from the dependency's latency entirely, usually via an async queue or by making the call non-blocking, so a slow downstream degrades throughput instead of taking the whole service down with it.
Structured elaboration
Why this happens (the mechanism): by Little's Law, the number of requests in flight L equals arrival rate λ times the time each request spends in the system W: L=λW. If a downstream call's latency goes from 50ms to 500ms while your request rate stays at, say, 200 requests/second, the in-flight count grows from L=200×0.05=10 to L=200×0.5=100, a 10x increase, purely from the latency change with no change in incoming traffic. If your thread or connection pool was sized for ~10-20 concurrent in-flight requests to that dependency, it's now exhausted, and requests start queueing on your side, which is exactly the symptom described.
Immediate mitigations (minutes, not a redesign):
| Mitigation | What it does | Trade-off it introduces |
|---|---|---|
| Tight timeouts | Caps how long you'll wait, preventing unbounded queue growth | Cuts off requests that might have succeeded a moment later; needs to be shorter than your own SLA to the caller |
| Circuit breaker | Stops calling the dependency once error/latency crosses a threshold, failing fast instead of queueing | Can trip on transient blips if thresholds are too sensitive; denies service even to calls that might succeed |
| Bulkhead (isolated pool) | Gives this dependency its own thread/connection pool so its slowdown can't exhaust pools shared by healthy dependencies | Reduces pooled efficiency (can't borrow capacity across dependencies); requires knowing sizing up front |
| Load shedding / fast 503 | Rejects excess requests immediately when queue depth crosses a threshold, protecting the instances still healthy | Directly reduces availability for shed requests; needs to shed selectively, not randomly, if some requests matter more |
Longer-term architectural fix:
- Decouple via an async queue: put a durable queue between the caller and the slow dependency so the caller can return quickly (accept-and-acknowledge) and the dependency is drained at its own sustainable pace, rather than the caller blocking on it synchronously. Trade-off: the caller can no longer return a synchronous success/failure for that operation; the interaction model has to change to something the client and product can tolerate (a "pending" state, a webhook, a poll).
- Idempotent retries with backoff and jitter: if retries are needed, they must be capped, exponential, and jittered so a fleet of callers doesn't retry in lockstep and re-create the exact overload it's recovering from. Trade-off: added complexity, and retries must be provably idempotent on the downstream side or they risk duplicate side effects.
- Capacity planning against the tail, not the average: provision the dependency (or the pool sized to call it) based on observed p99 latency, not p50, since it's the tail that determines when queues start building. Trade-off: costs more standing capacity for headroom that's idle most of the time.
Applying this to concrete variants of the same pattern: the reasoning above is the same whether the slow dependency is a payment-validation service (immediate: circuit breaker + fast-fail with a clear "try again" to the user rather than a silent hang; long-term: async payment confirmation via webhook), a message-queue consumer falling behind (immediate: shed or dead-letter the oldest low-priority messages, bulkhead the consumer pool by message type; long-term: scale consumers horizontally and partition by priority), a retry storm from a flood of client-side 503s (immediate: the client-side backoff-with-jitter above is the direct fix; long-term: make the shedding threshold adaptive so it doesn't itself become the trigger for a thundering herd), or a synchronous order-processing pipeline backing up (immediate: bulkhead the slow stage's pool; long-term: convert that stage to the async-queue pattern above).
Trade-offs & pitfalls
- Every immediate mitigation above trades some availability or correctness for stability: timeouts drop requests that might have succeeded, circuit breakers deny service during their open window, load shedding sacrifices some requests to save the rest. The point isn't to avoid the trade-off, it's to make it deliberately and visibly rather than let an unbounded queue make it for you via an eventual crash.
- A common wrong turn: adding retries as the first response to a slowdown. Naive retries without backoff amplify load on an already-struggling dependency and can turn a partial slowdown into a full outage (a retry storm).
- Circuit breakers and bulkheads need to be tuned against real traffic and latency distributions; thresholds copied from a different service's runbook are a common source of either false trips (unnecessary unavailability) or no protection at all (thresholds too loose to matter).
Design a DR architecture for a customer-facing web application that needs to meet a 99.99% SLA and an RPO under 5 minutes. Walk through your redundancy strategy, database replication approach, and failover mechanics.
Sample Answer
The design combines two layers that solve different failure modes: multi-AZ redundancy within a primary region to survive the common case (a host, rack, or single-AZ failure) with near-zero disruption, and a warm standby in a second region, kept current via continuous replication, to survive the rare but severe case (a full regional outage). Trying to hit 99.99% and a 5-minute RPO (recovery point objective: the maximum data, measured in time since the last durable copy, you can afford to lose) with only one of those layers doesn't work: multi-AZ alone doesn't survive a regional disaster, and cross-region alone (without in-region redundancy) means every AZ blip forces an unnecessary cross-region failover.
Architecture
flowchart TB
subgraph Primary["Primary region (us-east-1)"]
LB1[Load Balancer]
AZ1[App tier - AZ1]
AZ2[App tier - AZ2]
AZ3[App tier - AZ3]
DBP[(Primary DB\nsynchronous multi-AZ)]
LB1 --> AZ1 & AZ2 & AZ3
AZ1 & AZ2 & AZ3 --> DBP
end
subgraph DR["DR region (eu-west-1)"]
LB2[Load Balancer - standby]
AZS[App tier - warm, scaled down]
DBS[(Standby DB replica)]
LB2 --> AZS --> DBS
end
DBP -- "async replication\n(CDC / log shipping)" --> DBS
DNS[DNS / Global routing] --> LB1
DNS -. failover .-> LB2
Redundancy strategy. Inside the primary region, the app tier runs across three AZs behind a load-balanced, health-checked pool (an AZ failure just removes that AZ's capacity from rotation), and the database runs synchronous or semi-synchronous multi-AZ replication so an AZ failure promotes a same-region replica with no data loss. This layer alone handles the overwhelming majority of real infrastructure failures.
Database replication approach. In-region: synchronous, to keep in-region failover at RPO≈0. Cross-region to the DR site: asynchronous (via change-data-capture or log shipping), because synchronous cross-region commits would add tens of milliseconds of write latency for every transaction just to protect against an event, a full regional outage, that's far rarer than an AZ blip. The async lag is what has to stay under the 5-minute RPO budget, so it's monitored continuously with an alert threshold well below 5 minutes (for example, paging at 60 seconds of lag) so operators have room to react before the budget is actually at risk.
Failover mechanics. In-region (AZ failure): automatic, health-check-driven, sub-minute, no human involved. Cross-region (regional failure): health checks on the primary region trigger an automated runbook that promotes the DR replica to primary, scales the warm (already-running, just smaller) app tier in the DR region up to full capacity, and updates DNS/global routing to point at the DR region's load balancer. The DR app tier is warm, not cold, specifically so this promotion is a scale-up operation, not a from-scratch provision, which is what keeps the cross-region RTO (recovery time objective: how long restoring service takes after a failure) in the minutes range instead of the tens-of-minutes range a cold standby would need.
Worked example: why the SLA is bounded by the DR path, not the multi-AZ path
Pin illustrative per-AZ availability at 99.9% for a single AZ's app-tier capacity. With three independent AZs, the app tier only fails when all three fail simultaneously:
1−A3-AZ=(1−0.999)3=10−9⟹A3-AZ≈99.9999999%That's a downtime budget of about 525,600×10−9≈0.0005 minutes/year from independent AZ failures alone, far inside the 99.99% target's 52.56-minute/year budget. The conclusion this points to: independent-component math says multi-AZ redundancy alone should comfortably clear 99.99%, so in practice the SLA is bounded by things that math doesn't model, correlated failures, a bad deploy, a full regional outage, human error, not by an AZ dying in isolation. That's precisely why the DR region exists: it's insurance against the failure modes the availability formula can't see, not an incremental nudge on an already-tiny independent-failure number.
Trade-offs and pitfalls
The most common wrong turn is treating the independent-AZ math above as if it were the whole SLA story and skipping the cross-region layer as "redundant redundancy"; it's protecting against a different, uncorrelated failure class entirely. A second pitfall is under-provisioning the warm standby: if "warm" means a token instance that can't actually absorb full production load once promoted, the design has a theoretical RTO that doesn't survive contact with a real regional failover, so the standby's scaled capacity needs to be load-tested, not just health-checked. A concrete way this shows up: if the actual mandate is something like "cut our current 4-hour cold-standby RTO down to 15 minutes on a fixed budget," the honest trade-off conversation is that a full active-active design gets there fastest but costs the most and adds multi-writer complexity, while a well-automated warm standby (this design) gets to a low-single-digit-minutes RTO for a fraction of the cost, which is usually the better fit unless the traffic is latency-sensitive and truly global. Finally, none of this replaces DR drills: an automated runbook that's never been exercised against a real regional cutover is the same "theoretical RTO" problem as an untested pilot-light tier, just at higher stakes.
What does 'blast radius' mean when you're talking about a production failure? Name a few concrete engineering practices that reduce it, and what that costs you.
Sample Answer
Direct answer
Blast radius is the scope of impact when a component fails: how many users, tenants, or dependent services are affected, and how severely, not just whether the failure happened at all. Reducing blast radius means designing so a single failure touches the smallest possible slice of the system, which makes outages smaller, easier to detect, and faster to recover from, even if it doesn't reduce how often failures happen at all.
Practices that reduce it, and what they cost
| Practice | How it shrinks blast radius | What it costs |
|---|---|---|
| Circuit breakers | Stop repeated calls to a failing dependency, isolating the failure to the caller instead of letting it spread | Added latency and complexity in the failure path; a poorly tuned breaker can trip on transient blips |
| Finer-grained service decomposition | A failure or overload in one bounded service only affects its own consumers, not unrelated functionality | More services to deploy, monitor, and operate; cross-service calls add their own new failure modes |
| Bulkheads (per-tenant or per-dependency resource pools) | One tenant's or one dependency's exhaustion doesn't consume capacity meant for everyone else | More total resources provisioned (dedicated pools cost more than one shared pool sized for the average case) |
| Traffic shaping and rate limits | Caps how much load a single misbehaving client or spike can push into downstream systems | Legitimate bursty clients can get throttled unless limits are tuned carefully |
Worked example
Consider a service with 1,000 tenants sharing a single connection pool. If that pool exhausts, every tenant is affected. Now split that same total capacity into 10 isolated pools of 100 tenants each, so each pool serves 100 of the 1,000 tenants and only that pool's own tenants are affected if it exhausts:
1,000100=10% of tenants affected (isolated pools)vs.100% (shared pool)Splitting the same total capacity into 10 pools of 100 tenants each means a single pool's exhaustion now affects only 100 of the 1,000 tenants, 10% of the blast radius of the shared-pool design, for the same total resources. The cost is operational: 10 pools to monitor and size instead of one, and if traffic isn't evenly distributed across tenants, some pools may be under-utilized while others are tight, which the shared pool didn't have to worry about.
Trade-offs & pitfalls
Reducing blast radius is generally a trade of operational complexity and some resource inefficiency for smaller, more contained failures; it doesn't reduce the underlying failure rate of any individual component. The common mistake is treating blast-radius reduction as free: partitioning by tenant, region, or dependency multiplies the number of things to monitor and can hide a systemic bug (one that affects every partition equally) behind what looks like ten separate, unrelated small incidents instead of one clearly systemic one.
Unlock Full Question Bank
Get access to all 45 Fault Tolerance, High Availability, and Disaster Recovery interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.