Ch.16: System Design: Active-Active Writes Are a Trap
Outline
- 0:00 When both regions sell the last unit
- 0:57 Entry regions versus write authority
- 2:16 Requirements and failure model
- 3:26 Recovery point and recovery time objectives
- 4:42 What active-active hides
- 6:04 Pattern 1: Single global writer
- 7:20 Pattern 2: One home region per key
- 8:59 The inventory invariant
- 10:06 Pattern 3: Globally ordered writes
- 11:52 What a network partition does
- 12:44 Why last-writer-wins cannot save inventory
- 14:10 Pattern 4: Mergeable writes
- 15:36 Clocks, TrueTime, and consensus
- 17:00 One database, four promises
- 18:02 Failover is authority transfer
- 19:37 Observability for business invariants
- 20:51 The completed multi-region design
Transcript
0:00 One unit left. A limited-run graphics card, SKU-RED-1, and every region on the planet shows one in stock. A shopper in Virginia clicks buy. The same second, a shopper in Tokyo clicks buy. Both of them get the green checkmark. Both? Then one of those confirmations is a lie. Not yet, which is the trap: each region read its local replica, saw one, wrote zero, and confirmed. Replication catches up and every copy now stores zero. The multi-region database is converged, the dashboard is green, and the business sold one physical item twice.
0:33 So today we design one of these that can't lie like that. The mechanism is deciding, for each piece of data, who orders its writes, what can merge later, and which writes have to stop when the network splits. Four write patterns, one two-question rule, and by the end you'll be able to classify every table you own. The map is where we start. So here's the design that has to keep that promise, all of it up front, and I want you to read it with 1 question in mind: not where requests go, but who gets to decide.
1:08 Top layer a global traffic manager routing users to three application cells, Virginia, Frankfurt, Tokyo. Under each cell, a database gateway that knows which data lives where. And then the parts people skip. The left half is everything that's copied outward so reads never leave the region, the catalog above all. Right, that's the half nobody argues about. The right half is write machinery: user rows with 1 home each, an inventory range with 1 leader and cross-region voters, and background streams that merge what they can, with a queue for what they can't.
1:44 Plus a failover controller that can fence a failed region, and backups that don't depend on live replication at all. Every arrow into a database here answers "where can a request land." It deliberately doesn't answer "where is the write decided." Keeping those two questions separate is basically the whole chapter. We'll come back to this exact picture at the end, and every box will answer one question or the other. And everything today comes from the vendors' own docs and the original Spanner paper, links in the description.
2:13 Now we can earn it piece by piece. Before any of that earning, let's pin the functional requirements, because the map only makes sense against them. Route every user to a healthy nearby region. Serve catalog reads locally, everywhere. Read and write each user's profile in that user's home region, and move that home safely when their life moves. Then the sharp ones. Reserve scarce inventory without overselling, no matter which region the request enters. Create the order and its reservation together, atomically, so a retry can't buy the card without reserving it.
2:44 And accept append-only events locally in any region without coordination, such as a "notification seen" marker. Sure. Those three are the easy promises, though. The clauses people forget are the ones that bite: keep the allowed reads and writes running through a regional outage. Transfer write authority after a failure without ever creating two owners. Surface the conflicts we can't merge automatically. And re-verify the data when a lost region comes home. Notice the list already refuses to treat the data uniformly.
3:13 Catalog, profiles, inventory, events: four different promises hiding inside one word, "database". Which means the design job is really four smaller jobs, and we get to pick the right tool for each. And the non-functional requirements decide how strict each of those four promises is. First, invariant safety, the absolute 1: available inventory never goes below zero, and one reservation ID moves stock at most once. Ever. Then two numbers that anchor every region-loss conversation. Let's say Tokyo goes dark for an hour.
3:45 Recovery point objective: how much acknowledged data we're allowed to lose. Recovery time objective: how long we're allowed to be down. Yep, and for inventory and orders, our recovery point objective is zero; losing an acknowledged reservation isn't on the menu. There's also a stranger requirement, one that sounds like giving up: a region that can't prove it holds the needed votes, or the current ownership of a row, refuses strict writes. Refusing is correct behavior. Inventing a second history is the real outage.
4:17 And the quiet ones do real work too: latency budgets per data class, residency rules about where rows may physically live, and cost, because every extra voting region is real money and, I mean, real milliseconds on every write. Half the design reviews I've sat through were really about this list. And pinning it down is what sets up the fight with the most seductive phrase in distributed systems. So every design review of that list arrives at the same magic words. The loudest architect in the room leans back and says: just make it active-active, every region takes writes, no failover story, fast everywhere. It's seductive because it sounds like a design.
4:57 I'll own this one. I've pitched that exact sentence in a real design review, and when a colleague asked me where two writes to the same row get ordered, I had nothing. I didn't realize until that meeting that I'd been treating a deployment diagram as a consistency claim. That's the problem with the label. Active-passive at least states a deal: one region orders every write, the rest hold copies and wait. Active-active only promises a writable endpoint in every region. On ordering, on conflicts, on partitions: silent, silent, silent.
5:28 So we replace the label with four questions, and they run the rest of this chapter. Where do requests enter? Where does a write to this particular key get ordered? When do we tell the caller "done", before or after other regions have it? And when regions can't talk, which side keeps writing? Ask those four about the opening incident, and it stops being a mystery. Nobody was ordering the writes, and "done" came before replication. Four questions, four bad answers. The rest of the chapter is spent earning the good ones, one data class at a time.
6:03 So let's answer those questions for the easiest data class first, the catalog. Pattern one, the humble 1: a single global writer. Every write routes to one primary region, and asynchronous copies flow out to everyone else. That's Aurora's global database shape, a single primary with read-only secondaries, and AWS's docs put typical replication lag under a second. Typical. Not a contract. Yeah, and "typical" is the word that ends up quoted in a postmortem. For the catalog, though, it's honestly the right answer.
6:33 For example, product names, prices, and images get read constantly in every region, but they're written occasionally, by merchandising tools, from one place. No shopper's checkout is blocked on a description edit crossing an ocean. What you're buying is one unambiguous order for free, ordinary multi-row transactions, and the simplest recovery story on this map. What you're paying is remote write latency for far-away writers, plus a real failover procedure for the day the primary region dies. Even DynamoDB, where every global-table replica can technically accept writes, documents "write to one region" as a supported mode.
7:11 The passive regions still earn rent, local reads plus shorter recovery, and the result is a default you should have to argue your way out of. Now, profiles do make that argument, because profile writes sit on the user's own path. Saving a shipping address or a payment preference shouldn't wait on an ocean crossing. So pattern two moves authority closer without multiplying it: one home region per key. Let's say Priya lives in Tokyo. Her row is homed in Tokyo, her writes commit in Tokyo, and Virginia holds a replica for reads and recovery.
7:45 Which is a real product shape, not a whiteboard fantasy. DynamoDB documents it as "write to your region"; CockroachDB does it row by row with regional tables. And the home itself is just data: a home-region column plus an ownership epoch, a version number for authority itself. The gateways route every request by it. Okay, but people move. Imagine Priya flies to London and hits save. Does that write really cross the planet back to Tokyo? It does. That's the honest cost of this pattern: her writes pay the Tokyo round trip until her home officially changes.
8:21 What you don't do is let London just take the row because she happens to be there this week. Why shouldn't it, though? Moving the data to wherever the user is feels like the obvious kindness. Because "just take it" is how two regions end up believing they own the same row. Rehoming has to be a guarded handover. Guarded how? Freeze writes at the old home, replicate through the cutover point, advance the epoch, and only then does London get to write. Any replica still waving the old epoch gets refused.
8:52 That refusal is called fencing. File it away, because it comes back at the worst possible moment in this chapter. So catalog has one writer, profiles have one home each. Now the graphics card, where the chapter stops being gentle. Inventory carries an invariant: available stock never dips negative, and each reservation ID changes stock at most once. That second clause is what makes retries safe. For example, the same reservation request can arrive twice and still count exactly once. Right, and an invariant is a promise about a sequence, which means someone has to see every change in one order. Two regions independently going one-to-zero isn't a conflict you can merge away afterward. Neither row is wrong.
9:37 The order is wrong. That's worth sitting with for a second. When Virginia and Tokyo both wrote zero, each write was locally flawless. The failure doesn't live in either region. It only exists between them. Which is why this data class can't be "accept locally, sort it out later". There is no later. Sorting it out would mean un-selling a card someone already believes they bought, so this state gets one global order, no matter where its requests come from. Pattern three gives it that order: synchronous, globally ordered writes.
10:09 And here's where it gets really interesting, because I'd rather skip the textbook definition. Two shoppers, one graphics card. What actually happens on the wire? Okay, let's unpack this. Both requests enter locally: Virginia's cell takes reservation R-US, Tokyo's cell takes R-JP. But entering isn't deciding. The inventory range holding this SKU has one current leader, let's say it lives in Virginia right now, plus voting replicas in the other regions. So both reservations get routed to that leader, and Tokyo's request crosses an ocean before anyone says yes.
10:44 It does. That's the toll. The leader takes R-US first and runs it as, basically, a conditional transaction: decrement, but only if available is at least one, and only if this reservation ID hasn't already been used. So it doesn't wait for every region to answer, just enough of them? Right. A majority, not everyone. It ships the write to its voters, and once a majority acknowledges, it commits. Virginia's shopper gets a green checkmark that actually means something. And the Tokyo reservation? R-JP arrives a beat later, runs the same conditional check, reads zero, and fails cleanly.
11:21 "Sold out" in a couple of seconds, instead of a confirmation today and an apology email on Thursday. Oh, I like that trade. An honest no now beats a fake yes later. And this is Spanner's shape, synchronous Paxos replication with commits waiting on a voting quorum, and CockroachDB's, with Raft underneath. Notice neither of us quoted a latency number. The cost is the real network paths, caller to leader, leader to its farthest needed voter. Move the leader or the voters, and the number moves. Now cut the cable under that toll.
11:53 Tokyo can't reach Virginia or Frankfurt. What happens to inventory writes? The side that can still assemble a majority keeps going. Virginia plus Frankfurt hold two of the three votes, so reservations continue there. Tokyo holds one vote and a strong opinion. Strong opinions don't commit anything. Huh. So the vote count is the whole veto. And on a status page that looks like a failure, so let's say it plainly: Tokyo refusing strict writes is the design working. The alternative was Tokyo writing its own separate timeline, and we priced that in the first minute of the chapter.
12:28 Yeah, and Tokyo isn't dead, either. Catalog reads stay local, seen-events keep queueing locally, and only the writes that need global order wait for the world to heal. That gives the outage a shape we chose in advance, instead of one it chose for us. Before we crown the quorum, the tempting alternative deserves a fair hearing, because asynchronous multi-writer is sometimes right, and it's everywhere. Every region acknowledges its own writes immediately, and replication syncs in the background. Which is DynamoDB's eventually consistent global-table mode, and Cosmos DB's default for multi-write accounts.
13:03 So imagine two regions update the same item inside that replication window. Who wins? The database picks. Last-writer-wins: newest timestamp keeps the row, and the losing update is discarded. Hold on. Discarded where, exactly? Nowhere, and that's the detail that stopped me in AWS's own global-tables guide. No metric, no audit event, no conflict record. So a write you acknowledged to a customer just... stops having happened. Quietly, yes. You find out from your own books, or from the customer. Okay, but set the silence aside for a second.
13:37 Convergence itself still worked. Every replica agrees on zero. That's exactly the trap. Look at what it converged to: last-writer-wins picked a zero, and the bytes match everywhere, forever. Equal bytes don't prove correctness. Both reservation rows still exist. Oh, that's nasty. That one line rearranges the whole postmortem. And a smarter merge doesn't escape it: a counter that faithfully keeps both decrements converges to negative one, so the database stops hiding the oversell and starts displaying it.
14:09 So after all that, I don't want the takeaway to be "never". When is the asynchronous pattern actually legitimate? When the merge is the intended meaning, not damage control. For example, "user has seen this notification". Two regions each record it, the merge is a set union, and the union is simply true. Order never mattered, so nothing was lost. Same shape for event logs with globally unique IDs, or preference tags where the union is the actual spec. Yeah, and that's where conflict-free replicated data types earn their name.
14:42 CRDTs: structures built so concurrent updates merge to the same result in any arrival order. Redis's active-active mode ships them, with merge behavior defined per data type. The limits being the fine print people skip. That fine print is my rule of thumb from watching teams get burned: "use a CRDT" is not a design. Which type? Which operations? What does delete mean? A grow-only set has clean answers. A number that must never go below zero doesn't. That's an invariant wearing a counter costume.
15:14 Ha. So strip the costume off and send it back through pattern three's single order. And the shops that run multi-writer happily for years aren't disproving this, by the way. They've already done the classification: everything they let merge really merges, and the wall is the row that doesn't. The data model makes the call, not the enthusiasm in the meeting. Now, pattern three's toll invites a famous objection, so let's just have it. Google put atomic clocks and GPS receivers in its datacenters. If TrueTime can order transactions by time, why does Spanner still push every write through Paxos?
15:50 Because a timestamp doesn't store bytes. TrueTime's actual innovation is honesty: instead of pretending to know what time it is, it returns a bounded interval, "now is somewhere in this window", and the paper measures that window in single-digit milliseconds. Spanner stamps the transaction inside it, waits those few milliseconds of uncertainty out, and only then exposes the commit, so commit order matches real-world order. But the write itself... ...still has to survive machines dying. So it's still Paxos to a majority of voters, and still a coordinator walking any transaction that spans two ranges through its commit. The clock buys trustworthy ordering and clean global reads.
16:32 It doesn't refund a single cross-region round trip. Huh. So the clock's a better witness, not a shortcut. And CockroachDB reaches a similar place without the atomic hardware: hybrid logical clocks, ordinary time plus a logical counter, Raft underneath, a hard cap on clock drift. Same conclusion either way, then. Clocks help you order and read; quorum placement decides what a write costs. Nobody's outrun the speed of light with a better clock. So with all four patterns on the table, put them side by side and our "one database" resolves into a federation behind one API.
17:07 Catalog: one global writer, replicas everywhere. Pattern one. Profiles: one home per key, guarded rehoming. Pattern two. Inventory and its reservations: one ordered range, leader and quorum. Pattern three, the only writes paying the cross-region toll, because they're the only ones you can't apologize your way out of. Seen-events and preference tags: asynchronous, merge defined up front. Pattern four. And the map's quiet hero is the locality catalog, the table of tables. Every key range mapped to its pattern, its home, and its allowed read modes, consulted by the gateways on every request.
17:43 When someone asks, for instance, "is our database strongly consistent?", the honest answer is a lookup, not a slogan. And that's the shift: one system, four promises, each one chosen on purpose instead of inherited from a product label. Now let's break a region and find out if they hold. That four-pattern map leads to the real exam: break Virginia and see whether failover keeps write authority singular. The scary version isn't a clean crash. It's the one-way partition. Virginia's cell looks healthy to Virginia's users; Tokyo's looks just as healthy from Tokyo.
18:17 The severed thing is between them, and each side has to suspect the other is gone. I'll admit my reflex is wrong here. My reflex is Virginia's unreachable, promote Tokyo, keep the business running. Sure. Promote it based on what evidence, though? Tokyo can't distinguish "Virginia is dead" from "Virginia can't hear me". Imagine Tokyo self-promoting while Virginia's leader is still alive and committing: two owners, split-brain, the disaster this whole design exists to prevent. Traffic health is not write authority.
18:50 The load balancer can tell you Virginia looks fine from where it stands; only a quorum can tell you who's allowed to write. Oh, that's the whole thing, isn't it. And the sequence follows straight from it: the rehoming discipline from the profile section, under pressure. Fence first: revoke the old owner's epoch, so its writes bounce even if it's alive and confused. Then promote through a quorum or the failover machinery, reroute traffic, and make the recovering region catch up before it serves anything as current. And when Virginia rejoins, matching bytes isn't the bar.
19:22 It replays the log to the required position, and the invariant checker gets its pass, no negative inventory, no duplicate reservations, before it's a peer again. Routing never granted authority in the first place. The epoch did. And a design that fails over on paper still has to prove it somewhere that isn't production. Start with the number no server graph carries: which region holds each range's leader, and what is every write paying for it? Concrete case: a deploy moves the inventory leader to Frankfurt, and from that moment every Virginia write pays an extra cross-region hop without a single host metric blinking.
19:58 Wait, so a latency regression with every graph calm? Exactly that. Same family: stale reads aging past what the product promised, and a merge policy resolving conflicts that nobody is counting. So the dashboard tracks authority and invariants, not just machines: leader location per range, quorum latency by region pair, stale-read age, epoch rejections, and stock that never reads negative. A broken invariant should page somebody, not wait for the quarterly audit. And rehearse the authority transfer itself, on purpose.
20:29 First, kill the leader region, not the convenient one. Then partition one direction only. Make a region slow instead of dead, because slow is meaner than dead. Route a stale gateway at a row that moved, and watch the fence actually catch it. Until you've seen Tokyo refuse a write for the right reason, you only hope it will. So, back to the exact picture we opened on, and every arrow means something now. Traffic manager and cells: entry. Catalog fanning out from one writer. Priya's row behind its epoch. The ordered range with its leader and voters.
21:03 Seen-events merging in the background, and the failover controller holding the fence. Entry versus authority stopped being decoration about 15 minutes ago. It's the line every failure in this chapter tried to cross. So the next time a review says "make it active-active", here's the whole method in two questions. 1: does this key genuinely need writes from more than one region, or does it need fast local reads and one good home? Most data fails that question, and it's happier for it. And 2: if it truly does, do concurrent operations have a merge that preserves the business meaning? Yes means pattern four, asynchronous with a real, named merge. No means pattern three, one order, a quorum, and physics on the write path.
21:47 There's no secret 5th answer. That's what the label was hiding all along. Run SKU-RED-1 through it one last time: one reservation commits, one shopper gets an honest no in seconds, and no dashboard lies to anyone. Next time, we build a search engine: how a query finds the right 10 results out of a billion documents in well under a second. Thanks for listening to Learning Podcasts.