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.
Define cascading failure and walk through a realistic example: service C fails, B (which depends on C) gets overloaded, and A (which depends on B) starts degrading too. At each layer, what protection would you put in place to stop the cascade from propagating?
Sample Answer
Direct answer
A cascading failure is when one component's failure increases load or latency on the components that depend on it, and that increased load causes those components to fail too, propagating outward until a large part of the system is affected, even though only one component actually broke in the first place. The mechanism is almost always resource exhaustion: threads, connections, or memory tied up waiting on the failed component instead of being freed quickly.
Walkthrough: C fails, B overloads, A degrades
flowchart LR
A[API Gateway] -->|rate limit and timeout| B[Order Service]
B -->|bulkhead pool: payments| C[Payment Service]
C -.fails.-> B
B -->|circuit breaker opens| D[Fallback: queue order for async retry]
A -->|circuit breaker opens| E[Fallback: 503 with Retry-After]
B -->|isolated pool: other deps unaffected| F[Inventory Service]
- C (Payment Service) fails, hanging instead of returning errors quickly, perhaps due to a downstream outage of its own.
- B (Order Service) calls C without a tight timeout. Each call to C now blocks for far longer than normal, tying up a thread or connection from B's pool for the duration.
- B's resource pool exhausts. As more requests arrive at B, more threads get stuck waiting on C, until B has no capacity left to serve any request, including ones that don't even touch C.
- A (API Gateway) calls B, and B is now slow or unresponsive for everything, so A's calls to B start timing out or queueing too, degrading A's own capacity in turn.
Worked example: how fast does B's pool actually exhaust?
Little's Law relates the number of requests in flight to the arrival rate and the time each spends being processed:
L=λWSay B receives 500 requests per second, and under normal conditions each call to C takes 50ms:
Lnormal=500×0.05=25 concurrent in-flight requests25 concurrent requests is a light load on a typical connection pool. Now C hangs, and B's HTTP client has no explicit timeout of its own, falling back to a default of 30 seconds:
Lfailure=500×30=15,000 concurrent in-flight requests neededIf B's thread pool has 200 threads, the time to exhaust it entirely is:
texhaust=500200=0.4 sUnder 400 milliseconds. That's how quickly a single hung dependency with no timeout turns into total unavailability for a service handling 500 requests per second: the pool never gets close to steady-state at the 30-second hang time, it simply fills with stuck requests almost instantly and stays full.
Protections at each layer
- At B, calling C: a tight, explicit timeout (measured in low hundreds of milliseconds, not the client library's 30-second default) so a hung call fails fast and frees the thread quickly; a circuit breaker that opens after a run of failures or timeouts, so B stops even attempting calls to C once it's clearly down, and falls back to queueing the order for later processing; a bulkhead, a dedicated connection pool just for calls to C, so exhaustion from C-related calls doesn't consume the threads B needs to serve requests that don't touch C at all (like inventory checks).
- At A, calling B: the same pattern one layer up, a timeout on calls to B, a circuit breaker that trips once B's error rate or latency crosses a threshold, and a fallback (a fast 503 with
Retry-Afterrather than a hung request) so A's own capacity isn't consumed waiting on a B that's already struggling.
Trade-offs & pitfalls
Timeouts that are too aggressive cause false-positive failures under normal, brief latency variance; timeouts that are too loose don't prevent the cascade fast enough, as the Little's Law example shows. Bulkheads cost real resources (a dedicated pool per dependency uses more total connections or threads than one shared pool) in exchange for isolation, so they're worth applying to the dependencies most likely to fail or most likely to take down unrelated traffic if they do. The most common mistake is only protecting the first hop (B to C) and assuming that's sufficient; as the walkthrough shows, without protection at the A to B hop too, the failure still reaches A once B is degraded, just one layer later.
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 idempotency mean in the context of retries, and why does it matter? Walk through how you'd make a payment-creation endpoint safe to retry, including how you'd handle the idempotency key.
Sample Answer
Direct answer
Idempotency means performing the same operation multiple times has the exact same effect as performing it once. It matters for retries because network failures make it impossible for a client to reliably tell "the request failed" apart from "the request succeeded but the response was lost"; without idempotency, a client that retries after a timeout risks creating a duplicate side effect, like charging a customer twice for one order.
Making a payment-creation endpoint safe to retry
The standard mechanism is a client-generated idempotency key attached to the request:
- The client generates a unique key (a UUID) once per logical operation, before the first attempt, and sends it on every retry of that same logical operation in a header such as
Idempotency-Key. - The server does an atomic check-and-set against a persistent store keyed by that idempotency key: if the key is new, it proceeds with the charge; if the key already exists, it returns the previously stored result instead of processing the charge again.
- The check-and-set has to be atomic (a single transactional operation, not a read followed by a separate write) so two near-simultaneous retries can't both see "key doesn't exist" and both proceed.
- The stored result includes enough to reconstruct the original response (status, charge ID, amount) and a status field (
in_progress,succeeded,failed) so a retry that arrives while the first attempt is still executing gets told to wait or gets the eventual result, rather than racing ahead. - Keys are kept with a bounded TTL (commonly 24 to 72 hours) since indefinite retention is unnecessary once a client has almost certainly given up retrying, and TTL bounds the storage cost of the idempotency table.
sequenceDiagram
participant Client
participant API as API Region A
participant Store as Idempotency Store
participant PG as Payment Gateway
Client->>API: POST charges Idempotency-Key K1
API->>Store: check-and-set K1 in_progress
Store-->>API: new key proceed
API->>PG: create charge
PG-->>API: charge succeeded
API->>Store: save result for K1
API-->>Client: 200 OK response lost in transit
Client->>API: retry POST charges Idempotency-Key K1
API->>Store: check K1
Store-->>API: found status succeeded
API-->>Client: 200 OK cached result no new charge
Worked example
In the sequence above, the server actually completes the charge and writes the success result to the idempotency store, but the client's connection drops before the 200 response arrives, so from the client's point of view the request timed out. The client retries with the same key K1. The server's check-and-set finds K1 already marked succeeded, so it returns the stored response (the original charge ID and amount) directly and never calls the payment gateway again. Exactly one charge exists, regardless of how many times the client retries.
Harder extension: retries across a cross-region failover
The same duplicate-request risk gets worse if the retry lands on a different region than the original attempt. Say the first request goes to Region A, and before the response comes back, DNS or Anycast reroutes the client (as part of a regional failover) so the retry with the same idempotency key goes to Region B. If the idempotency store is region-local and not replicated, Region B has never heard of K1, sees it as a new key, and processes a second charge, exactly the failure the mechanism was supposed to prevent.
The fix is that the idempotency store itself has to be as available and as replicated as the failover design assumes the rest of the system is: either a globally consistent store (accepting the added write latency) or, more commonly for payments specifically, delegating idempotency to the payment gateway itself, which usually supports its own idempotency keys and is already a single global system of record regardless of which region initiated the call. Relying on the gateway's own dedupe as the backstop means even a fully region-local idempotency store failing open during a failover doesn't result in a real double charge.
Trade-offs & pitfalls
Idempotency keys add a write to the hot path (the check-and-set) and a storage system that has to be highly available, since if the idempotency store itself is down, you're forced to choose between blocking the write entirely or risking a duplicate. TTL choice is a real trade-off: too short and a legitimately slow client retry after the TTL expires creates a duplicate; too long and the storage grows unnecessarily and stale in-progress records from crashed requests linger. The most common mistake is only deduplicating the write itself while forgetting downstream side effects (an email receipt, a webhook fired to a third party) that happen inside the same logical operation and need to be gated by the same check, not fired unconditionally every time the handler runs.
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 44 Fault Tolerance, High Availability, and Disaster Recovery interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.