Learn/System Design/Multi-region Design & Conflict Resolution
Advanced~18 min read

Multi-region Design & Conflict Resolution

Choose regional ownership, consistency contracts, conflict repair, safe promotion and failback.

System DesignDistributed SystemsProduction

Regions solve different requirements

A second region may reduce reader latency, provide disaster recovery or satisfy data residency. These objectives need different architectures. Replicating a database abroad can violate the intended locality even when all application requests stay local.

List which failures must be survived: process, host, zone, region or shared control-plane failure. Multiple regions using one identity service, deployment credential or traffic manager still share dependencies. Specify how much traffic the surviving region must handle and what functionality can degrade.

Compare three ownership models

ModelNormal writesMain benefitFailure decision
Primary plus standbyOne regionClear authoritative orderFence old primary and promote an eligible replica
Regional owner per tenantEach tenant has one writer regionLocality without same-record concurrent writersTransfer tenant ownership and routing generation
Multi-writerSeveral regions modify shared recordsLocal writes during some partitionsDefine conflicts or coordinate invariants

A multi-region frontend can still use a single writer database. An active-active label does not tell you whether the same record can be written in both places. Make that distinction explicit.

Latency and acknowledgment contract

Suppose a local database commit takes 8 ms and a distant region round trip is 90 ms. Waiting for that region changes the write latency budget; adding local application replicas cannot remove the WAN dependency. These are scenario inputs, not provider performance claims.

Specify whether acknowledgment means received in memory, logged durably or applied for reads. Promotion eligibility depends on that contract. Synchronous replication protects the acknowledged writes only under the chosen durability and failover policy; promoting a different stale replica or losing all durable copies defeats the assumption.

Sessions and read-your-writes

A user edits a profile in region A, then is routed to B. An asynchronous replica may show the old profile. Options include returning the changed state from the write, temporarily routing relevant reads to A, or attaching a commit position and waiting until B reaches it.

Bound that wait and define the timeout response. A short sticky-routing interval is a heuristic if replication can take longer; it is not a proof. Monotonic reads require that the next read not move to a version older than the one already observed. Authorization changes and inventory decisions can demand stricter contracts than profile display.

Partition timeline and split brain

text
A is owner at epoch 12
Network separates A from the promotion authority
B receives ownership at epoch 13
A still has processes and queued requests
Protected writes must reject epoch 12

Traffic routing alone is not fencing: old connections and delayed jobs can still reach A. Use infrastructure that prevents the old writer from committing, or a storage protocol that enforces ownership epochs. If safe ownership cannot be established, pausing writes can be the correct availability trade-off.

Do not promote just because one health check timed out. Separate evidence of application unhealthiness, replication freshness and authority to promote. Automation reduces response time but can also execute an unsafe decision faster.

Conflicts must preserve business meaning

DataCandidate mergeWhat can go wrong
Independent set additionsSet unionRemoval semantics need a different design
Per-writer positive counterMerge each component by maximum, then sumConcurrent reset/decrement is not covered
Profile fieldsField-specific version/merge policyLast-write-wins can erase a valid edit
Last available seatSerialized conditional reservationMerging two accepted sales cannot manufacture another seat

A grow-only counter illustrates a CRDT: A's component changes from 4 to 5, B's from 7 to 9; component-wise maximum gives 5+9=14 after convergence. Updates must belong to unique writer components and not decrement. This is useful for compatible counters, not a universal transaction replacement.

Wall-clock last-write-wins depends on clock assumptions and discards losing values. Preserve conflicted versions where human review matters. Vector clocks can expose concurrent histories, but deciding which result the business wants remains application work.

Repair and anti-entropy

Repair compares replica state and reconciles differences even when no user reads the affected key. Read repair only fixes keys that are accessed. Merkle-tree comparisons can narrow differing key ranges by hierarchical hashes, reducing comparison transfer when most data agrees.

Keep tombstones, schema versions and ownership metadata in the repair contract. Monitor time since successful repair and backlog, not only current replication lag. A disconnected replica returning after deletion must not resurrect its stale value.

Failback is a migration

After B becomes primary, restarting A does not make A authoritative again. Rebuild or catch up A from the accepted history, compare state and transfer ownership using the same guarded process. Reconcile external effects, such as charges completed while a region was isolated, before re-enabling affected workflows.

Test failover with reduced capacity, cold caches and a surge of reconnects. Reserve headroom or explicitly shed lower-priority traffic. Restoration, verification and traffic movement contribute to recovery time; the database election alone does not measure business recovery.

Worked regional capacity decision

Two regions normally serve 6,000 requests/second each. Each can sustainably serve 8,000 under the latency SLO. Losing one leaves 12,000 requests/second for an 8,000-capacity survivor. The plan must add capacity in time, shed 4,000 requests/second, or change the objective. Replication alone has not achieved full-load regional resilience.

Exercise

Design a global job board with regional user data and a public catalog. Specify which data can be replicated globally, where a user writes, how progress reads observe writes and what happens during a region loss. Draw a stale-client write and a failback timeline. Explain why public catalog availability and account mutation safety can use different policies.

Section navigation

Keep it in your account.

Your progress is saved securely and available when you return.

Sign in Create a free account