Multi-Region and Geo-Distributed Systems Questions
Running a system across regions and continents: multi-region replication, data residency and sovereignty, geo-routing and CDN edge distribution, cross-region consistency and quorum placement, and conflict resolution when two regions accept writes. Covers regional failover and split-brain prevention, recovery objectives (RTO/RPO), region-by-region rollout and blast-radius containment, and the latency, cost, and consistency tradeoffs of going global. Global distribution strategy across the service and data tiers.
Explain Recovery Time Objective (RTO) and Recovery Point Objective (RPO). How do synchronous replication, asynchronous replication, and regular backups influence the RTO and RPO you can guarantee across regions? Provide concrete example pairings of replication choices to achievable RTO/RPO targets.
Sample Answer
RTO (recovery time objective: the maximum acceptable downtime after a failure before the outage itself becomes the real problem) answers "how long can we be down," and RPO (recovery point objective: the maximum acceptable data loss, measured in time, meaning how old the recovered data is allowed to be) answers "how much data can we afford to lose." Replication choice sets the RPO you can promise, and how automated the cutover (the moment
traffic actually switches to the newly promoted region) is sets the RTO; a backup-only strategy caps both at whatever the backup cadence and restore time allow.
How the mechanism shapes RPO
- Synchronous replication (a write is only acknowledged once every region in the write quorum has durably stored it) drives RPO to effectively zero, because no acknowledged write exists on only one region.
- Asynchronous replication (the primary acknowledges a write before a remote replica has it) leaves an RPO equal to whatever unreplicated data existed at the moment of failure, typically the size of the replication lag window.
- Periodic backups or snapshots set RPO equal to the backup interval: a nightly backup means up to 24 hours of data loss in the worst case.
How the mechanism shapes RTO
RTO depends less on the replication mode itself and more on how automated the failover is: synchronous replication makes an immediate promotion safe (the standby already has everything), though cross-region synchronous writes add latency to every write and can slow down under network stress. Asynchronous replicas can fail over just as fast mechanically, but may need a reconciliation step if any writes were in flight. Restoring from a backup is inherently the slowest path, since it involves finding the backup, restoring it, and replaying anything since.
Concrete pairings
| Replication choice | Typical RPO | Typical RTO | Where it fits |
|---|---|---|---|
| Synchronous cross-region, automated failover | Near zero | Seconds to minutes | Payments ledgers, anything where losing an acknowledged write is unacceptable |
| Synchronous in-region plus asynchronous cross-region | Seconds to low minutes | Low minutes | Most transactional services that need durability but can't pay cross-region write latency on every commit |
| Asynchronous replication only | Minutes, the lag window | Minutes plus reconciliation time | User profile or session data where a brief rollback is tolerable |
| Nightly backups, no live replica | Up to 24 hours | Hours, restore plus replay time | Analytics, batch reporting, anything reprocessable from source |
Worked example
If a database replicates asynchronously with a steady-state lag of 8 seconds and the primary region fails at an arbitrary moment, two different numbers come out of the same data and it matters which one you are quoting.
Replication lag is a STOCK, not a rate: at any instant it is the age of the unreplicated tail, measured in seconds of writes. Because it is a stock sampled at a random instant, the EXPECTED loss is the time-weighted average of the lag. Take a day that sits at 8 seconds of lag for 99% of the time and spikes to 90 seconds for the other 1%: the expected loss is 0.99 x 8 + 0.01 x 90 = 8.8 seconds. That is the number that describes a typical bad day, and it is genuinely the average lag over the day, not an alternative to it.
The number you can COMMIT to is a different statistic entirely. An objective is a bound, not an expectation, and a failure does not politely arrive at an average moment. If the failure lands inside one of those spikes, the loss is 90 seconds, and the spikes are not independent of failures: the same network or load event that spikes replication lag is often what takes the region down, so conditioning on a failure actually shifts the sample toward the tail rather than away from it. State the RPO against the worst observed lag (or a stated high percentile, with the percentile named), and state the expected loss separately if you want a planning number. Quoting the average as the guarantee is how an 8-second promise becomes a 90-second incident with a straight face.
Trade-offs and pitfalls
Tightening RPO toward zero by moving to synchronous cross-region replication is not free: it adds a wide-area round trip (commonly tens to a couple hundred milliseconds between continents) to every write's latency and can reduce availability, because a write now blocks if the remote region is unreachable. The common mistake is quoting an RPO number measured under normal conditions and never re-verified during an actual regional outage drill, when replication lag is exactly the thing most likely to have spiked.
Design a failure-injection (chaos engineering) framework to safely simulate a regional outage for production traffic. Describe safety gates, blast radius control, automated rollback, metrics and SLO checks to assert resilience, coordination with stakeholders, and runbook steps executed during and after the experiment.
Sample Answer
A safe chaos-engineering framework for a regional outage is a controlled experiment with three non-negotiable properties: the blast radius grows in small, reversible steps; every step is guarded by an automated gate tied to real SLOs (service-level objectives: the internal targets defining acceptable service behavior); and a rehearsed rollback executes in seconds, not minutes, the moment a gate trips.
Safety gates and preconditions
Before starting: confirm synthetic canary traffic is healthy in every region, confirm no risky deploy has landed in the last 72 hours, and require explicit sign-off from both an SRE lead and the owning product team, recorded against the experiment so there is a clear record of who approved running a fault against production.
Blast radius control
Start at the smallest meaningful scope: a single, non-critical service, a single region, 1 percent of that region's traffic, during a low-traffic window. Only widen scope (more services, more traffic percentage, a full region) after the previous scope has run cleanly, and always exclude anything stateful and hard to reverse, such as a primary database, from the earliest stages.
What fault is actually being injected: draining is not an outage
Be explicit that shifting load-balancer weight away from a region is an EVACUATION, not an outage. Conflating the two is the most common way a regional-outage drill passes cleanly while the real event goes badly. A graceful drain announces itself in advance: connections finish, health checks stay green throughout, in-flight writes complete, replication keeps streaming, and no leader election is ever triggered. What a drain does test is real and worth testing on its own, namely whether the surviving regions have the compute headroom, connection-pool capacity, and cache warmth to carry the extra load. What it structurally cannot test is everything that distinguishes an outage: how long detection takes, what happens to requests in flight at the instant the region vanishes, whether writes acknowledged a second before the loss survive, whether leader election completes and in what time, whether clients holding long-lived connections reconnect, and whether OTHER regions that call into the now-dead region degrade gracefully or cascade.
So run them as a ladder, not as one experiment. First the evacuation, at the percentages above, to prove capacity, because that is what tells you the surviving regions will not fall over when the second stage works. Then an unannounced HARD fault under the same blast-radius discipline: blackhole the region's traffic at the network layer, or fail its health-check endpoint without draining first, so that detection time sits inside the measurement rather than outside it. Start the hard fault at the same smallest scope (one non-critical service, one region, a low-traffic window) and widen only once measured detection and failover times match what the runbook claims they are.
Automated rollback
An orchestrator (infrastructure-as-code plus an automation layer, not a person manually editing routing tables mid-incident) performs the actual fault injection, for example shifting load-balancer weight away from a region, and continuously evaluates health checks. The moment a threshold breaches, it reverts the routing change automatically and only then pages a human, so recovery does not wait on someone noticing a dashboard. For the hard-fault stage the undo is a different action and has to be built and rehearsed separately: withdrawing a network blackhole, restoring a health-check endpoint, or re-advertising a route is not the same operation as moving a load-balancer weight, and an abort path that only knows how to move weights will fail exactly when it is needed.
Metrics and SLO checks
Watch user-facing latency (p95 and p99: the 95th and 99th percentile response times, meaning 95 percent or 99 percent of requests are faster than that number), error rate, and a downstream signal like queue depth or database write error rate; define a concrete gate, for example "error rate above 0.5 percent sustained for 2 minutes" or "p99 latency more than double the pre-experiment baseline," and treat either as an automatic abort, not a judgment call made under pressure.
Coordination with stakeholders
Notify the on-call team, the affected product owner, and customer support before starting, with the exact scope and abort criteria in writing; post a live status update during the run so nobody discovers the experiment by noticing a spike in their own dashboard first, and circulate a written summary afterward regardless of outcome.
Runbook: during and after
During: run the pre-flight checklist that is part of the runbook (the written, step-by-step procedure the team follows to execute and respond to the experiment), get sign-off, start at the smallest scope, monitor continuously, and either proceed to the next scope step or auto-abort based on the gate. After a clean run: let the system soak at full scope for a further period, commonly 24 hours, before declaring success, and publish the findings. After an aborted run: open an incident regardless of whether real users were affected, preserve the metrics and traces from the window, and drive concrete fixes (a missing retry, an under-provisioned failover path, a replica quorum misconfiguration, where too few replicas were required to agree before a write counted as committed) before attempting the experiment again.
Worked example
An outage simulation against a secondary region starts at 1 percent of that region's traffic for 5 minutes; health checks stay within the pre-defined error-rate gate (under 0.5 percent), so the orchestrator ramps to 5 percent for the next 15 minutes, then to 25 percent for 30 minutes. At the 25 percent step, error rate crosses 0.6 percent sustained for just over 2 minutes; the automated gate trips, the orchestrator reverts the routing change within seconds, and an incident is opened even though no real customer-visible outage resulted, because the gate existing to trip is the point of the exercise.
Trade-offs and pitfalls
Conservative, small-step phasing is slower to reach a meaningful conclusion but is what makes running this safely against real production traffic possible at all; skipping straight to a full-region simulation to save time is how a chaos experiment turns into a real incident. The recurring mistake is defining gates only on the service being tested and missing a downstream dependency that silently absorbs the failure differently, such as a queue that backs up quietly for twenty minutes before anything alerts; include downstream saturation metrics (signals that a resource, like a queue or a connection pool, is filling up faster than it drains) in the gate set from the start, not after the first surprise.
Your system uses asynchronous cross-region replication and you occasionally observe data loss during failover. The dev team demands near-zero data loss but global write latency must stay under 200ms for 90% of users. As an SRE, propose architecture(s) to minimize data loss while meeting latency constraints, including trade-offs (semi-sync, write shards, quorum writes), and rollback paths.
Sample Answer
Direct answer
The data loss happens because asynchronous replication acknowledges the client's write before it has landed in any other region, so anything still in flight when the primary fails over is gone. The fix is not "make everything synchronous" (that would blow the 200ms budget), it is to make the minimum durability guarantee synchronous: acknowledge a write to at least one nearby peer region before telling the client it succeeded, and reserve full quorum writes for the specific data that most needs zero loss.
Proposed architecture
- Semi-synchronous replication to a nearby peer. The primary waits for an acknowledgment from one geographically close region (tens of milliseconds away) before acking the client, instead of waiting for every region or none at all. This shrinks the "in-flight and unreplicated" window from "everything written since the last async batch" to "the single write currently in flight to one peer."
- Write sharding by home region. Route each write to the region that owns that shard (by user or tenant), and make its synchronous peer a nearby region rather than a global one. This keeps the added latency bounded by a short regional hop instead of a transoceanic one.
- Quorum writes for the highest-value data. Define a replica set of N = 3 (the primary region plus its two nearest peers) with a write quorum W = 2: a write is durable once two of three regions have it. Combine with a read quorum R = 2 so reads survive losing any one region without a global round trip.
- Rollback and reconciliation path. On failover, promote whichever region has the most advanced applied log position, not just "the next one in a list." Anything that was in flight and never reached a second region is captured by writing every mutation first to a durable, synchronously-written write-ahead log (or an equivalent message queue) so it can be replayed against the new primary once it recovers, using an idempotency key so replay never double-applies a write.
Worked example
Rough network round-trip times (RTT): us-east <-> us-west about 60ms, us-east <-> eu-west
about 90ms, us-east <-> ap-southeast about 200ms. If the synchronous peer is the nearest one (60-90ms round trip), the write path has to wait for
a full round trip to that peer, not half of it (half would be about 30-45ms, an underestimate),
so it adds roughly 60-90ms waiting for the peer's ack: the primary has to send the write across
and get the acknowledgment back before it can tell the client the write is durable enough. Add typical application and serialization overhead of 50-80ms and the total lands around
110-170ms (a half-round-trip miscalculation charges only 30-45ms for the network and lands at 80-125ms, understating the total by exactly the half round trip it dropped, 30ms at the near end and 45ms at the far end, and making the budget look far roomier than it is), comfortably inside
the 200ms p90 (90th-percentile) budget. Trying to make the synchronous peer a third continent away (200ms round trip) adds that entire
200ms on its own, not half of it (100ms alone would still leave headroom; the real full 200ms
does not), and would blow the budget before any serialization overhead is even added, which is exactly why that far region stays asynchronous-only (a disaster-recovery copy, not a consistency participant).
Trade-offs and pitfalls
- This design does not reach zero data loss, it shrinks the loss window from "anything since the last async flush" to "the single write in flight to the second node," which is the honest trade the 200ms budget forces.
- Promotion logic must fence the old primary (revoke its ability to accept writes, for example by expiring its lease) before the new one starts accepting writes, or you risk split-brain: two regions each believing they are the write leader and both accepting conflicting writes.
- Sharding by home region pushes complexity onto any operation that spans two users' shards; those need their own coordination story (a saga or a narrower synchronous scope), not silent fan-out.
You're evaluating multi-region database topologies for an e-commerce platform: a strongly-consistent distributed SQL (e.g., CockroachDB), managed SQL with read replicas per region, or per-region DBs with asynchronous replication. Compare each option for consistency, latency, operational complexity, failover, and migration difficulty. Recommend an approach for target local latency 10ms and cross-region acceptable latency 2s.
Sample Answer
Direct answer
Given a tight 10-millisecond local latency target and a generous 2-second cross-region target, a distributed SQL database (in the style of CockroachDB, offering strong consistency via a consensus protocol per range
of data, where a "range" is simply a contiguous chunk of a table that the database replicates
and coordinates writes for as one unit) configured with geo-partitioned tables and pinned replica placement is the strongest fit: it can hit the 10-millisecond local target for reads and writes served by a nearby, correctly-placed replica, while still comfortably fitting genuinely cross-region operations inside the lenient 2-second budget, all without giving up strong consistency.
Structured elaboration
| Option | Consistency | Latency | Operational complexity | Failover | Migration difficulty |
|---|---|---|---|---|---|
| Distributed SQL (e.g. CockroachDB-style) | Strong (linearizable/serializable) globally, via per-range consensus | Sub-10ms locally only if replicas and leaseholders (the one replica currently responsible for | |||
| coordinating reads and writes for a given range) are explicitly pinned near the requester; | |||||
| naive placement can silently blow the 10ms target | High: requires understanding and actively configuring replica and leaseholder placement, zone constraints | Automatic at the range level via consensus, no manual promotion step needed | High: application and schema often need real redesign around partitioning keys | ||
| Managed SQL with regional read replicas | Strong on the single writer, eventually consistent (lagging) on replicas | Local reads fast; writes from a remote region pay a full round trip to the single writer, comfortably inside 2s | Moderate: standard single-writer operational model plus lag monitoring | Manual promotion of a replica, minutes-scale | Low: closest to a familiar single-region relational setup |
| Per-region databases with independent async replication | Weakest: eventual, with real risk of exceeding the 2s cross-region acceptable window under load | Fastest local latency of the three, no consensus overhead | High ongoing complexity to reconcile divergent regional writes | Region-local, but reconciling with peers after a failure is nontrivial | Highest: essentially building and maintaining a bespoke multi-master system |
Worked example
A same-region, cross-availability-zone (AZ) network round trip is typically on the order of 1-2 milliseconds. A range's three replicas pinned within one region across different AZs reach a local consensus quorum (a quorum is simply the minimum number of replicas that must agree before a write counts as committed) in about one of those round trips, not inside one. A majority of three is two, and the leaseholder counts as one of the two itself, so the commit cannot return until at least one follower's acknowledgment has travelled out and come back: a full AZ round trip is the floor, not a budget the commit fits underneath. The realistic figure is that round trip plus the write-ahead-log flush on the leaseholder and on the acknowledging follower, so roughly 2-4 milliseconds end to end. That floor is worth stating out loud in an interview, because any claimed commit latency below the round trip to the replica whose acknowledgment is being waited on is physically impossible, and quoting such a number is a reliable sign the figure came from somewhere other than the network path. Even at 2-4 milliseconds there is still comfortable headroom inside the 10-millisecond local target for query planning, application overhead, and the occasional retry.
A cross-region operation, say between a U.S. East and a European region at a typical wide-area round trip in the 70-90 millisecond range, uses roughly:
2000ms85ms≈4% of the 2-second cross-region budgetleaving roughly a 23x safety margin (2000ms / 85ms is about 23.5x, comfortably more than the
roughly 20x a rounder mental estimate might suggest). That headroom means occasional genuine cross-region transactions (a multi-region query, or a lease transfer during rebalancing) are entirely affordable, as long as they stay off the primary user-facing hot path that has to hit the 10-millisecond local number.
Trade-offs and pitfalls
The single biggest risk with a distributed SQL deployment here is trusting the database's default replica-placement algorithm, which typically optimizes for even distribution and fault tolerance across the whole cluster, not for keeping a given user's data physically close to where they actually connect from. Left on defaults, a request can end up needing a wide-area round trip to reach the current leaseholder for its data, silently failing the 10-millisecond target even though the database is technically "strongly consistent" and "working." Hitting the stated targets requires deliberately configuring geo-partitioned tables and zone/locality constraints so each user's data has a replica, and ideally the leaseholder, physically near them.
Design a globally-distributed architecture for a food-delivery-style marketplace platform that supports low-latency order placement and local dispatch in 200+ cities, provides high availability, and enforces strong consistency for financial operations while allowing eventual consistency for non-critical data. Discuss partitioning, cross-region replication patterns, failure modes, data sovereignty, and the trade-offs of active-active versus active-passive regional topologies.
Sample Answer
Direct answer
Go active-active for the read and dispatch path (browsing, order placement, courier matching)
because that traffic is naturally single-city and latency-critical, but keep exactly one region
of record per financial transaction so money movements are never applied twice from two regions
independently. Most of this platform doesn't need cross-region availability at all; the design
problem is knowing which slice does.
Structured elaboration
Partitioning. Partition primarily by city/market, not by user id: almost every read (nearby
restaurants, available couriers) and write (place order, assign courier) is scoped to one city.
Route each city to its physically nearest region and keep that city's live dispatch state
entirely inside that one region.
Cross-region replication patterns. Live dispatch state stays regional; only a read-only copy
replicates asynchronously to other regions for global search and disaster recovery, never to serve
that city's own live traffic. The financial ledger instead replicates via change-data-capture
(CDC, streaming every committed insert/update straight off the database's write log) into a single
global ledger store, so the money view is reconciled once centrally rather than merged from
several regional copies.
Failure modes. A region outage during a city's peak hours stalls courier assignment for that
city; mitigate with a standby region per city that can take over dispatch within a stated recovery time
objective (RTO, how long the takeover itself may take, measured from detection to the standby
serving real traffic), accepting in exchange a recovery point objective (RPO, the maximum
acceptable data loss measured in time) of a few seconds of dispatch state that had not replicated
yet. The two budgets are separate and they trade against each other: synchronous replication drives
RPO toward zero but puts a cross-region round trip on every dispatch write, while asynchronous
replication keeps the write local and pays for it with a non-zero RPO. Quote both numbers per city,
because a 30-second RTO with a 10-minute RPO and a 10-minute RTO with a 1-second RPO are very
different products for a courier who is mid-delivery. Local connectivity flakiness (a courier's app losing signal) is handled at the edge
with a short-TTL cache of "nearby restaurants," not a cross-region call per request.
Data sovereignty. Markets under strict data-localization law must store personal identifiable
information (PII: customer addresses, courier identity documents) in-region only; only aggregated,
tokenized analytics (each real value swapped for a random reference id that cannot be
reverse-mapped without a separate, tightly controlled lookup table) may leave. That means "home region" has to be a hard placement constraint per
market in the architecture, not just a latency preference.
Active-active vs. active-passive.
| Active-active | Active-passive | |
|---|---|---|
| Best for | Dispatch/order placement (already single-city) | Financial ledger, refunds |
| Consistency risk | Two regions can apply conflicting writes to the same record | One writer at a time, no conflict window |
| Cost | Higher (duplicate live capacity everywhere) | Lower, standby is warm not fully serving |
Worked example
In Tokyo's dinner rush, ~3,000 active couriers ping their location every 5 seconds, about
600 pings/sec. Assignment always happens in Tokyo's home region regardless of topology, so making
dispatch active-active across ap-northeast and a US region would double that write volume for zero
latency benefit. Contrast the ledger: two regions independently applying a refund to the same order
before either sees the other's write can double-refund the customer. A single region-of-record
with synchronous replication to a nearby standby (say Tokyo to Osaka, on the order of 5-15ms round
trip) avoids that at a cost that's affordable precisely because ledger writes are a small fraction
of the 600+ pings/sec of dispatch traffic.
Trade-offs & pitfalls
Defaulting the whole platform to active-active "because it's more available" ignores that most
traffic is already naturally single-region; the real availability risk is that city's home region
going down, which active-passive with a tight RTO and a seconds-scale RPO handles at far lower
replication cost. The
pitfall is conflating "distribute for latency" with "distribute for consistency," which is exactly
how the double-refund bug above gets shipped.
Unlock Full Question Bank
Get access to all Multi-Region and Geo-Distributed Systems interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.