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.
Most of your traffic is reads, but you occasionally get writes from any region, and you want to route reads to the nearest region for latency. Walk through the replication and consistency strategy that makes this work.
Sample Answer
Direct answer
Deploy a read replica in every user-facing region and route reads to the nearest one for latency, while anchoring writes to a single-writer-per-shard model: each account or entity has one home region that owns writes for it (sharded by key, not globally centralized), and a write originating from any other region gets forwarded to that entity's home region. Reads stay fast everywhere because they never leave the local region; writes pay a forwarding cost only when they originate somewhere other than the entity's home region, which for most workloads is the minority case.
Architecture
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
- Read routing: DNS-based or edge-proxy latency routing sends each client to its nearest region; that region serves reads from its local replica, optionally backed by a local edge cache for hot keys to cut load further.
- Write routing: a lightweight write-forwarder in each region inspects the target entity's shard key, determines which region owns it, and forwards the write there if it isn't local; the write commits in its home region and the forwarder returns the result (or an idempotency-tracked async acknowledgment) to the originating client.
- Replication: the write's home region streams committed changes via change-data-capture (CDC) to every read replica, asynchronously; replicas apply changes in the order the CDC stream delivers them, tracking a per-record last-writer timestamp so replay order and true causal order stay consistent.
Consistency model
This design is eventually consistent for reads: a read served from Region B immediately after a write committed in Region A's home shard may not reflect that write yet, bounded by CDC replication lag rather than by any hard guarantee. That's an explicit, necessary trade for the latency goal, since making every read wait for a cross-region round-trip to confirm it has the absolute latest value would defeat the entire point of routing reads to the nearest region. Where a specific read genuinely needs to see its own very-recent write (a user immediately viewing an item they just created), the standard fix is read-your-own-writes: route that specific read to the entity's home region (or to a replica known to have caught up past a specific CDC watermark) instead of the nearest replica, rather than weakening the consistency model for every read to satisfy the rare case.
Write conflicts are structurally rare by design, because each entity has exactly one home region and therefore exactly one writer at any time; there's no multi-master merge problem to solve because there's no multi-master. The forwarding hop is the cost of ruling that problem out entirely rather than solving it after the fact.
Worked example: tracing one write end to end
Pin a concrete case: user_id=482's home region is us-east (that's where its shard's writer lives), but the client happens to be connected to the nearest edge in eu-west. The eu-west write-forwarder inspects the shard key for user_id=482, sees it belongs to us-east, and forwards the write there; assume a cross-region round trip of about 90ms for that forward, so the write is accepted and committed in us-east roughly 95ms after the client sent it (90ms network plus a small local processing cost). From there, the CDC stream carries that committed change out to every read replica asynchronously; assume that hop adds about 300ms before eu-west's local replica has applied it (its own network hop plus normal stream batching, separate from the synchronous forward that carried the write there). So the total time from write-acceptance to eu-west's replica reflecting the change is about 95ms+300ms=395ms, call it roughly 400ms. A read served from eu-west's local replica 1 second (1,000ms) after the write was accepted already reflects it, comfortably past the ~400ms it takes to land; a read served only 50ms after acceptance would not yet reflect it, since 50ms is well inside that ~400ms propagation window, which is exactly the read-your-own-writes gap the design accepts and works around by routing that specific kind of read to the home region instead.
Trade-offs & pitfalls
The biggest latency cost this design accepts is on writes that originate far from an entity's home region: a user in Region C writing to an entity whose home shard is in Region A pays a full cross-region round trip for that write, even though every other user's reads and most other writes stay fast. If write locality doesn't naturally match user geography (an entity created in one region gets written to mostly by users somewhere else over time), this cost compounds instead of amortizing away, which is worth checking against real traffic patterns before committing to a static shard-to-region mapping. Replication lag is the other pitfall: CDC-based replication is asynchronous by nature, so a region that falls behind (network partition, replica overload) serves increasingly stale reads without necessarily surfacing an error, which is why lag needs to be an actively monitored metric with alerting, not just an assumed-small property of the pipeline. Finally, resist the temptation to solve the "occasional write from any region" requirement with full multi-master writes accepted locally everywhere; that reintroduces exactly the conflict-resolution complexity (concurrent writes to the same entity from two regions, needing merge logic or last-write-wins with its own correctness risks) that single-writer-per-shard was specifically chosen to avoid, in exchange for a write-latency win that the stated 90-percent-read workload doesn't actually need.
What's the trade-off between designing for high availability across multiple availability zones within one region versus spreading across multiple regions entirely? Walk through failover time, replication latency, and operational complexity for each.
Sample Answer
Direct answer
Multi-AZ (spreading across availability zones within one region) gives faster failover and lower replication latency because the zones are close together on a fast, low-latency network, but it only protects against zone-level failures, not a regional outage. Multi-region gives protection against a much larger class of failures (an entire region going down) at the cost of higher replication latency, slower failover, and meaningfully more operational complexity.
Comparison
| Dimension | Multi-AZ | Multi-region |
|---|---|---|
| Fault domain protected | Power, rack, single-zone network failure | Entire region: cloud provider outage, regional disaster |
| Typical failover time | Tens of seconds (load-balancer health checks converge quickly on a local network) | Tens of seconds to a few minutes (DNS propagation, cross-region orchestration) |
| Replication mode | Synchronous or near-synchronous is feasible (low round-trip time) | Usually asynchronous (round-trip time too high for sync writes without unacceptable latency) |
| RPO | Near zero, achievable with sync replication | Non-zero, bounded by async replication lag |
| Data residency | Straightforward, all data stays in one region/country | Requires deliberate placement to satisfy jurisdictional rules |
| Network/DNS complexity | Simple, one VPC, local load balancer | Geo-DNS or global load balancer, cross-region routing, split-brain avoidance |
Worked example: failover time
For a multi-AZ setup, failover time is dominated by load-balancer health-check detection: a 10-second check interval with a 3-consecutive-failure threshold to avoid flapping, plus roughly 10 seconds of connection draining:
tAZ=(10×3)+10=40 sFor a multi-region setup using DNS-based failover, detection can use a tighter 5-second health-check interval (since it's typically backed by dedicated external health checkers rather than the load balancer itself), still with a 3-consecutive-failure threshold, plus a DNS TTL of 30 seconds that a client's resolver may need to fully expire before it re-resolves to the healthy region:
tregion=(5×3)+30=45 sThese land in a similar ballpark for a well-tuned setup on paper, but the multi-region number is the optimistic case: it assumes every client's resolver honors the 30-second TTL exactly, which enterprise resolvers and some mobile carriers don't always do, so real-world multi-region failover tail latency for the slowest clients is meaningfully worse than the AZ case, even though the median case is comparable.
Trade-offs & pitfalls
Choose multi-AZ when the goal is protecting against common, smaller-scale infrastructure failures with minimal added latency and complexity; choose multi-region when the requirement is surviving a full regional outage or meeting geographic redundancy or compliance mandates, and accept that this means designing for a non-zero RPO and a longer failover tail. A frequent mistake is defaulting straight to multi-region for availability reasons alone without first exhausting what multi-AZ already buys, since multi-region adds real replication-latency and consistency costs that many systems don't actually need to pay to meet their real availability target.
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.
Compare the standard DR strategy tiers: backup-and-restore, pilot light, warm standby, and active-active multi-site. For each, what's the typical RTO/RPO range, and what does it cost you?
Sample Answer
The four standard DR tiers form a spectrum from cheapest-and-slowest to most-expensive-and-fastest, and each one trades infrastructure spend for recovery speed (RTO, recovery time objective: how long restoring service takes) and data freshness (RPO, recovery point objective: how much data, measured in time, you could lose): backup-and-restore keeps only backups running, pilot light keeps a minimal always-on core, warm standby keeps a scaled-down full copy running, and active-active multi-site keeps a full copy running and serving live traffic.
Comparing the four tiers
| Tier | What's running in DR | Typical RTO | Typical RPO | Relative cost |
|---|---|---|---|---|
| Backup-and-restore | Nothing; only backups exist in storage | Hours to a day+ (provision infra, restore data) | Hours (since the last backup) | Lowest: storage cost only |
| Pilot light | Core data store kept replicated and running; app/compute layer absent until needed | Tens of minutes to a few hours (scale up compute, deploy app) | Minutes (continuous replication to the core) | Low-moderate: one small always-on component |
| Warm standby | A scaled-down but fully functional copy of the whole stack, running continuously | Minutes (scale up capacity, redirect traffic) | Seconds to low minutes (near-real-time replication) | Moderate-high: a live, if smaller, second environment |
| Active-active multi-site | Full-scale copy in both/all sites, serving live traffic simultaneously | Near-zero (traffic reroutes, nothing to "start") | Near-zero to seconds (synchronous or tightly-bounded async replication) | Highest: full duplicate capacity plus distributed-write complexity |
The RTO/RPO ranges above are the typical shape of the trade-off, not a fixed number for any specific system: the exact figures depend on data volume, automation maturity, and how the replication is actually implemented within each tier.
Worked example: a budget-constrained startup
A mid-sized SaaS with a fixed infrastructure budget doesn't have to pick one tier for the whole system; the standard move is to mix tiers by criticality. Say the product has three logical components: authentication/billing (must never meaningfully go down, since it blocks every paying customer from doing anything), the core application (needs to come back reasonably fast but a short outage is tolerable), and internal admin tooling (only the ops team notices if it's down for a few hours).
A budget-conscious allocation: active-active for auth/billing (the one component where the cost premium is justified because its outage blocks revenue entirely, and it's usually small enough in infrastructure footprint that duplicating it fully is affordable), pilot light for the core application (keep the database replicated continuously so RPO stays low, but only spin up the app-server fleet in DR when actually needed, since that's the majority of the compute cost), and backup-and-restore for admin tooling (cheapest tier, acceptable because nobody customer-facing is blocked by it being down for hours). This gets the highest-blast-radius component the fastest recovery while keeping the overall DR bill proportional to what each component actually costs the business if it's down, instead of buying active-active everywhere by default.
Trade-offs and pitfalls
The most expensive mistake in this space isn't picking the "wrong" tier, it's picking a tier and never testing failover into it: a pilot-light setup that's never actually been promoted to full capacity under load is a theoretical RTO, not a real one, and the first real DR event is a bad time to discover the app layer doesn't actually scale up cleanly from zero. A related pitfall is under-provisioning a warm standby's capacity: "scaled down" often means it can absorb DR traffic at reduced performance, and teams sometimes forget to validate that the scaled-down size can actually handle 100% of production load once promoted, not just serve health checks. Finally, active-active's real cost isn't just the duplicate infrastructure line item, it's the ongoing engineering cost of keeping a multi-writer data model correct, which is easy to underestimate when comparing tiers purely on an RTO/RPO/dollar table.
Two regions running asynchronous replication get network-partitioned, and both keep accepting writes. When the partition heals, how do you detect the divergence and reconcile the conflicting writes?
Sample Answer
Direct answer
This is the multi-master conflict case: because replication was asynchronous and both regions kept accepting writes, there are now two divergent histories that both look locally valid. Detecting divergence means comparing state between regions, via hashes or version metadata, once connectivity returns, not waiting for a customer to report bad data. Reconciling means classifying each conflict by whether it's safely auto-mergeable or needs a compensating action; financial and other invariant-sensitive writes should never be silently overwritten.
Structured elaboration
Detecting divergence. Attach causal metadata to every write, a vector clock (an array with one counter per region, incremented only when that region processes a write, so comparing two vectors shows whether one write causally happened after the other, every counter at least as high, or whether the two happened independently) or at minimum a per-region monotonic sequence number plus timestamp, so that on reconnect two versions of the same key can be compared to tell, mathematically, whether one supersedes the other or whether they are genuinely concurrent, meaning both sides wrote independently during the partition.
sequenceDiagram
participant RegionA
participant RegionB
participant Reconciler
Note over RegionA,RegionB: Partition: both regions accept writes independently
RegionA->>RegionA: Write key X, vector clock [A:5,B:3]
RegionB->>RegionB: Write key X, vector clock [A:4,B:4]
Note over RegionA,RegionB: Partition heals
Reconciler->>RegionA: Fetch change log with vector clocks
Reconciler->>RegionB: Fetch change log with vector clocks
Reconciler->>Reconciler: Compare vector clocks, classify conflicts
Reconciler->>Reconciler: Auto-merge commutative writes
Reconciler->>Reconciler: Apply compensating transaction for invariant conflicts
Reconciler->>RegionA: Apply resolved state
Reconciler->>RegionB: Apply resolved state
Reconciling. Not every conflict resolves the same way:
| Conflict type | Resolution |
|---|---|
| Non-overlapping keys, touched on only one side | Trivial union merge, no real conflict |
| Commutative or CRDT-safe operations (counter increment, set-add). CRDT: conflict-free replicated data type, a data structure designed so concurrent updates always merge to the same result automatically, no coordination needed | Merge automatically, order doesn't matter |
| Same-field conflicting writes with no safe merge rule | Deterministic policy where business semantics allow it, such as latest-timestamp-wins, logged to a resolved-log for audit |
| Financial or invariant-sensitive conflicts | Never blind-overwrite; apply a compensating transaction, an explicit correcting entry linked to the original, so history stays auditable and reversible |
Worked example: detecting a real conflict with vector clocks
Two regions independently update the same key, "user:42.email," during the partition. Region A's write carries vector clock VA=[A:5,B:3], Region B's write carries VB=[A:4,B:4].
A vector clock V1 dominates V2, meaning V1 happened after V2 and can safely overwrite it, only if every component of V1 is greater than or equal to the corresponding component of V2, with at least one strictly greater. Checking:
VA[A]=5>VB[A]=4butVA[B]=3<VB[B]=4Neither vector dominates the other, A leads on its own component, B leads on its own, so this is a genuine concurrent conflict, not a case where one write simply supersedes the other, and it must go through the merge policy rather than being resolved by picking the higher-looking clock.
Worked example: scoping the reconciliation workload
Assume the partition lasted 10 minutes, each region accepted 150 writes/sec across a 5-million-key space:
n1=n2=150×600=90,000 writes per side E[colliding keys]=Nn1×n2=5,000,00090,000×90,000=5,000,0008,100,000,000=1,620Of 180,000 total writes across both sides, an estimated 1,620 keys, 0.9%, are true conflicts needing the merge policy; the rest merge by simple union. Sizing the reconciliation pipeline, and any manual-review queue, around this number, not around all writes during the partition, is what keeps the process fast.
Trade-offs & pitfalls
- Silently applying last-writer-wins to every conflict is the most common shortcut, and it's the wrong default for anything with a business invariant, balances, inventory counts: it loses data without a trace and without anyone knowing which write "won."
- Scanning the entire keyspace for divergence, instead of the range touched during the outage window, wastes time and delays recovery.
- A merge policy that isn't tested against the conflict-type classification before the incident means the policy is being designed live, under pressure, which is when mistakes compound.
- Compensating transactions must themselves be idempotent and auditable; a reconciliation "fix" that isn't traceable back to the original conflicting writes creates a second, harder-to-debug data quality problem.
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.