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.
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.
Data residency requirements mean you can't freely replicate a customer's data across borders. Walk through how you'd design a multi-region architecture that still gives you reasonable availability within those constraints.
Sample Answer
Direct answer
Default to regional partitioning: keep each country or region's primary data store entirely inside that region's boundary, route each user's traffic to their home region, and use a lightweight global control plane (metadata, feature flags, non-personal configuration) that's free to replicate everywhere because it doesn't contain the restricted data. Availability within a residency-constrained region then comes from the same tools you'd use anywhere: redundant replicas and a fast local failover path, just constrained to stay inside the boundary rather than failing over to a healthy replica in a different country.
Design options and their trade-offs
| Option | How it works | Availability profile | When it fits |
|---|---|---|---|
| Regional partitioning (baseline) | Separate DB clusters and object stores per region; requests route to the local region only | Availability depends entirely on in-region redundancy (multi-AZ within the region); no cross-border failover possible by design | Default choice: strongest compliance guarantee, simplest to audit, and the right starting point before considering exceptions |
| Consent-based replication | Per-record consent flag; consenting users' data replicates to a central region for global features (analytics, search) | Same as regional partitioning for non-consenting users; consenting users gain cross-region availability for the replicated subset | Global features (cross-region search, aggregate analytics) that need some data outside the home region, and where consent is a legally sufficient basis |
| Dual writes with residency filters | App writes to the local primary and, for allowed fields only, to a remote replica via residency-aware middleware that filters or pseudonymizes before crossing the border | Higher availability for the allowed subset (near real-time cross-border access) but adds write-path complexity and a second consistency boundary to reason about | Narrow fields that are provably safe to replicate (pseudonymized identifiers, non-PII aggregates) where the latency or availability win justifies the added complexity |
| Contractual and technical mitigations | Data processing agreements, standard contractual clauses, and customer-managed encryption keys stored in-country, used instead of full architectural separation | Doesn't change the underlying availability profile; it's a legal and access-control layer on top of whichever architecture you chose | When full technical separation is genuinely impractical and legal has validated the mitigation is sufficient for the jurisdiction in question; not a substitute for architecture, a supplement to it |
Achieving reasonable availability inside the constraint
Once data can't leave a region, availability has to come from redundancy that stays inside the boundary: multi-AZ deployment within the region (the residency constraint is about crossing the country border, not about running a single instance), automated in-region failover with the same detection-and-promotion mechanics as any single-region HA design, and continuous backups stored within the same jurisdiction (a backup that lands in a different country to satisfy some other requirement, like disaster-recovery diversity, would itself violate residency unless that specific transfer is legally permitted). The trade-off to be explicit about with stakeholders: a residency-constrained region cannot achieve the same worst-case availability ceiling as a design free to fail over anywhere on Earth, because the blast radius of "this entire region's infrastructure provider has an outage" can't be absorbed by routing to a healthy region elsewhere. That's a real, unavoidable cost of the constraint, not an engineering gap to be closed.
Worked example: one customer's data, traced through the design
Take a specific case: a customer in Germany whose profile (name, address, government ID fields) is subject to EU residency rules. Under regional partitioning, that profile lives entirely in an eu-central data store, and the German customer's traffic is routed only to eu-central; a customer-service agent or batch job in us-east has no path to that record at all, by design, and if eu-central has an outage, availability for that customer depends only on eu-central's own in-region redundancy, not on any other region. Now add one narrow exception under dual writes with residency filters: the product wants a global leaderboard that ranks customers by loyalty tier across all regions. The app writes the full profile to eu-central as always, and separately, the residency-aware middleware extracts and pseudonymizes just one field, a loyalty-tier flag ("gold", "silver", "bronze", with no name, address, or ID attached), and replicates only that pseudonymized value to us-east so the global leaderboard can read it. Every other field (name, address, government ID) never crosses the border; only that one pseudonymized, non-identifying field does, which is exactly the kind of narrow, provably-safe field this option is meant for. If the filtering logic had a bug that let the raw name field slip through in that same replication path, that would be a real data-protection incident, not just a defect, which is why this option carries the most scrutiny of the four.
Trade-offs & pitfalls
The most common mistake is treating residency as a per-field flag applied after the fact rather than a first-class part of the data model from the start; retrofitting field-level residency onto a schema that was built assuming free replication usually means re-auditing every existing replication path and every downstream consumer (analytics jobs, search indexes, caches) that may have already copied restricted data outside the boundary. Consent-based replication is attractive because it enables global features, but it pushes real operational complexity onto consent-state management: a user who revokes consent after their data has already replicated needs a defined, auditable deletion path in every region it reached, not just a flag flip in the primary. Dual writes with residency filters are the riskiest option operationally, because the filtering logic itself becomes a compliance-critical piece of code; a bug that lets one unfiltered field slip through the residency-aware middleware is a real data-protection incident, not just a bug, so this option deserves the most scrutiny and testing of the four and should generally be adopted last, only for fields where the business case clearly justifies the added risk.
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 bulkhead pattern, and how does it stop one failing dependency or noisy tenant from taking down the whole system? Give a concrete example of where you'd draw the isolation boundary.
Sample Answer
Direct answer
The bulkhead pattern partitions a system's resources (thread pools, connection pools, CPU, or entire nodes) into isolated compartments, named after a ship's watertight bulkheads, so that one failing dependency or one noisy tenant can only exhaust the resources in its own compartment, not the resources every other caller depends on. Without bulkheads, a single slow or misbehaving dependency can consume every available thread or connection in a shared pool, and a completely healthy code path fails simply because it couldn't get a thread to run on.
Where to draw the isolation boundary
A concrete example: an API gateway calls three downstream services, an inventory service, a recommendations service, and a payments service, all through one shared thread pool. If recommendations starts responding slowly, every thread in the shared pool eventually ends up blocked waiting on recommendations calls, and inventory and payment requests start timing out too, even though nothing is wrong with either of them. The fix is a dedicated, bounded thread pool (or connection pool) per downstream dependency: recommendations gets its own pool of, say, 10 threads, so a recommendations outage can stall at most those 10 threads and its own queue, while inventory and payments keep running normally on their own separate pools.
The boundary should sit wherever one caller's failure or slowness shouldn't be able to spill onto another caller's request. Common places to draw it:
- Per-downstream-dependency, as in the example above: each external service or database gets its own pool so a slow one can't starve calls to a fast one.
- Per-tenant, in a multi-tenant system: each tenant (or tenant tier) gets a capped share of connections or CPU so one noisy or abusive tenant can't degrade service for everyone else on shared infrastructure.
- Per-criticality-tier: payment and auth paths get reserved capacity separate from lower-priority paths like analytics or notifications, so a spike in low-priority traffic can't crowd out the paths that actually matter.
Trade-offs & pitfalls
Bulkheads trade utilization for isolation: reserved capacity that a compartment isn't currently using sits idle rather than being available to a busier compartment, so a poorly sized bulkhead can cause localized throttling even while the system as a whole has spare capacity. Sizing is the actual hard part in practice, not the pattern itself: too small and a legitimate burst of normal traffic gets rejected by its own bulkhead; too large and the isolation becomes theoretical, because if every pool is sized close to the shared pool's original total, a single compartment can still consume enough of the machine's real resources (CPU, memory, file descriptors) to degrade its neighbors even though the pool counters look fine. Bulkheads are also a different tool from a circuit breaker and the two are frequently confused: a bulkhead limits how much of a shared resource one dependency can consume (a capacity boundary), while a circuit breaker stops sending requests to a dependency once it's clearly failing (a decision to stop calling at all); they're complementary, since the bulkhead caps the damage while the circuit breaker is deciding whether to keep trying, and production systems typically use both on the same dependency together. The same reasoning extends beyond web request threads: an ML-serving platform running GPU inference for multiple models on shared hardware applies the identical idea by pinning each model (or tenant) to a dedicated slice of GPU memory and compute, so one model that starts issuing runaway-batch-size requests can't starve GPU capacity away from every other model sharing that hardware.
You're deciding between active-active and active-passive for a service that needs 99.999% availability. Walk through the cost and operational trade-offs of each, and which failure modes each one actually protects against.
Sample Answer
Direct answer
At 99.999% availability the annual downtime budget is under 5.3 minutes, tight enough that the choice between active-active and active-passive comes down to arithmetic, not preference. Active-passive's failover time is paid per incident, so with more than a couple of incidents a year it tends to blow the budget on its own, while active-active's failure handling (rerouting at a load balancer, not promoting a new primary) costs seconds per incident and leaves headroom. The trade is that active-active buys that headroom by taking on multi-master data consistency complexity that active-passive avoids entirely.
Structured elaboration
| Dimension | Active-active | Active-passive |
|---|---|---|
| Steady-state cost | Roughly 2x baseline compute, both regions serve production traffic | Roughly 1.2x to 1.5x baseline, standby is sized for failover, not full production |
| Data consistency | Requires multi-master conflict resolution (CRDTs, vector clocks, or app-level merge) or a globally consistent database | Single writer region, standby is a replication target with no write conflicts to resolve |
| What it protects against | Single-node, zone, and whole-region compute failure with near-zero client-visible interruption | Same failure classes, but recovery is bounded by failover execution time, not instantaneous |
| What it does NOT protect against | A bad deploy or corrupted write replicated to all active regions simultaneously | A stale or never-exercised standby that fails when finally promoted |
| Operational overhead | Continuous testing of split-brain and conflict-resolution paths | Regular, realistic failover drills to catch standby rot |
Worked example
99.999% availability budget:
525,600×0.00001=5.256 minutes per yearTake a payment-authorization service with a realistic incident cadence: assume, as a stated planning assumption from historical infrastructure incident rates, 4 unplanned failures per year serious enough to require failover, a mix of zone and instance-level failures.
Active-passive scenario. Automated failover, no human step, completes in 90 seconds per incident:
4×90=360 seconds=6.0 minutes per yearThat is over budget, 6.0 versus 5.256 minutes, even with fast, fully automated failover. To fit the budget with the same 4 incidents a year, each failover would need to complete in:
45.256×60=78.8 secondsThat is an extremely tight bound for anything involving promoting a standby and re-pointing traffic, and it leaves zero room for a fifth incident.
Active-active scenario, same service. A load-balancer health check reroutes traffic away from a failed region in about 10 seconds per incident, no promotion, no data cutover, just routing away from the unhealthy target:
4×10=40 seconds=0.67 minutes per yearThat leaves roughly 4.6 minutes of budget headroom for a larger, rarer event.
The same arithmetic applies to a 50-million-user authentication service. Authentication is read-heavy and largely stateless per request (token validation), making it one of the easier services to run active-active, since there is little write-conflict surface to design around. Most of the active-active complexity budget there goes toward the credential and session write path (password changes, new logins), not the read-heavy validation path that dominates traffic.
Trade-offs & pitfalls
- The math above assumes failover time is the only downtime source. A bad deploy that reaches both active regions simultaneously is a failure mode active-active does not protect against, and it is exactly the failure mode a canary-then-passive-region rollout is designed to catch; active-passive's "stale" region is sometimes a feature during a rollout, not just a liability.
- Active-passive standby rot is the single most common cause of a failed real failover: a standby that has never been promoted in a drill is not a tested capability, it is a hope.
- Active-active's data consistency cost is frequently underestimated. Teams budget for infrastructure duplication but not for the engineering time to design idempotent, conflict-resolvable writes, which is often the larger cost.
- Do not assume automated failover alone makes active-passive equivalent to active-active for a five-nines target; the arithmetic above shows per-incident RTO dominates the budget once incident count exceeds one or two a year.
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.