KnowSys

Running in Many Regions

Follow one shared shopping list, edited by a sister in Mumbai and a brother in New Jersey, as its app grows from one cloud region to several: why one region fails you, how a standby copy takes over and what it loses, and what happens when two regions accept writes at once, from home regions and cross-region consensus to last-writer-wins, vector clocks and CRDTs, as Redis, DynamoDB, Spanner and CockroachDB do it.

⏱ 65 min read◆ AdvancedAssumes: chapter 28 (replication) and chapter 27 (consensus); chapter 26 (clocks), chapter 30 (failure detection) and chapter 35 (DNS and anycast) help
Start reading

Asha lives in Mumbai and keeps a shared shopping list for her parents' flat in an app called Tokri. Her brother Kabir, who lives in New Jersey, is on the same list: when he hears on a video call that their mother is out of tea, he opens the app and adds it, and the next time Asha passes a shop, it's there. Tokri is a small company. Like a great many small companies, it runs everything in one place: its servers and its database sit in Amazon's us-east-1 region in northern Virginia, because that's where the tutorial they followed put them.

At about 12:20 in the afternoon in Mumbai on 20 October 2025, Asha opened the list and got a spinner, then an error. Tokri wasn't the problem. A few minutes earlier, a fault had hit the system that manages DNS records (the entries that turn a name into server addresses) for Amazon DynamoDB, the database Tokri keeps its lists in, and it had left the name dynamodb.us-east-1.amazonaws.com with no addresses behind it at all. Every program in that region that tried to reach DynamoDB failed. Amazon's summary of the event puts DynamoDB's errors at nearly three hours, and other services in the region that depended on it, such as launching new servers, stayed impaired for about fourteen. For those three hours Tokri was down for every customer it had, and there was nowhere else for it to go.

Even on a good day, Asha pays for the choice of Virginia. Every tap she makes travels to the far side of the planet and back, and a round trip that long takes roughly a fifth of a second before Tokri's servers have done any work at all. Kabir, a few hundred kilometres from the data centres, gets his answers in a few milliseconds.

So Tokri wants to run in more than one region: near Asha and near Kabir, and in a way that survives losing either one. That raises the question this chapter answers: how do you run one app in several regions so that losing a region doesn't take you down and every user is close to a copy, without losing or scrambling anyone's writes? We'll start with a standby copy that waits for disaster, find out what it loses and why switching to it is harder than it looks, then let both regions take writes and deal with what happens when Asha and Kabir edit the same list at the same moment, eight thousand miles apart.

01What one region can't do

1.1What a region is

A cloud provider's region is a cluster of data centres in one metropolitan area, run as one unit with its own copy of the provider's services. Amazon's us-east-1 is in northern Virginia, ap-south-1 is in Mumbai, eu-west-1 is in Ireland. Inside a region, the data centres are grouped into availability zones: separate buildings, or groups of buildings, with their own power and cooling, a few kilometres apart and joined by fast private fibre. A zone is designed to fail on its own, so an app spread across three zones survives a fire or a power cut in one.

Aerial view of many large flat-roofed data centre buildings and electrical substations among roads and housing near Ashburn, Virginia
Data centres near Ashburn in Loudoun County, Virginia, seen from a plane in November 2025. Northern Virginia holds one of the largest concentrations of data centres in the world, and us-east-1 is spread over many buildings like these. They're separate buildings with separate power, but they share software, operators and a control system, which is why a whole region can fail at once.Photo: Theodore Christopher, CC0, via Wikimedia Commons

Tokri already uses three zones, so a single burning building wouldn't take it down. What zones don't protect against is everything the buildings share. They run the same software, are changed by the same automated systems, are reached through the same DNS names, and depend on the same regional services. When one of those shared pieces breaks, every zone breaks with it.

One region: three zones, one set of shared services
US-EAST-1UsersEndpointssame DNS namesZone Aown powerZone Bown powerZone Cown powerShared servicessame software
Step 1. A fire or power cut takes out Zone A. Zones B and C have their own buildings and power, so Tokri keeps serving from them.
1 / 2

1.2Regions do fail

They do, and the published accounts show how. On 7 December 2021, an automated job that added capacity to one of AWS's internal services set off a flood of connection attempts on AWS's internal network in us-east-1. The devices connecting that internal network to AWS's main network were overwhelmed, and the clients kept retrying, which kept them overwhelmed. According to AWS's summary, the trouble started at about 7:30 AM Pacific time; the network devices didn't fully recover until 2:22 PM, and EventBridge was still working through its backlog of events until 6:40 PM. Launching new servers, the management console, and changes to DNS records were impaired for hours. AWS's own Service Health Dashboard, the page customers check to see whether AWS is having a problem, had trouble failing over to its standby region, and customers couldn't open support cases for nearly seven hours.

On 20 October 2025 it was DynamoDB, as in Tokri's story. AWS's summary traces it to "a latent race condition in the DynamoDB DNS management system": two copies of an automated process were applying DNS plans at about the same time, a delayed one overwrote a newer plan with a much older one, and a clean-up step then deleted that older plan, which left the regional endpoint with "an incorrect empty DNS record". The outage ran from 11:48 PM Pacific on 19 October to about 2:40 AM for DynamoDB itself, and until 2:20 PM for services built on it.

Two details in that second report matter for this chapter. First, customers who used DynamoDB global tables, which keep copies of a table in several regions, "were able to successfully connect to and issue requests" against their copies in other regions; those copies fell behind on changes from Virginia until it recovered. A second region helped the people who had one. Second, even some of AWS's own global services had pieces that lived only in us-east-1: console sign-in for some users in other regions failed, because it depended on something there.

1.3Two more reasons: distance and law

Outages are rare. Distance is permanent. Every one of Asha's requests crosses about 13,000 km of the Earth's surface each way, and section 2 shows that this costs about 200 ms per round trip however good the network is. A user in Virginia gets the same app in a few milliseconds.

The third reason is law. Some data has to stay in a particular country. India's central bank, the RBI, directed in April 2018 that operators of payment systems store "the entire data relating to payment systems operated by them" in a system only in India. Rules like that, called data residency requirements, mean some records must live in a given country's region whatever that does to latency or cost, and they turn "where does this row live?" into a design question.

So Tokri has three reasons to add a region: survive the loss of one, get close to users, and keep some data where the law says. Before deciding how, we need to know what the distance between regions costs, because every design in this chapter is a different way of paying it.

02How far away the second region is

2.1The speed of light in glass

Data between regions travels as pulses of light in optical fibre. Light in glass is slower than in a vacuum, about two thirds of its vacuum speed, or roughly 200,000 km per second, which is 200 km per millisecond. That's a physical limit no engineering can beat. A message from Mumbai to Virginia and back covers at least twice the great-circle distance between them (the shortest path over the Earth's surface), so its round-trip time (RTT), the time for a message to go and its reply to come back, can never be below that distance divided by 200 km per millisecond.

Real round trips are always longer, because cables don't follow great circles. They run along coasts, through the Red Sea and the Mediterranean, across the Atlantic, and through routers and repeaters along the way.

Microsoft publishes the median round-trip times between its cloud regions, measured continuously by probes and summarised monthly. The script below compares five of them, from the July 2026 table, with the floor the speed of light sets. Its regions are Central India (Pune, about 150 km from Mumbai), South India (Chennai), Southeast Asia (Singapore), West Europe (the Netherlands) and East US (Virginia). Its km function is the standard formula for distance over a sphere.

Compare real round-trip times between regions with the speed-of-light floor
python
Python
from math import radians, sin, cos, asin, sqrt
 
cities = {"Pune": (18.52, 73.86), "Chennai": (13.08, 80.27),
          "Singapore": (1.35, 103.82), "Amsterdam": (52.37, 4.90),
          "Virginia": (38.90, -77.45)}
 
def km(a, b):                        # great-circle distance on a 6,371 km sphere
    (la1, lo1), (la2, lo2) = [(radians(x), radians(y)) for x, y in (cities[a], cities[b])]
    h = sin((la2-la1)/2)**2 + cos(la1)*cos(la2)*sin((lo2-lo1)/2)**2
    return 2 * 6371 * asin(sqrt(h))
 
FIBRE = 200_000                      # km/s: light in glass, about 2/3 of c
real = {("Pune", "Chennai"): 20, ("Pune", "Singapore"): 53, ("Pune", "Amsterdam"): 139,
        ("Virginia", "Amsterdam"): 85, ("Pune", "Virginia"): 201}   # Azure median RTTs, ms
 
print(f"{'route':22s} {'distance':>9s} {'floor RTT':>10s} {'real RTT':>9s} {'ratio':>6s}")
for (a, b), r in real.items():
    d = km(a, b)
    floor = 2 * d / FIBRE * 1000
    print(f"{a+' - '+b:22s} {d:>6,.0f} km {floor:>7.0f} ms {r:>6d} ms {r/floor:>5.1f}x")
output
C++
route                   distance  floor RTT  real RTT  ratio
Pune - Chennai            914 km       9 ms     20 ms   2.2x
Pune - Singapore        3,784 km      38 ms     53 ms   1.4x
Pune - Amsterdam        6,966 km      70 ms    139 ms   2.0x
Virginia - Amsterdam    6,216 km      62 ms     85 ms   1.4x
Pune - Virginia        12,966 km     130 ms    201 ms   1.6x

Read the last column first. Real round trips are 1.4 to 2.2 times the physical floor, which is what you'd expect from cables that bend around continents. Pune and Virginia are nearly 13,000 km apart as the crow flies, so even a perfect straight fibre would take 130 ms there and back; the real median is about 200 ms. Pune to Chennai, two regions in the same country, costs roughly 20 ms.

Those two numbers set the scale for the rest of this chapter. A design that makes Asha wait for Virginia on every write costs her about 200 ms per wait. A design that makes Mumbai wait for a nearby region costs a tenth of that.

2.2What 200 ms does to an app

One round trip is only the start. A browser or phone opening a fresh secure connection first spends one round trip on the TCP handshake and another on the TLS handshake (chapter 35 walks through both), so Asha's first request to Virginia costs roughly three round trips, 600 ms, before Tokri's code runs. After that, every request on the warm connection costs one more round trip.

It gets worse when the app server and the database are in different regions. GitHub found this out in October 2018: after a failover moved its primary databases to the US West Coast while many of its application servers stayed on the East Coast, applications were, in GitHub's words, "unable to cope with the additional latency introduced by a cross-country round trip for the majority of their database calls". A page that makes 20 database queries one after another, each crossing 70 ms, spends 1.4 seconds waiting before it can respond.

2.3The links between regions fail too

Nor are the cables in the map permanent. Between 23 January and 4 February 2008, undersea cables were damaged in six places in the Mediterranean, the Persian Gulf and near Malaysia. The cuts near Alexandria, on two of the main cables between Europe and South Asia, disrupted internet traffic to India and the Middle East while operators rerouted traffic over other cables and ships went out to repair the breaks.

A relief map of the Middle East and South Asia with six numbered markers showing where undersea cables were damaged in January and February 2008, two near Alexandria in Egypt
The 2008 undersea cable cuts. Markers 2 and 3, near Alexandria, were on SEA-ME-WE-4 and FLAG, two of the main paths between Europe and India. The regions at both ends kept running; what broke was the path between them. A multi-region design has to decide what each region does when it can't hear from the others.Map: Zayani, on a Central Intelligence Agency base map, public domain, via Wikimedia Commons

When the path between two regions breaks while both keep running, each region can still serve its own users but can't reach the other. This is a network partition, a split of the system into groups that can't talk to each other. From inside one region, a partition looks exactly like the other region being dead, because in both cases the messages stop arriving (chapter 30 explains why silence can't be told apart from death). That ambiguity will come back to bite us in section 4.

So a second region is at least tens and often hundreds of milliseconds away, and the link to it sometimes breaks. Probably the simplest way to use one is to keep a copy of everything there and not send users to it until the first region dies.

03A standby copy: active-passive

3.1The idea

Tokri's first multi-region design keeps us-east-1 exactly as it is, serving every user, and adds a second region, say us-west-2 in Oregon, that serves nobody. Oregon holds a copy of the database that is kept up to date from the first, plus whatever servers it would need to take over. If Virginia fails, Tokri switches everyone to Oregon. This is active-passive: one region is active and handles all traffic; the other is passive and waits. Switching traffic and writes from the failed region to the standby is called failover.

Keeping the standby's database up to date is replication, which chapter 28 covers inside one region. The active database is the primary, the only copy that accepts writes. It sends a stream of its changes to the standby, the replica, which applies them in the same order. The question that matters across regions is when the primary tells the writer "done".

3.2Wait for the copy, or don't

Suppose Kabir adds tea to the list. Virginia's primary writes it locally. Now it has two choices.

With synchronous replication, the primary sends the change to Oregon and waits for Oregon to confirm it has stored it before telling Kabir "saved". If Virginia then vanishes, Oregon has the tea. But every write now costs a round trip to Oregon, about 70 ms in Microsoft's table for its own Virginia and California regions, and if Oregon is slow or unreachable, writes in Virginia stop, even though Virginia is fine. For a standby across an ocean, it's a few hundred milliseconds per write.

With asynchronous replication, the primary tells Kabir "saved" as soon as its own copy has the change, and sends it to Oregon in the background. Writes are as fast as with one region and Virginia keeps working when Oregon doesn't. The price is that Oregon is always a little behind. How far behind it is, in time or in writes, is called replication lag.

Nearly every active-passive setup across regions is asynchronous, because the synchronous version charges every user for an event that may never happen. Amazon's managed versions say how far behind they usually run: Aurora global databases replicate "with typical latency of under a second", and DynamoDB global tables, in their default mode, replicate "typically within a second or less".

3.3RPO and RTO

Asynchronous replication means that when Virginia dies, the last moments of writes may die with it. Two numbers describe how bad a failover is allowed to be, and every disaster-recovery plan is written in terms of them.

The recovery point objective (RPO) is how much recent data you can afford to lose, measured in time. An RPO of 5 seconds means "after a failover, we may lose at most the last 5 seconds of writes". With asynchronous replication, what you lose is whatever lag there was at the moment of failure.

The recovery time objective (RTO) is how long you can afford to be down, from the moment the region fails until users are working again in the other one. It includes noticing the failure, deciding to fail over, and doing it.

Your turn: design it before reading on

Tokri takes about 300 writes a second at its evening peak, and Oregon normally runs about one second behind Virginia. Virginia fails at the peak. How many writes are lost, and is that the worst case?

These two objectives are independent. A nightly backup copied to another region gives an RPO of up to a day and an RTO of however long it takes to rebuild everything from it. Synchronous replication to a fully running second region can give an RPO of zero and an RTO of minutes. What sits between them is mostly a question of how much of Oregon is already running when the failure happens.

3.4How warm to keep the standby

AWS's disaster-recovery guide names four strategies, from backup and restore up to multi-site active/active, and calls the active-passive form of the last one, a full-size copy that takes no traffic, a hot standby. The names have become the standard vocabulary. Here are the four levels a standby region can sit at.

StrategyWhat's running in the second regionWhat failover involvesTypical RPO / RTO
Backup and restoreNothing but copies of backupsBuild the whole environment from code, restore data from backupHours of data; hours to rebuild
Pilot lightDatabases, kept current by replication; servers defined but switched offStart servers, deploy anything missing, scale upSeconds of data; tens of minutes
Warm standbyA small but complete, working copy of productionScale it up to full sizeSeconds of data; minutes
Hot standbyA full-size copy, running, taking no trafficSwitch trafficSeconds of data; minutes or less

In the guide's own one-line distinction between the middle two, "pilot light cannot process requests without additional action taken first, whereas warm standby can handle traffic (at reduced capacity levels) immediately." The name pilot light comes from a gas boiler: a small flame that's always lit so the main burner can be turned on quickly.

?Why not always run a hot standby?

Because it doubles the bill for something that sits idle. And the cheaper levels hide a risk: scaling up at failover time depends on the cloud provider's control plane, the APIs that create, configure and resize resources, as opposed to the data plane that serves requests through resources that already exist. Control planes are exactly what tends to be impaired during a regional event; both AWS outages above made launching new servers fail for hours. And if the failure is region-wide, every other customer is trying to scale up in the same standby regions at the same moment. AWS's guide says to rely only on data-plane operations during failover for that reason, and calls a standby already sized for full load "statically stable": it needs nothing created to take over.

Decision

How warm should Tokri keep its standby region?

Pilot light
Replicated database in Oregon; app servers defined in code but not running.
  • Costs little beyond the database copy
  • Data loss limited to replication lag
  • Failover needs servers launched during an incident
  • Tens of minutes of downtime, more if launches are impaired
chosen
Warm standby
A small working copy in Oregon: a few servers, the replicated database.
  • Can take some traffic at once
  • Can be tested with real requests any day
  • Still has to scale up under load during the incident
  • Costs a few servers around the clock
Hot standby
Full-size copy in Oregon, running, idle.
  • Takes over with no new resources
  • Shortest recovery time
  • Pays for two full regions and uses one
  • At that cost, active-active is usually better value

A warm standby is the usual first step for a company Tokri's size: cheap enough to keep, and alive enough to test continuously by sending it a little real traffic. Its weak spot is the scale-up; sizing it to carry the minimum acceptable service (all reads and writes, maybe without image uploads or recommendations) without scaling removes most of that risk.

04Failing over

4.1The steps

Tokri now has a warm standby in Oregon, a second or so behind Virginia. When Virginia fails, four things have to happen, and each one can go wrong.

Active-passive, and a failover from Virginia to Oregon
ACTIVE: VIRGINIAPASSIVE: OREGONasync changesAsha, KabirphonesTraffic steeringDNS or anycastApp serversus-east-1Primary DBVirginiaApp serversus-west-2, fewReplica DBOregon, ~1 s behindHealth checksfrom several places
Step 1. Normal operation. Every request goes to Virginia. The primary database replicates each change to Oregon in the background, about a second behind.
1 / 6

4.2Noticing, and steering traffic

Detecting a dead region uses health checks: small requests sent every few seconds to a page on the app that answers only if the app and its database are working. Amazon's Route 53, for example, checks an endpoint every 30 seconds by default (10 seconds if you pay for fast checks) and marks it unhealthy after 3 failures in a row, from checkers in several regions. So detection alone takes between half a minute and a minute and a half, by design: chapter 30 explains why declaring something dead too fast does more harm than declaring it dead too slowly.

Steering traffic means changing where users' requests go. There are two common ways.

The first is DNS. Users find Tokri's servers by looking up a name such as api.tokri.app, and the answer is the address of a region. With DNS failover, the DNS service returns Virginia's address while Virginia's health checks pass and Oregon's address when they don't. The catch is caching. Each DNS answer comes with a TTL (time to live), the number of seconds resolvers and clients may reuse it, so after the switch some users keep going to Virginia until their cached answer expires. Failover records usually have short TTLs, a minute or less, and some clients still hold answers longer than they're told to.

The second is anycast, where the same IP address is announced from many places on the internet, and each user's packets go to whichever announcement is nearest by the internet's routing (chapter 35 covers how). Services such as AWS Global Accelerator and Cloudflare give an app one anycast address, receive traffic at the edge of their network near the user, and forward it to a healthy region. Since the address never changes, there's no DNS cache to wait out; switching regions is a change inside the provider's network.

?Why not fail over automatically as soon as the checks fail?

Because a failover isn't free, and a false alarm costs as much as a real one. Promoting Oregon throws away the writes in the lag window, and switching back later is a second risky operation. AWS's guide says automatically triggered failover "should be used with caution", notes that a failover on a false alarm incurs those losses for nothing, and says that manually triggered failover "is therefore often used", with every step after the decision automated so that the human only presses one button.

4.3Split-brain, and fencing

The most dangerous failure in this section is the one where Virginia isn't dead. Suppose the network between Virginia and the health checkers breaks, as in the 2008 cable cuts, but Virginia's servers keep running and users near Virginia can still reach them. The checkers declare Virginia dead, Oregon is promoted, and now two primaries accept writes for the same data, each believing it's the only one. This is split-brain: one system has split into two halves that each think they're in charge.

Asha's write to remove milk lands in Virginia; Kabir, routed to Oregon, adds bread. Each primary has writes the other lacks, and there's no simple way to merge two databases that both believed they were the only truth.

The defence is fencing: making sure the old primary can't accept writes before the new one starts. One way is to cut it off physically, by shutting its servers down or revoking its network access through the cloud provider's API; clusters call this STONITH, "shoot the other node in the head". Another is an epoch number: every promotion increments a number held in a store that both regions consult, every write carries the epoch of the primary that made it, and storage rejects writes from an older epoch. A deposed primary's writes then fail instead of quietly diverging. Chapter 26 covers the same idea as fencing tokens.

A failover with an epoch number
Old primary (Virginia)Epoch storeNew primary (Oregon)Storagewrite, epoch 7promote meepoch 8write, epoch 8write, epoch 7
Step 1. Before the failure, Virginia is primary with epoch 7, and storage accepts its writes.
1 / 5

The epoch store has to be something that can't itself split in two, which means a consensus system running across at least three regions (section 7 explains why three). That's a recurring pattern: even an active-passive design needs one small, strongly consistent piece spread over several regions to decide safely who's in charge.

4.4Failing back, and GitHub's 24 hours

GitHub's October 2018 incident shows what the RPO costs after a failover, even without split-brain. At 22:52 UTC on 21 October, maintenance on network equipment cut the link between GitHub's US East Coast data centre and its main network hub for 43 seconds. In those 43 seconds, GitHub's automated database failover tool, which ran across several sites and used the Raft consensus protocol (chapter 27) to make its decisions, saw the East Coast as unreachable and promoted database replicas on the West Coast to primary. Applications started writing to the West Coast.

When the link came back, the East Coast databases held a few seconds of writes that had never replicated west, and the West Coast held nearly 40 minutes of writes that the East didn't have. Neither copy was complete. The team couldn't fail back to the East Coast without throwing away the West's writes, and couldn't stay on the West Coast without the cross-country latency problem from section 2.2. They chose to protect the data, rebuilt the East from backups and replication, and reconciled the stranded writes afterwards. In all, the site was degraded for 24 hours and 11 minutes. GitHub's post-incident analysis gives the full timeline.

Predict before you read on

Tokri's failover from Virginia to Oregon took 4 minutes and lost about 1 second of writes. Virginia is healthy again an hour later. What's the safest way back?

4.5Why failovers fail when you need them

A failover path that's never used rots. Configuration drifts, the standby is missing a secret that was added in the main region last month, a quota in Oregon is too low to scale up, someone hard-coded the Virginia database address in one service, and the runbook refers to a dashboard that was renamed. None of this shows until the day it's needed, which is the worst day to discover it.

So teams practise. A game day is a planned exercise where a team deliberately fails something in production, or in a faithful copy of it, and watches what happens. Google described its annual, company-wide, multi-day Disaster Recovery Testing event, DiRT, in ACM Queue in 2012. AWS's Fault Injection Service can now pause a DynamoDB global table's replication to and from one region to simulate losing it, so the failover can be rehearsed without a real outage.

Even done well, active-passive leaves two problems untouched. Asha still waits 200 ms for every write, because all writes go to Virginia. And Oregon's servers spend their lives idle. Both point the same way: let every region take writes.

05Both regions taking writes

5.1Why it's attractive

Tokri now adds a region in Mumbai, ap-south-1, and makes it a full peer of Virginia. Asha's phone talks to Mumbai and Kabir's to Virginia; each region serves its own users' reads and writes locally, and each sends its changes to the other in the background. This is active-active: more than one region serving traffic, including writes, at the same time.

Everything gets better at once. Asha's taps take a few milliseconds instead of 200. Both regions do useful work, so the second region's cost buys capacity as well as insurance. And if Mumbai fails, its users can be sent to Virginia, which is already running at full size, so "failover" shrinks to moving traffic. AWS's guide observes that "most customers find that if they are going to stand up a full environment in the second Region, it makes sense to use it active/active."

Active-active: each region takes writes and replicates to the other
MUMBAIVIRGINIAchanges, both waysAshaMumbaiKabirNew JerseyApp serversap-south-1Lists DBMumbai copyApp serversus-east-1Lists DBVirginia copy
Step 1. Asha adds tea to the list. Her request goes to Mumbai, a few milliseconds away, and Mumbai's copy of the database stores it at once.
1 / 4

5.2The problem: two writes, two regions, one list

The last step of that diagram is the whole difficulty. In active-passive there was one primary, so all writes had one order. Now two regions accept writes to the same list without waiting for each other, so for about a tenth of a second after any write, the other region doesn't know about it, and anything written there in that window was written without seeing it.

Here's the concrete version. The list's title is "Mum's house". Asha renames it "Mum and Dad" at the same moment that Kabir renames it "Mum's place". Each region stores its own user's rename and tells them it's saved. Then the changes cross the ocean. Mumbai receives "Mum's place", Virginia receives "Mum and Dad". If each region applies whatever it receives last, Mumbai ends with Kabir's title and Virginia with Asha's: two copies of one list, permanently different, and each sibling sees the other's choice. Two writes like this, made in different places with neither aware of the other, are concurrent, and a pair of concurrent writes to the same data that the system has to choose between or combine is a conflict.

There are only three ways out, and the rest of this chapter takes them in turn:

  1. Don't let the conflict happen: give each list one home region that takes all its writes (section 6).
  2. Don't commit until the regions agree: use a consensus protocol across regions so there's one order again, and pay a round trip on every write (section 7).
  3. Let it happen, then merge: accept both writes and have every region apply the same rule to reach the same answer (sections 8 and 9).

06One home per list

6.1Partitioning by owner

The cheapest way to avoid conflicts is to make sure they can't arise. Give every list a home region, the one region allowed to accept writes to it. Asha created the list, so its home is Mumbai. Asha's writes land in Mumbai directly; Kabir's writes to the same list are forwarded from Virginia to Mumbai. Mumbai is the single primary for this list and gives all its writes one order, exactly as in active-passive, but for a different list, Kabir's own grocery list, the home is Virginia. Every region is a primary for some data and a replica for the rest.

AWS's guide calls this a "write partitioned" strategy: "assigns writes to a specific Region based on a partition key (like user ID) to avoid write conflicts." It's the same partitioning chapter 29 uses to spread data over machines, with regions as the partitions.

Home-region routing: each list's writes go to the region that owns it
forwarded writereplicateAshaMumbaiKabirNew JerseyApp, Mumbaiknows each list's homeApp, Virginiaknows each list's homeMumbai DBhome of 'Mum's house'Virginia DBhome of Kabir's own listHome directorylist → region
Step 1. Asha adds tea to 'Mum's house'. Mumbai's app looks up the list's home, finds Mumbai, and writes locally. Asha waits a few milliseconds.
1 / 4

Notice who pays. Its owner writes locally. Everyone else pays one round trip to the home on each write, and reads locally but slightly stale. For data that's mostly touched by one person, such as a user's profile, settings, drafts or private lists, almost nobody pays anything. For a shared list, whoever isn't at home pays the long trip.

6.2Choosing and moving the home

The home can be chosen in several ways: where the user signed up, where they are most of the time, or by law (Indian users' payment records homed in India). Some systems move the home when the user moves, which is sometimes called "follow the user": if Asha spends three months in London, her data's home can migrate to Ireland. Moving a home is a small failover for one key: stop writes to it, let the new home catch up, record the new home in the directory, resume.

Some databases do this for you. CockroachDB's REGIONAL BY ROW tables add a hidden column, crdb_region, to every row, and store and serve each row from that region; the column "defaults to the region of the gateway node that inserted the row", so a list created by Asha in Mumbai lives in Mumbai unless told otherwise. Updating crdb_region moves the row.

?Why not use home regions for everything?

Two reasons. First, when a home region fails, its data is in exactly the position active-passive left us in: someone has to fail that slice over to another region, with the same detection, fencing and RPO questions as section 4. Home regions shrink the blast radius, because only Mumbai-homed lists are affected, but don't remove it. Second, some data has no natural owner. A shared list edited equally from two continents, a counter of how many people bought something, a global product catalogue: for these, every choice of home makes half the writers pay the long trip. The next two sections handle data like that.

07Agreeing before committing

7.1Consensus across regions

The second way out is to give up on writing locally and make regions agree on every write before it counts. Chapter 27 showed how a group of servers agrees on one order of writes with a consensus protocol such as Raft or Paxos: a leader proposes each write, and it's committed once a majority of the servers, also called a quorum, has stored it. Any two majorities share at least one server, so no two leaders can commit conflicting writes, and the group keeps working as long as a majority is alive and connected.

Put one replica in each of three regions, say Mumbai, Singapore and Virginia, and you have a system where every write has one order, nothing committed is ever lost when one region fails, and the survivors carry on without a human. The cost is that a write commits only when the leader has heard back from one other region. If the leader is in Mumbai, the nearest other replica is Singapore, about 50 ms away, so every write waits roughly 50 ms at minimum. Writes from Virginia have to reach the leader in Mumbai first, which adds the 200 ms trip.

?Why three regions and not two?

Because a majority of two is two. With replicas in only Mumbai and Virginia, a write needs both, so losing either region, or the link between them, stops all writes; this is no better than synchronous replication. With three, any two form a majority, so one region can fail and the other two continue. That's why every system in this section asks for three regions as a minimum to survive the loss of one, and why the third can be a cheap witness: a replica that votes and stores the log of writes but serves no reads. CockroachDB's survival goals page says it directly: "To survive region failures, you must add at least 3 database regions."

7.2Spanner and CockroachDB

Google Spanner is the best-known database built this way. In a Spanner multi-region configuration, two regions hold read-write replicas and one is the default leader region, where the leaders live; a third region holds a witness. Each write needs "a write quorum that's composed of a majority of voting replicas", including one in the leader region. In exchange, Google offers a 99.999% availability target for multi-region configurations against 99.99% for a single region. Reads don't need a quorum; reads in regions far from the leader can be served by local replicas if the app accepts data a little out of date, and Google suggests a staleness of at least 15 seconds for them.

Spanner also lets a transaction read a consistent snapshot of the whole database at any timestamp, from any region, which depends on TrueTime: a clock API backed by GPS receivers and atomic clocks in every data centre that returns the current time as an interval guaranteed to contain the true time, about 4 ms wide on either side in the 2012 paper. Before a write becomes visible, Spanner waits until its timestamp is certainly in the past everywhere, a short pause called commit wait. Chapter 26 works through it in detail. For multi-region design, the point is that TrueTime lets Spanner order writes made in different regions by their timestamps without extra messages, and that the price of doing so is a few milliseconds per write, small next to the cross-region round trip the write pays anyway.

A Spanner multi-region configuration
DEFAULT LEADER REGIONSECOND REGIONWITNESS REGIONAppa writeLeaderplus a replicaTwo replicasread-writeWitnessvotes, no reads
Step 1. A write goes to the leader, which lives in the default leader region.
1 / 4

CockroachDB uses the same structure, Raft groups spread over regions, with ordinary clocks and a configured maximum offset between them instead of TrueTime. It makes the trade-offs per table:

Table localityWhere it's fastWhat it costsGood for
REGIONAL BY TABLEReads and writes in the table's home regionSlow from every other regionData used mostly in one place
REGIONAL BY ROWEach row fast in its own home regionSlow for a row away from its homePer-user data: section 6, built in
GLOBALReads fast in every regionWrites slower everywhereRarely changed data read everywhere: a product catalogue, currencies

Look closer at the last row. A GLOBAL table lets every region read current data locally, without asking the leader, by having writes use a "non-blocking transaction protocol" with a commit timestamp set slightly in the future, far enough ahead that by the time any region's clock reaches it, every region has the write. Reads then never wait; writes wait out that future timestamp, which depends on the maximum clock offset, so CockroachDB recommends lowering the offset to 250 ms for multi-region clusters. A catalogue that changes a few times an hour and is read millions of times a day is exactly what that trade fits.

Survival is also a per-database choice. With the default ZONE survival goal, a database survives losing an availability zone; writes commit within one region, and a region's loss makes its data unavailable until it returns. With REGION, CockroachDB keeps five replicas of each piece of data across three or more regions, and "write latency will be increased by at least as much as the round-trip time to the nearest region."

7.3DynamoDB's strong mode

Amazon added the same option to DynamoDB global tables, generally available from 30 June 2025. A global table created in multi-Region strong consistency (MRSC) mode is, in its documentation's words, "synchronously replicated to at least one other Region before the write operation returns a successful response", and strongly consistent reads in any region "always return the latest version of an item". It must be deployed in exactly three regions: three full replicas, or two replicas and a witness that AWS runs for you. As of October 2026 the docs list fifteen regions where it's available, Mumbai among them, in any combination of three, including across continents.

Its documentation is candid about the price. Writes and strongly consistent reads "require cross-Region communication", so their latency depends on the distance between the chosen regions. A write to an item that's being modified in another region at the same moment fails with a ReplicatedWriteConflictException, and the caller retries. Transactions across several items aren't supported in this mode, and neither is automatic expiry of items. The reward is an RPO of zero: a write that returned success is in two regions.

That leaves data that's written from everywhere, where waiting 50 to 200 ms per write isn't acceptable, and where a little temporary disagreement is fine, as long as everyone ends up agreeing. Tokri's shared list is exactly that. For data like this, both regions accept the write immediately, and we need a rule for combining what they accepted.

08Accepting both writes: last-writer-wins

8.1Keep the newest

The simplest rule is to stamp each write with the time it was made and, when two versions of the same item meet, keep the one with the later timestamp and discard the other. This is last-writer-wins, or LWW. If two stamps are equal, a fixed tie-break, such as the region's name, decides, so every region picks the same winner. Since every region applies the same rule to the same set of versions, all regions end up with the same value however the messages arrive. Systems that eventually agree like this are called eventually consistent.

DynamoDB global tables use exactly this in their default mode, multi-Region eventual consistency (MREC): "If the same item is modified in multiple Regions simultaneously, DynamoDB will resolve the conflict by using the modification with the latest internal timestamp on a per-item basis." Cassandra does it per column.

Here's what that does to the renaming conflict from section 5.2, step by step.

A conflict under last-writer-wins
'Mum and Dad''Mum's place'Asharenames at 10:00:05Mumbai copyclock correctVirginia copyclock 3 s slowKabirrenames at 10:00:07LWW rulekeep higher stamp
Step 1. At 10:00:05 by a correct clock, Asha renames the list 'Mum and Dad'. Mumbai stamps it 10:00:05 and says saved.
1 / 5

8.2What LWW loses

LWW has two separate problems, and the diagram shows both.

The first is clocks. Servers' clocks are kept close to true time by NTP, but "close" between regions means milliseconds at best and much more on a bad day (chapter 26 has the numbers). When two writes are closer together than the clocks' disagreement, LWW's choice between them is decided by whichever clock is ahead, not by which write came later. DynamoDB's documentation says its timestamps are "internal" and doesn't say how they're generated.

The second is worse, because it happens even with perfect clocks: LWW throws away one of two writes that were both meant to count. Suppose the list item "tea" has a quantity field. Asha changes it from 0 to 1 and Kabir, at the same moment, from 0 to 1, both meaning "add one pack". Two packs is the right answer. LWW keeps one write and returns 1. The same happens if the whole list is stored as one item and Asha adds tea while Kabir adds biscuits: one whole list wins, and one addition vanishes. Both writes could have been kept, but LWW only knows how to choose.

?Why do systems use it, then?

Because it's simple, it always converges, it needs no extra metadata, and for many kinds of data it's exactly right: a user's profile photo, the last-known location of a delivery driver, a setting that's changed by one person at a time. If one write fully replaces the other and nobody minds which survives a near-simultaneous pair, LWW is the right rule. The trouble starts when writes are meant to add up, as with quantities and lists, and LWW quietly loses some of them.

8.3Knowing that writes were concurrent

Before we can do better than "pick one", a region needs to know whether two versions conflict at all. If Virginia's version of the list was written after Kabir had seen Asha's change, then Kabir's version already includes hers and should replace it. Only if neither saw the other are they concurrent and in need of merging. Timestamps can't tell these apart; something has to track who has seen what.

That something is a version vector: a small table, stored with each item, holding one counter per region, counting the writes from that region that this version includes. Asha's write in Mumbai turns {mumbai: 0, virginia: 0} into {mumbai: 1, virginia: 0}. If Kabir edits that version after it reaches him, his write becomes {mumbai: 1, virginia: 1}, and since every entry is at least as large as in Asha's version, Virginia's version descends from hers and replaces it. If instead Kabir edited the older version, his is {mumbai: 0, virginia: 1}: each version has an entry larger than the other's, so neither includes the other, and they're concurrent. Applied to processes instead of regions, with one counter per process and used to order events instead of versions, the same idea is a vector clock.

A timeline diagram with three horizontal lines labelled A, B and C, arrows for messages between them, and boxes under each event listing counters such as A:2 B:4 C:1; a blue shaded region labelled Cause and a pink region labelled Effect
Vector clocks for three processes A, B and C. Each event carries one counter per process, and each message carries the sender's vector, which the receiver merges by taking the larger value in each entry. Events in the blue area happened before B's event with A:2 B:4 C:1; events in the pink area happened after it; the two white areas marked Independent are concurrent with it, and only comparing vectors can tell them apart.Diagram: translated by Duesentrieb from German Wikipedia, CC BY-SA 3.0, via Wikimedia Commons

Amazon's 2007 Dynamo paper (the ancestor of DynamoDB's design, though not of its current code) kept both concurrent versions, called siblings, and handed them to the application to merge on the next read; its shopping cart merged by taking the union of items, which occasionally brought deleted items back. Chapter 26 covers version vectors and that story in detail. Detecting concurrency is half of the job. The other half is the merge, and writing a correct merge for every kind of data by hand is where most teams go wrong. The next section is about data types that bring their own.

09CRDTs: data types that merge themselves

9.1What makes a merge safe

Suppose every region keeps its own copy of some piece of state and, from time to time, sends a copy to the others, which combine it with their own using a function merge. For all regions to end up identical regardless of the order, timing or repetition of these messages, merge needs three properties:

  • Commutative: merge(a, b) equals merge(b, a), so it doesn't matter which region's state arrives first.
  • Associative: merging a with b and then with c gives the same as merging b with c and then with a, so it doesn't matter how messages are grouped.
  • Idempotent: merge(a, a) equals a, so receiving the same state twice changes nothing, and retries are harmless.

A data type whose merge has all three, and whose updates only ever move the state "forward" in a way merge respects, is a conflict-free replicated data type, or CRDT. Marc Shapiro, Nuno Preguiça, Carlos Baquero and Marek Zawirski formalised them in 2011. Their guarantee, called strong eventual consistency, is that any two replicas that have received the same set of updates are in the same state, without any coordination at all. Chapter 28 introduced two of them briefly; here we build four from scratch, each on Tokri's list.

9.2Counters: G-counter and PN-counter

Start with the tea quantity. A single number can't be merged: if Mumbai has 2 and Virginia has 1, the merged value could be 2, 3 or 1 depending on history the numbers don't record.

So stop storing one number. A grow-only counter, or G-counter, stores one number per region: how many increments that region has made. Each region only ever increases its own entry. The counter's value is the sum of the entries, and merge takes the larger value of each entry.

EventMumbai's copyVirginia's copyValue seen
Start{M: 0, V: 0}{M: 0, V: 0}0, 0
Asha adds 2 in Mumbai{M: 2, V: 0}{M: 0, V: 0}2, 0
Kabir adds 1 in Virginia{M: 2, V: 0}{M: 0, V: 1}2, 1
States exchanged and merged{M: 2, V: 1}{M: 2, V: 1}3, 3

Taking the larger value per entry is commutative, associative and idempotent, so the sum comes out at 3 whichever order things arrive in, and a duplicated message doesn't count twice. That's the reason for one entry per region: each entry has only one writer, so it never has a conflict of its own.

A G-counter can only grow. To allow decrements, a PN-counter keeps two G-counters, P for increments and N for decrements, and reports P minus N. When Asha changes her mind and removes one of her two packs, Mumbai's N entry goes from 0 to 1. Both halves merge as before, and the value comes out as (2 + 1) − 1 = 2 everywhere.

9.3A register: LWW, made explicit

The list's title is a single value that's replaced, not added to: there's no sensible way to combine "Mum and Dad" with "Mum's place". For this, a last-writer-wins register is a CRDT: it stores the value together with a (timestamp, region) pair, and merge keeps the one with the larger pair. That's commutative, associative and idempotent, so it converges. It's the same rule as section 8, and it still loses the earlier of two concurrent titles, but here it's a deliberate choice for one field that should hold one value, not a rule applied blindly to the whole list. (Some systems offer a multi-value register instead, which keeps all concurrent values and lets the user pick, as a file-sync app does when it creates a "conflicted copy".)

9.4Sets: the OR-set

The items on the list form a set, and sets are where naive merges go wrong. Taking the union of two copies is commutative, associative and idempotent, but it can never remove anything: a deleted item comes straight back from the other copy, which is the Dynamo cart problem. Keeping a separate set of removed items and subtracting it is no better: an item removed once can then never be added again.

The observed-remove set, or OR-set, solves this by giving every add a unique tag. Adding "milk" stores the pair ("milk", tag a1). Removing "milk" doesn't record "milk is gone"; it records the specific tags of "milk" that this region can currently see, here a1, as removed. An item is in the set if it has at least one tag that hasn't been removed. Merge takes the union of the added pairs and the union of the removed tags.

Walk through the case that breaks everything else. The list holds milk with tag a1, in both regions. Kabir, who has just bought milk, removes it in Virginia, which marks a1 removed. At the same moment Asha, who has heard their mother needs more, adds milk again in Mumbai, which creates ("milk", b7). After merging, both regions have tags a1 and b7 for milk, and only a1 is removed. Milk is on the list. The remove deleted only the add it had seen, and the concurrent add it couldn't have seen survives. When an add and a remove of the same item are concurrent, the add wins; when the remove came after the add, it sticks.

Predict before you read on

With an OR-set, Asha adds 'eggs' in Mumbai. A minute later, after the add has replicated, Kabir removes 'eggs' in Virginia, and at the same time Asha removes 'eggs' in Mumbai. What does the list show once everything has merged?

9.5All four on one list

The script below builds a PN-counter, an LWW register and an OR-set in a few lines each, with merge as the only way state passes between regions. It starts both regions in sync with milk on the list, then runs the edits from this section concurrently: Asha adds two packs of tea and removes one, adds milk and tea, and renames the list; Kabir adds a pack of tea, removes milk, and renames the list two seconds after Asha, on a clock that's three seconds slow. Then it merges in both directions. The last line applies plain last-writer-wins to the tea quantity, as if each region had stored one number. (Python's uuid module makes the unique tags.)

Merge a PN-counter, an LWW register and an OR-set from two regions
python
Python
import uuid
 
# PN-counter: per region, a count of increments (P) and of decrements (N)
def pn_add(c, region, n):
    P, N = dict(c[0]), dict(c[1])
    if n >= 0: P[region] = P.get(region, 0) + n
    else:      N[region] = N.get(region, 0) - n
    return (P, N)
def pn_merge(a, b):
    mx = lambda x, y: {r: max(x.get(r, 0), y.get(r, 0)) for r in sorted(x.keys() | y.keys())}
    return (mx(a[0], b[0]), mx(a[1], b[1]))
def pn_value(c): return sum(c[0].values()) - sum(c[1].values())
 
# LWW-register: keep the value with the highest (timestamp, region)
def lww_merge(a, b): return max(a, b)
 
# OR-set: each add gets a unique tag; a remove deletes only the tags it has seen
def or_add(s, item):    return (s[0] | {(item, uuid.uuid4().hex[:6])}, s[1])
def or_remove(s, item): return (s[0], s[1] | {t for (i, t) in s[0] if i == item})
def or_merge(a, b):     return (a[0] | b[0], a[1] | b[1])
def or_value(s):        return sorted({i for (i, t) in s[0] if t not in s[1]})
 
# Both regions start in sync: the list holds milk, and 0 packs of tea
start_items = or_add((set(), set()), "milk")
mum = {"tea": ({}, {}), "items": start_items, "title": (100, "mumbai", "Mum's house")}
vir = {"tea": ({}, {}), "items": start_items, "title": (100, "mumbai", "Mum's house")}
 
# Concurrently, before either region hears from the other:
mum["tea"] = pn_add(mum["tea"], "mumbai", +2)       # Asha: +2 packs of tea
mum["tea"] = pn_add(mum["tea"], "mumbai", -1)       # Asha: one fewer after all
vir["tea"] = pn_add(vir["tea"], "virginia", +1)     # Kabir: +1 pack
vir["items"] = or_remove(vir["items"], "milk")      # Kabir: milk is bought, remove it
mum["items"] = or_add(mum["items"], "milk")         # Asha: Mum needs more milk, add it
mum["items"] = or_add(mum["items"], "tea")
mum["title"] = (105, "mumbai", "Mum and Dad")       # Asha renames at real time 105
vir["title"] = (104, "virginia", "Mum's place")     # Kabir renames at real time 107,
                                                    # but Virginia's clock is 3 s slow
print("before sync  Mumbai:  ", pn_value(mum["tea"]), or_value(mum["items"]), mum["title"][2])
print("before sync  Virginia:", pn_value(vir["tea"]), or_value(vir["items"]), vir["title"][2])
 
for k, merge in (("tea", pn_merge), ("items", or_merge), ("title", lww_merge)):
    merged_m, merged_v = merge(mum[k], vir[k]), merge(vir[k], mum[k])
    mum[k], vir[k] = merged_m, merged_v
print("after sync   Mumbai:  ", pn_value(mum["tea"]), or_value(mum["items"]), mum["title"][2])
print("after sync   Virginia:", pn_value(vir["tea"]), or_value(vir["items"]), vir["title"][2])
print("tea counter state:", mum["tea"])
 
# The same quantity edits as a plain number with last-writer-wins
plain_mumbai, plain_virginia = (103, "mumbai", 1), (106, "virginia", 1)
print("plain number + LWW, packs of tea:", lww_merge(plain_mumbai, plain_virginia)[2])
output
C++
before sync  Mumbai:   1 ['milk', 'tea'] Mum and Dad
before sync  Virginia: 1 [] Mum's place
after sync   Mumbai:   2 ['milk', 'tea'] Mum and Dad
after sync   Virginia: 2 ['milk', 'tea'] Mum and Dad
tea counter state: ({'mumbai': 2, 'virginia': 1}, {'mumbai': 1})
plain number + LWW, packs of tea: 1

Before the sync, the regions disagree on everything: Virginia's list is empty because Kabir removed the only milk it knew about. After it, they agree on everything, and merging in opposite orders made no difference. Look at each line in turn:

  • Tea is 2, Asha's net one plus Kabir's one, and the state line shows why: Mumbai's increments 2, Virginia's 1, Mumbai's decrements 1. A plain number under LWW ends at 1 and loses a pack without telling anyone.
  • Milk is on the list, because Asha's new add carried a tag Kabir's remove never saw. Tea is there too.
  • The title is "Mum and Dad", though Kabir renamed it later. That's the LWW register doing what section 8 warned about, kept here on purpose for one field where losing a near-simultaneous rename is acceptable.

9.6What CRDTs cost, and what they can't do

The OR-set's removed tags are tombstones: records that something was deleted, kept so the deletion can win against old copies. In this toy version they grow forever. A real system can drop a tombstone only once every region has seen the remove, which means tracking what each region has seen, usually with version vectors; that bookkeeping is where most of a production CRDT's complexity lives. A CRDT value is also bigger than the plain value: one entry per region per counter, one tag per add.

The deeper limit is that CRDTs converge, but they can't enforce a rule that needs the regions to agree before acting. "Never sell more tea than is in the warehouse" can't be checked by Mumbai alone if Virginia might be selling the last pack at the same moment; both see one left, both sell it, and the merged stock count is −1, perfectly converged and wrong. "Only one user may have the username asha" has the same shape. Rules like these, called invariants, need a single order of the writes that might break them, which means a home region (section 6) or consensus (section 7) for exactly those writes.

10What the real systems do

10.1Redis Active-Active

Redis Software and Redis Cloud offer Active-Active databases, and they are a CRDT system you can use through ordinary Redis commands. Each participating cluster, at least two and normally up to ten, holds a full copy and accepts reads and writes locally; a process called the syncer streams each cluster's changes to the others. Its documentation says Active-Active databases "only use conflict-free replicated data types", and that they follow "a strong eventual consistency model, which means that local values may differ across replicas for short periods of time, but they all eventually converge to one consistent state." Internally, it tracks what each copy has seen with vector clocks.

Redis Active-Active: two clusters, one counter
App, MumbaiApp, VirginiaRedis clusterMumbai, full copyRedis clusterVirginia, full copySyncerstreams changes
Step 1. INCRBY counter 7 in Mumbai. The write is local and fast, and Mumbai now reads 7.
1 / 3

What makes it useful to study is that the documentation spells out the merge rule for each Redis type, and they map exactly onto section 9:

Redis typeMerge rule, per the docsSection 9 equivalent
String set with SET"Last write wins", by the operating system's wall clock; "the only case with Active-Active databases where OS time is used to resolve a conflict"LWW register
String used with INCR, INCRBY, DECRBYCounters "accumulate the total counter operations across all member Active-Active databases"PN-counter
Set"OR-Set" behaviour: "remove can only remove instances it has already seen and in all other cases element add wins"OR-set
HashFields added or removed with OR-set rules; each field's value behaves as a string or counterOR-set of registers and counters
Sorted setOR-set for membership; a score set concurrently by ZADD is last write wins, and ZINCRBY is summed like a counterOR-set plus registers and counters
ListConcurrent inserts are all kept in an arbitrary but consistent order; DEL removes only the elements it observedA sequence CRDT

Their own counter example is the one we worked through: INCRBY counter 7 in one region and INCRBY counter 3 in another both read their own value until a sync, after which both read 10.

The documentation is equally clear about what this doesn't give you. Two concurrent RPOPs from the same list in two regions can return the same element, because "Lists in Active-Active databases guarantee that each element is POP-ed at least once, but cannot guarantee that each element is POP-ed only once"; the docs recommend popping from only one region if an item mustn't be handled twice. When two regions set different expiry times on a key concurrently, the longer one wins. Deleted keys stay in memory as tombstones until every instance has seen the delete. Memory planning should allow for the CRDT metadata: the docs say Active-Active "requires double the memory of regular replication, which can be up to two times (2x) the original data size per instance".

An optional causal consistency setting guarantees that if a region had seen operation A before it performed B on the same key, every region applies A before B. Its documented cost is that each instance relays the order of operations it received to all the others, which increases network traffic "by a factor of (N-2)" for N instances.

10.2DynamoDB global tables

DynamoDB global tables offer both ends of this chapter in one product. In the default MREC mode, every replica accepts writes, changes replicate asynchronously "typically within a second or less", and conflicts on an item are resolved by last-writer-wins on DynamoDB's internal timestamp. AWS describes the RPO of this mode as "equal to the replication delay between replicas, usually a few seconds". Two more details from the documentation matter for design. A strongly consistent read returns the latest version only if the item was last written in the same region, and a conditional write checks its condition against the local copy, so a check like "only if stock is above zero" can pass in two regions at once. And a transaction is atomic only in the region that ran it: other regions may briefly see some of its writes and not others.

DynamoDB MREC: one item in stock, sold twice
BuyerMumbaiBuyerVirginiaReplicaMumbai, stock 1ReplicaVirginia, stock 1Replicationasync, about 1 s
Step 1. Both buyers send "take one if stock is above zero" at the same moment. Each condition is checked against the local copy, which says 1, so both succeed.
1 / 3

In MRSC mode, covered in section 7.3, writes reach a second region before succeeding and every strongly consistent read is current, in exchange for cross-region latency on writes, retries on ReplicatedWriteConflictException, and no transactions. A table's mode is chosen when it's created and can't be changed. The documentation doesn't describe how MRSC works internally.

10.3Side by side

Redis Active-ActiveDynamoDB MRECDynamoDB MRSCSpanner / CockroachDB multi-region
Where writes are acceptedEvery regionEvery regionEvery region, coordinatedLeader per piece of data
Write latencyLocalLocalCross-region round tripRound trip to a majority
ConflictsMerged per type (CRDTs)Last writer wins, per itemRejected; the caller retriesCan't happen: one order
Data lost if a region diesUnreplicated recent writesA few seconds (AWS's figure)NoneNone
Invariants across regionsNoNoYes, per itemYes, with transactions

11The parts people forget

11.1IDs and sequences

With the database sorted out, a set of smaller problems appears, each of which can quietly break a multi-region launch. The first is IDs. Suppose Tokri keeps its orders in a relational database that numbers rows with an auto-incrementing column: the database hands out 1, 2, 3 in order. With two regions both accepting inserts, Mumbai and Virginia each hand out 1,042 to different orders within the same second, and when the rows replicate, two orders have one ID.

The old fix for two-primary MySQL is to give each region its own slice of the numbers: setting auto_increment_increment to 2 and auto_increment_offset to 1 in one region and 2 in the other makes one region hand out odd numbers and the other even ones. It works, and it's brittle: adding a third region means renumbering the scheme.

The more common fix is to stop asking a central counter at all. Twitter's Snowflake, published in 2010, packs a 64-bit ID from three parts: 41 bits of milliseconds since a chosen start date ("gives us 69 years"), 10 bits of machine ID ("up to 1024 machines"), and 12 bits of sequence number within the millisecond ("rolls over every 4096 per machine"). Machines with different IDs can never collide, and IDs still sort roughly by creation time. UUIDv7, standardised in RFC 9562 in May 2024, does the same in 128 bits with a 48-bit millisecond timestamp followed by random bits, so it needs no machine IDs at all.

11.2Unique names

The second problem is uniqueness. Asha signs up for a Tokri username "asha" in Mumbai at the same moment someone in Virginia signs up for "asha". Each region checks its own copy, finds the name free, and creates the account. After replication there are two users named asha, and no merge can fix it: one of them has to be told, after the fact, that their username is gone.

This is the invariant problem from section 9.6, and it has the same fixes. Give the username table a home region, so every signup checks and claims names in one place, paying a round trip for users far from it; signups are rare, so that's usually fine. Or keep usernames in a strongly consistent multi-region store (CockroachDB, Spanner, a DynamoDB MRSC table) so the check and the claim are one coordinated write. What doesn't work is a uniqueness constraint in each region's database, because each region's constraint only knows its own rows.

11.3Caches and queues

Caches have to be thought about per region. Each region should have its own cache, near its app servers, because a cache across an ocean defeats the point. But when a list changes in Mumbai and replicates to Virginia, Virginia's cache still holds the old version until something removes it. The usual fix is to have each region invalidate its own cache when it applies a replicated change, not only when its own users write. Managed products can get this wrong in ways that are easy to miss: DynamoDB's documentation notes that writes replicated into a region bypass its in-memory cache, DAX, so "DAX caches can become stale" there until their entries expire.

Queues and event streams need the same care. If a write in Mumbai puts a "list changed" event on a queue in Mumbai to send a notification, and the replicated write also triggers an event in Virginia, Kabir gets two notifications. Decide where each kind of side effect happens, usually in the item's home region, and make consumers idempotent: each event carries an ID, and processing the same ID twice does nothing the second time. Redis's at-least-once list pops from section 10.1 are this problem in miniature.

11.4The cost of moving bytes

Traffic between regions is billed. On AWS, as of the September 2026 price list, sending data from us-east-1 to Mumbai costs $0.02 per GB, and from Mumbai to us-east-1 it costs $0.086 per GB; between us-east-1 and Ohio it's $0.01. Replication sends every write to every other region, so the bill grows with write volume times the number of regions. Here's Tokri's, if it replicates 2 TB of changes a month in each direction:

Mumbai to Virginia, 2 TB2,000 GB × $0.086$172 / month
Virginia to Mumbai, 2 TB2,000 GB × $0.02$40 / month
Three regions (add Ireland), 2 TB each way per pairMumbai 4,000 GB × $0.086 + Virginia, Ireland 8,000 GB × $0.02$504 / month
Two regions, transfer aloneabout $210 / month

That's probably modest for a company Tokri's size, and it's on top of the database's own charges for applying each replicated write in each region. It becomes significant for systems that replicate large objects or high-volume streams everywhere. Usually the savings come from replicating only what other regions need (a list's home region needs every change; a region that only reads it may need it only when someone opens it), and from compressing replication traffic.

11.5Testing a region's loss

Section 4.5's advice applies with more force to active-active, because "failover" is now a traffic shift that happens rarely enough to rot. AWS's guide frames the questions for active-active testing as: "Is traffic routed away from the failed Region? Can the other Region(s) handle all the traffic?" The second one is easy to get wrong: two regions each running at 70% of capacity can't absorb each other's load. Run a game day that drains one region of traffic entirely, watch the survivor's capacity and latency, and use fault injection to pause replication to one region and see what users of the others experience.

And keep backups. Replication copies mistakes as faithfully as good writes: a bad deployment that deletes every list in Mumbai deletes them in Virginia a second later. AWS's guide notes that for data corruption or deletion, "the recovery point will always be at some point before the disaster was discovered", which is only true if there's a point-in-time backup to recover from.

12Choosing a design

12.1Tokri's choices, per kind of data

No single design fits all of an app's data, and in practice a multi-region app mixes them, choosing per table or even per field. Here's where Tokri ends up:

DataRule for concurrent writesDesignWhy
A user's profile, settingsOne writer at a time; replaceHome region per user, LWW elsewhereAlmost only the owner writes it
Shared list items and quantitiesShould add upCRDTs (OR-set, PN-counters), every region writesEdited from everywhere; nothing to lose by merging
List titleReplace; either is fineLWW registerOne value; a lost near-simultaneous rename is acceptable
UsernamesMust be uniqueStrongly consistent multi-region table, or one homeAn invariant across all regions
Warehouse stockMust not go below zeroHome region per warehouseAn invariant, and each warehouse is in one place
Product catalogueRare writes, constant readsGlobal table: slow writes, local readsReads dominate
Payment records for Indian usersMust stay in IndiaHome region in India, by lawData residency
Decision

How should Tokri run in several regions?

Active-passive
One region takes all traffic; a warm standby takes over after a failover.
  • One primary: no conflicts
  • Simplest to reason about
  • Cheapest second region
  • Far users wait a round trip on every request
  • Failover loses the lag window and is risky if untested
chosen
Active-active, per-data choices
Every region serves its users; home regions, CRDTs, LWW or consensus chosen per table.
  • Local latency for most reads and writes
  • Losing a region means shifting traffic, not promoting a copy
  • Both regions earn their cost
  • Every table needs a conflict decision
  • IDs, uniqueness, caches and queues all need rework
Consensus for everything
A Spanner- or CockroachDB-style database with every table replicated across three regions.
  • One order of writes, transactions, RPO of zero
  • No merge logic in the app
  • Every write pays a round trip to a majority
  • Needs three regions minimum

Tokri's data is mostly lists that people edit from wherever they are, and a lost item there is a small annoyance, while a 200 ms wait on every tap is felt all day. So it goes active-active with CRDTs for the lists, and keeps a small strongly consistent store for the few things that must never be doubled: usernames, stock and the failover epoch. A bank making the same decision would weigh it the other way, and accept slower writes everywhere to get one order of every transaction.

12.2Rules that hold up

  • Decide what a region's loss should cost before building anything: the RPO and RTO you need, written down, per kind of data.
  • Keep app servers next to the database they use. Cross regions once per request, not once per query.
  • Fence before you promote. A deposed primary that can still write is how split-brain starts.
  • Choose a conflict rule per kind of data, explicitly: home region, consensus, CRDT or LWW. "Whatever replication does" is a rule too, usually the wrong one.
  • Use three regions for anything that needs a majority, even if the third is only a witness.
  • Generate IDs without a central counter, and route uniqueness checks to one place.
  • Watch replication lag and alert on it; your real RPO is your lag during incidents.
  • Practise losing a region on a schedule, and keep backups that replication can't overwrite.

12.3Symptom, cause, fix

SymptomLikely causeFix
Users far from the primary find every page slowEvery write, or every query, crosses regionsAdd a region near them; keep app and database together (§2.2, §5)
Edits "disappear" a few seconds after being savedLWW discarding a concurrent write, or a skewed clock winning racesA CRDT for data that should merge; check clock offsets (§8, §9)
A quantity or counter is lower than the sum of everyone's changesRead-modify-write of a plain number under LWWPN-counter, or route that counter to one home (§9.2)
Deleted items come backUnion merge without tombstonesOR-set semantics (§9.4)
Two accounts with the same usernameEach region checked uniqueness against its own copyOne home for the check, or a strongly consistent table (§11.2)
Duplicate-key errors on replicated insertsAuto-increment IDs issued in two regionsSnowflake-style or UUIDv7 IDs (§11.1)
After a failover, writes from both regions can't be reconciledSplit-brain, no fencingEpoch numbers or STONITH before promotion (§4.3)
Failover takes far longer than plannedStandby needs scaling during an incident; untested runbookSize the standby for minimum service; game days (§3.4, §4.5)
One region serves stale data after another region's writeCache entries not cleared when replicated changes arriveInvalidate on apply, in every region (§11.3)

13Summary

  1. A region fails as a unit. Availability zones protect against a building; shared software, control systems and DNS can take a whole region down for hours, as us-east-1 showed in December 2021 and October 2025.
  2. Distance between regions is set by the speed of light in glass. About 200 km per millisecond, so Mumbai to Virginia is at least 130 ms round trip and about 200 ms in practice; the link itself can also break.
  3. Active-passive keeps a standby copy, usually asynchronously. Writes stay fast, and the replication lag at the moment of failure is the data you lose: the RPO. The RTO is detection plus decision plus promotion plus steering.
  4. Failover is four steps, each fallible. Health checks are deliberately slow; DNS has caches to wait out and anycast doesn't; a deposed primary must be fenced, or split-brain follows; failing back is a second failover.
  5. An untested failover is a hope. GitHub's 43-second partition became 24 hours; game days and fault injection are how teams find the rot first.
  6. Active-active puts users near a writable copy, and makes conflicts possible. There are three ways out: one home per item, consensus before commit, or accept and merge.
  7. Home regions avoid conflicts by giving each item one writer. The owner writes locally, others pay a round trip, and a home's loss is a small active-passive failover.
  8. Cross-region consensus gives one order and an RPO of zero. Spanner, CockroachDB and DynamoDB MRSC need three regions and charge each write a round trip to a majority.
  9. Last-writer-wins converges by discarding writes. Clock skew decides close races, and writes meant to add up get lost even with perfect clocks; version vectors can at least detect concurrency.
  10. CRDTs merge without coordination, but can't protect invariants. G- and PN-counters, LWW registers and OR-sets converge in any order, as in Redis Active-Active; stock levels and unique names still need one place to decide.
  11. IDs, uniqueness, caches, queues, transfer costs and backups all change with a second region, and each has a standard fix.

14Build this

Run two "regions" on your laptop, take one away, and watch what each design loses.

  • Start two Redis servers on different ports, A and B, with a script that copies each one's changes to the other through a queue you can delay by 200 ms and pause. Use tc qdisc add dev lo root netem delay 100ms on Linux, or a sleep in the copier, for the distance.
  • Active-passive: write only to A, replicate to B, and kill A at a random moment under a steady write load. Count the writes acknowledged by A that B never received, and compare with the lag your copier reported just before.
  • Split-brain: pause the copier without killing A, promote B, and keep writing to both for 30 seconds. Try to merge the two afterwards. Then add an epoch number that B increments on promotion and that the copier checks, and repeat.
  • Last-writer-wins: let both accept writes to the same keys with timestamps, and skew one side's clock by 2 seconds. Count lost increments to a shared counter.
  • CRDTs: replace the counter with a PN-counter and the item list with an OR-set (the TryIt in section 9.5 is a start), stored as Redis hashes. Merge in random orders, with duplicated and delayed messages, and check that both sides always end identical and that no increment is lost.

15Interview questions

beginnerWhat are RPO and RTO, and how do they relate to synchronous and asynchronous replication?›

RPO, the recovery point objective, is how much recent data you can afford to lose in a disaster, measured in time. RTO, the recovery time objective, is how long you can be down before service is restored in another place.

Asynchronous replication acknowledges writes before the remote copy has them, so the data lost in a failover is whatever the replication lag was at that moment: an RPO of seconds, which grows during incidents. Synchronous replication waits for the remote copy, so the RPO can be zero, at the cost of a cross-region round trip on every write and a stall if the remote region is unreachable. RTO depends mostly on how much of the standby is already running (pilot light, warm, hot) and how fast you detect, decide and steer traffic.

intermediateWalk through a regional failover for an active-passive system. What can go wrong at each step?›

Detect: health checks from several locations fail several times in a row; too sensitive and you fail over on a blip, too slow and you're down longer. Decide: a human or automation declares the region dead; automatic failover on a false alarm loses the lag window for nothing, so many teams automate everything except the button. Fence: make sure the old primary can't write, by shutting it off or with an epoch number storage enforces; skip this and a partitioned-but-alive primary creates split-brain. Promote: the replica becomes primary, losing unreplicated writes. Steer: DNS or anycast sends users to the new region; DNS caches delay this, and the standby may need to scale up through a control plane that's also impaired.

Then fail back as a second, planned failover: rebuild the old region as a replica, let it catch up, pause writes briefly, switch.

intermediateTwo users in different regions register the same username at the same time in an active-active system. How do you prevent duplicates?›

You can't with per-region checks or CRDTs, because each region only knows its own writes until replication, and no merge rule can give one name to two people. Uniqueness is an invariant, so the check and the claim need a single order: route all username claims to one home region (the far user pays a round trip, which is fine for a rare operation), or store usernames in a strongly consistent multi-region store such as Spanner, CockroachDB or a DynamoDB table in multi-Region strong consistency mode, where the conditional write is coordinated across regions. A cheaper variant reserves names by hashing them to a home region, so the load spreads but each name still has exactly one place that decides.

intermediateWhat does last-writer-wins lose, and when is it still the right choice?›

It discards all but one of a set of concurrent writes. With clock skew, the winner of a close race is whichever region's clock is ahead, so a later write can lose to an earlier one. And even with perfect clocks, writes meant to combine, like two increments of a counter or two additions to a list stored as one value, lose one of them, silently.

It's right when one write fully replaces another and either outcome is acceptable: a profile photo, a setting, a device's last-seen location. DynamoDB global tables in their default mode and Redis Active-Active strings use it for that reason.

deepDesign a G-counter and an OR-set. Why do they converge, and what do they cost?›

A G-counter is a map from region to count; each region increments only its own entry, the value is the sum, and merge takes the per-entry maximum. Maximum is commutative, associative and idempotent, and each entry has one writer, so replicas that have seen the same updates agree in any order and duplicates don't double-count. A PN-counter pairs two of them for increments and decrements.

An OR-set tags each add uniquely; a remove records the tags of that element it can see; merge unions both the adds and the removed tags; an element is present if any of its tags isn't removed. A concurrent add has a tag the remove didn't see, so add wins over a concurrent remove, while a later remove sticks. The cost is metadata: one entry per region per counter, a tag per add, and tombstones that can only be collected once every region has seen the remove, which needs version-vector bookkeeping. Neither can enforce invariants like "stock stays at or above zero".

deepWhy do Spanner, CockroachDB and DynamoDB's strong mode all require three regions, and what does a write cost?›

They commit a write when a majority of voting replicas has it. With two regions the majority is both, so losing either region or the link between them stops writes. With three, any two are a majority, so one region can be lost and the other two continue, and since any two majorities overlap, no two sides of a partition can both commit. The third can be a witness that votes but serves no reads.

Each write waits for the leader's region plus the nearest other voting region, so at least one round trip between them: tens of milliseconds for regions on one continent, around 200 ms between India and the US East Coast, plus the trip to the leader for writes from elsewhere. Spanner adds a commit wait of a few milliseconds from TrueTime's uncertainty, small next to that round trip. Reads can be local if the app accepts slightly stale data, or if the table is designed for it, like CockroachDB's global tables, which move the cost onto writes.

16Go deeper

check yourself
Mumbai to Virginia is about 13,000 km. What's the lowest possible round-trip time over fibre, and what do real networks see?›

About 130 ms, at 200 km per millisecond there and back. Real medians are about 200 ms, because cables follow coastlines and seas.

What's the difference between pilot light and warm standby?›

Pilot light has the data replicated but servers switched off, so it can't serve until they're started; warm standby is a small working copy that can take traffic at once.

What does fencing prevent, and name one way to do it.›

A deposed primary that's still alive continuing to accept writes (split-brain). Shut it off through the provider's API, or have storage reject writes carrying an old epoch number.

Two regions each run INCRBY 5 on the same key in Redis Active-Active before syncing. What does GET return afterwards?›

The sum of all increments, 10 above whatever the starting value was, in both regions: counters merge as PN-counters.

Why can't a CRDT keep a warehouse's stock at or above zero?›

Two regions can each see one item left and each sell it without coordinating; the merge converges to −1. The check needs one place to decide: a home region or consensus.

AWS, Disaster Recovery of Workloads on AWS

Backup and restore, pilot light, warm standby and active-active, with RPO and RTO for each, data plane versus control plane, and the write-global, write-local and write-partitioned patterns. docs.aws.amazon.com

Shapiro, Preguiça, Baquero, Zawirski, Conflict-free Replicated Data Types (2011)

The definitions of state- and operation-based CRDTs and strong eventual consistency, with counters, registers and sets. lip6.fr

Corbett et al., Spanner (OSDI 2012)

Paxos groups across data centres, TrueTime's interval and commit wait, and how they give externally consistent transactions worldwide. research.google

DynamoDB global tables: how they work

MREC and MRSC side by side: last-writer-wins, replication latency, the three-region rule, witnesses, conflicts, transactions and fault injection. docs.aws.amazon.com

Redis Active-Active: develop and data types

The merge rule for every Redis type, with timelines, plus expiry, tombstones, causal consistency and at-least-once pops. redis.io

CockroachDB, table localities and survival goals

Regional-by-table, regional-by-row and global tables, and what surviving a region costs in replicas and write latency. cockroachlabs.com

Where you meet these ideas in the wild:

AWS us-east-1, 20 October 2025

DynamoDB errors for nearly three hours and dependent services for about fourteen, with replicas elsewhere reachable but lagging. (summary)

A race in DNS automation emptied a regional endpoint; global tables in other regions kept serving.
AWS us-east-1, 7 December 2021

Hours of impaired launches, console sign-in and support, from a capacity change on internal infrastructure. (summary)

Retry storms on an internal network, and a status page that struggled to fail over.
GitHub, 21 October 2018

The clearest public account of what unreplicated writes on both sides of a failover cost. (analysis)

A 43-second partition, an automatic cross-country failover, and 24 hours of reconciliation.
Microsoft, Azure network round-trip latency statistics

The table behind section 2's comparison; useful for choosing which regions to pair. (learn.microsoft.com)

Real median round trips between dozens of regions, updated every six to nine months.
28 · Replication & Consistency Models

Primaries, replicas, lag and consistency models inside one region, and a first look at last-writer-wins and CRDTs. Read it

27 · Consensus: Raft, Paxos & Leases

How majorities agree on one order of writes, which cross-region consensus spreads over three regions. Read it

26 · Time, Clocks & Ordering

Why clocks disagree, vector clocks and version vectors, TrueTime and commit wait, and fencing tokens. Read it

30 · Failure Detection & Membership

Why silence can't be told from death, and how to choose timeouts that don't fail over on a blip. Read it

35 · DNS, TLS & the Edge

DNS caching and TTLs, anycast, and the handshakes that make a far region's first request cost three round trips. Read it

22 · Redis

The single-node Redis that Active-Active databases extend with CRDT metadata. Read it