KnowSys

Designing Amazon's Cart and Checkout

At midnight Prime Day opens, and Leo adds a deal-priced pair of headphones to his cart on his laptop and his phone, then checks out, while millions of people do the same. We'll design the system that keeps his cart and takes his order: from a single database to Dynamo's always-writeable cart, vector clocks and merged carts, DynamoDB's Paxos groups, a checkout workflow with inventory claims, payment authorisation and idempotency keys, an asynchronous order pipeline, and the load shedding, cells and shuffle sharding that keep a peak from becoming an outage.

⏱ 60 min read◆ IntermediateAssumes: chapter 28 (replication) and chapter 29 (partitioning) help, chapter 31 (distributed transactions) helps, the Stripe case study (chapter 52) helps
Start reading

It's a minute past midnight and Prime Day has just opened. Leo has been watching one deal since the evening: a pair of noise-cancelling headphones that normally cost $249, offered as a Lightning Deal at $149 for as long as the stock lasts. A bar under the price says 0% claimed. He clicks Add to Cart on his laptop. Then, because the laptop seems slow and he doesn't trust it, he taps Add to Cart on the same deal on his phone. His phone shows the headphones in his cart, and the bar on the laptop already says 31% claimed. He adds a cable on the laptop, removes it again, opens the cart on his phone, and taps Place your order. A few seconds later an email arrives: Ordered: headphones.

A huge warehouse floor seen from above, with a long conveyor carrying boxes on the left and rows of blue roll cages on the right
Inside BRS2, an Amazon fulfilment centre in Swindon, UK, in 2025. Every order placed at midnight ends up as a box on a sorter like the one on the left. Everything in this chapter happens before an order gets here, in the few seconds between Add to Cart and Order placed.Photo: Auledas, CC BY 4.0, via Wikimedia Commons

Every one of those taps had to work while millions of other people were tapping too. The cart had to accept Leo's writes from two devices at once, possibly reaching different servers, and end up with one sensible cart. There was a fixed number of deal units, so the system had to decide who got one without selling more than it had. And "Place your order" quietly started a chain of work: check the price, set a unit aside, ask Leo's bank whether the card is good, record the order, and later pick, pack and ship it. Any step could fail, time out or run twice.

In this case study we'll design the system behind Leo's cart and checkout the way an engineer would: start with the most obvious design, find exactly where it breaks, and fix it, step by step. We'll keep coming back to one question: when millions of people hit the same deals at the same moment, how does Amazon make sure Leo's cart never refuses him, and his order is taken exactly once, without selling headphones it doesn't have? Along the way we'll go from one database to a ring of storage nodes, vector clocks, Paxos groups, a checkout workflow, conditional writes, idempotency keys and the cells that stop one failure from reaching everyone.

01What we're building, and how big

1.1What it has to do

This case study covers the path from "I want this" to "your order is placed". It comes down to a short list:

  1. Keep a cart for each shopper: add an item, remove it, change a quantity, from any device, and show the same cart everywhere.
  2. Price the cart, including deals whose price depends on time and on stock.
  3. Hand out limited stock for deals like Leo's, without promising more units than exist.
  4. Take the payment: check that the card is good when the order is placed, and charge it when the goods ship.
  5. Create the order exactly once, however many times the button is pressed or the request is retried.
  6. Pass the order on to the warehouses that pick, pack and ship it.

And the qualities it needs while doing that:

  • Writable, always: adding to the cart should work even when servers, disks and network links are failing, because a refused Add to Cart is a lost sale.
  • Exact where money and stock are concerned: one order per Place your order, one charge per order, no more deal units sold than the seller has.
  • Fast: every page is assembled from many services, so each one must answer quickly, even at the worst moment of the year.
  • Contained: a failure in one part, or for one group of customers, mustn't take the whole store down.

Notice that the first and second qualities pull in opposite directions. "Always accept the write" and "never sell what you don't have" can't both hold for the same piece of data during a network failure, and much of this chapter is about deciding which data gets which promise.

1.2How big is it?

Amazon has published two kinds of numbers about this, years apart. Its Dynamo paper, presented at SOSP in 2007, describes the shopping cart service during a busy holiday shopping season: it "served tens of millions requests that resulted in well over 3 million checkouts in a single day". That paper also says Amazon's platform served "tens of millions customers at peak times using tens of thousands of servers located in many data centers around the world".

Newer numbers come from the posts AWS publishes after each Prime Day, describing how Amazon's own retail systems used AWS services during the event. They measure the infrastructure, not the cart alone, but they give the shape of the peak:

Prime DayDynamoDB peakOther figures from the same post
2019 (48 hours)45.4 million requests a second7.11 trillion DynamoDB calls; Aurora processed 148 billion transactions
2021 (66 hours)89.2 million requests a secondFrom the DynamoDB paper: "trillions of API calls" from Alexa, the Amazon.com sites and the fulfilment centres
2023 (2 days)126 million requests a secondCloudFront peaked at over 500 million HTTP requests a minute; SQS at 86 million messages a second
2024 (2 days)146 million requests a secondAurora processed over 376 billion transactions; 733 fault-injection experiments run beforehand
2025 (4 days)151 million requests a secondCloudFront served over 3 trillion HTTP requests in the Prime Day week; SQS peaked at 166 million messages a second
2026 (23–26 June)192 million requests a secondOver 59 trillion DynamoDB requests over the four days; SQS peaked at 213 million messages a second

These are requests from all of Amazon's systems, including Alexa, the warehouses and every page of the store, so most of them have nothing to do with carts. Amazon doesn't publish its cart or checkout traffic.

Your turn: design it before reading on

Suppose 20 million people shop in the first hour of Prime Day, and each one adds, removes or changes cart items ten times. How many cart writes a second is that on average? And what happens if half of them arrive in the first five minutes?

02Version 1: one application, one database

2.1The obvious design

The most obvious design is one web application in front of one relational database. It has a carts table with a row per item in each cart, an inventory table with a stock count per product, and an orders table. Add to Cart inserts a row. Place your order runs a single database transaction: read the cart, check and decrement the stock, insert the order, and commit. If anything fails, the transaction rolls back and nothing has happened.

Version 1: one application and one database
Leo's laptopLeo's phoneWeb applicationall the codeOne databasecarts, stock, ordersCard network
Step 1. Leo clicks Add to Cart on his laptop. The application inserts a row into the carts table.
1 / 4

This is roughly how Amazon.com began. In a 2006 interview with Jim Gray in ACM Queue, Werner Vogels, Amazon's CTO, described it as "a monolithic application, running on a Web server, talking to a database on the back end", an application called Obidos that held "all the business logic, all the display logic, and all the functionality". For years, he said, the scaling work went into making the back-end databases hold more items, customers and orders, "until 2001 when it became clear that the front-end application couldn't scale anymore".

2.2Where it breaks

It breaks in three places, and each one matters for Leo.

The first is size. At the peak, every shopper's every click is a write or a read against the same database, and one database server has a ceiling. You can buy a bigger machine for a while, and you can add read copies, but every write to a cart, a stock count or an order still lands on the one primary.

The second is people. Many teams change one application and one schema. A change to the orders table for one team's feature can break another team's code that reads it, so every change needs coordination, and the whole thing ships together. Vogels described exactly this: the parts that needed to scale couldn't be scaled independently, because they were "shared by many different teams and processes".

The third is availability, and it's the one that hurts at midnight. When the database is slow or down, every part of the store is slow or down with it, including Add to Cart. Leo clicks, waits and sees an error. A database is right to refuse a write it can't make durable, but from Amazon's side, a refused Add to Cart is a sale that may never happen.

2.3Split it into services

Amazon's fix for the first two problems, starting around 2001, was to split the application into services: separate programs, each owned by one team, each with its own data, each reachable only through a published interface over the network. Vogels put the rule plainly in 2006: "No direct database access is allowed from outside the service, and there's no data sharing among the services." He also put the ownership rule plainly: "You build it, you run it."

There's a famous, harsher version of this story. In 2011 Steve Yegge, who had worked at Amazon and then at Google, wrote a long internal post about platforms and accidentally shared it publicly. In it he recalled a mandate from Jeff Bezos, "back around 2002 I think, plus or minus a year", that all teams must expose their data and functionality through service interfaces and communicate only through them, with no direct reads of another team's data store, and that anyone who didn't would be fired. Amazon has never published the memo, so treat Yegge's version as one engineer's recollection. Vogels' 2006 description is the primary source for what the rule was.

So the cart becomes a cart service with its own store, inventory becomes an inventory service, and so on. Each page is then assembled from many of them. In 2006 Vogels said the Amazon.com gateway page called "more than 100 services", and the Dynamo paper says a page request "typically requires the rendering engine to construct its response by sending requests to over 150 services".

Splitting into services solves the people problem and lets each service scale on its own. It doesn't solve the third problem by itself. A cart service still needs somewhere to keep carts, and if that store refuses writes when a server or a network link fails, Leo still sees an error at midnight. So the next step is a store that never refuses him.

03A cart that always accepts writes

3.1Why the cart can't say no

The Dynamo paper opens with the cart. Customers, it says, "should be able to view and add items to their shopping cart even if disks are failing, network routes are flapping, or data centers are being destroyed by tornados." The paper is careful about what that requirement means: the cart service needs a store it "can always write to and read from", and "an 'Add to Cart' operation can never be forgotten or rejected".

Here's why that's harder than it sounds. To survive a server failure, the store keeps each cart on several machines, called replicas. A common way to keep replicas in agreement is to accept a write only after a majority of them have stored it, so that any later read, which also asks a majority, is sure to reach at least one replica with the latest write. Chapter 28 works through why. But a majority rule has a cost: if Leo's phone can reach only one of three replicas, because the other two are down or cut off by a network failure, the write is refused.

For a bank balance, refusing is right. For a cart it's the wrong trade. If the cart briefly has two versions, the worst that happens is that Leo sees a cable he'd removed, and removes it again. If the cart refuses his write, he probably doesn't buy the headphones. Dynamo's authors put the conclusion in one phrase: Dynamo targets "the design space of an 'always writeable' data store", and pushes "the complexity of conflict resolution to the reads", so that a write is never rejected.

3.2Where Leo's cart lives

Dynamo spreads data over many storage nodes by consistent hashing, which chapter 29 covers in detail. In short: hash each key onto a circle of numbers, give each node several positions on the circle, and make each key belong to the first node clockwise from it. Leo's cart has a key, say his customer ID, and the hash of that key lands somewhere on the ring.

Dynamo stores each key on N nodes, the first one clockwise and the next N − 1 distinct nodes after it. This list of nodes for a key is called its preference list. A typical N was 3. Every node knows the whole ring, by gossiping with the others, so any node can work out Leo's preference list without asking anyone.

A ring of eight nodes n1 to n8; two keys, foo and bar, hash to points on the ring and are each stored on the next three nodes clockwise
A consistent hash ring of eight nodes, from Apache Cassandra's documentation of its Dynamo-style design. The key foo hashes between n1 and n2 and is stored on n2, n3 and n4, its preference list with N = 3. Leo's cart key works the same way.Image: Apache Cassandra documentation, Apache License 2.0

When the cart service writes Leo's cart, the request goes to a coordinator, usually the first node in his preference list. It writes the new version locally, sends it to the other N − 1 nodes, and reports success once W nodes in total have stored it. Reads work the same way with R. The paper gives (N, R, W) = (3, 2, 2) as the configuration several instances of Dynamo used, and says a typical service-level agreement was that 99.9% of reads and writes finish within 300 milliseconds.

3.3When a home node is down: sloppy quorums and hints

Now suppose one of the three nodes in Leo's preference list is down at midnight. A strict rule would still work with two of three, but if two are unreachable, a strict quorum write fails. Dynamo doesn't use a strict rule. It writes to "the first N healthy nodes from the preference list", walking further round the ring past nodes it can't reach. A quorum counted this way, partly from stand-ins, is called a sloppy quorum.

A stand-in keeps the replica with a hint attached, a note saying which node it belongs to, and keeps hinted replicas in a separate local store. When it sees the intended node come back, it hands the replica over and deletes its copy. This is hinted handoff. Chapter 28 animates both ideas on a different example; here's the same mechanism on Leo's cart.

Leo's Add to Cart while one of his cart's home nodes is down
DYNAMO RING, N = 3, W = 2put(cart)handoffLeo's phoneCart servicestatelessNode AcoordinatorNode Bhome replica, downNode Chome replicaNode Dstand-in, keeps hint
Step 1. Leo taps Add to Cart. The cart service hashes his customer ID and sends the write to node A, the first node in his preference list.
1 / 5

Dynamo's authors are clear about the cost. A write acknowledged by stand-ins and a read answered by the home nodes needn't overlap, so a read can miss the newest version until the hints are delivered. For the rarer case where a stand-in itself fails before handing its hint over, Dynamo also runs a background repair: replicas compare their data using Merkle trees, trees of hashes where each parent is a hash of its children, so two nodes can find which keys differ by exchanging a few hashes instead of all their data.

Dynamo also places each key's replicas in more than one data centre, so that a whole data centre can fail without losing carts. According to the paper, its client applications "received successful responses (without timing out) for 99.9995% of its requests" over the two years before publication, with no data loss.

?Why not just set W to 1?

You can. With W = 1, a write succeeds as long as any single node in the system can store it. Dynamo's paper notes that this "ensures that a write is accepted as long as a single node in the system has durably written the key", and also that most Amazon services set W higher, because a write stored on one node is one disk failure away from being lost. W trades durability against availability, and (3, 2, 2) was the common middle.

So writes are never refused. That leaves the problem the paper said it had pushed to the reads: Leo's laptop and phone may now have written two different carts to two different sets of nodes, and the next read has to make sense of both.

04Two devices, two carts: vector clocks

4.1Leo's two carts

Go back to midnight. Leo's cart holds the headphones and a cable. Then the network between two of Amazon's data centres hiccups. His phone's request reaches node B on one side, and he adds a charger there. His laptop's request reaches node A on the other side, and he removes the cable there. Both writes are accepted, because the cart is always writeable. Now there are two versions of his cart, and neither one has seen the other's change.

Leo's cart splits into two versions, and comes back together
Leo's laptopNode Aone data centreNode Banother data centreLeo's phoneNext read of the cartcart service{A:1}headphones, cable{A:1}headphones, cable
Step 1. Before midnight both nodes hold the same cart, version {A:1}: headphones and a cable.
1 / 6

Those labels on the versions are what make this work, and they're called vector clocks.

4.2What a vector clock records

zoomAmazonCart serviceDynamoVector clock

A vector clock is a list of (node, counter) pairs attached to every version of an object. When a node coordinates a write, it takes the clock of the version the client says it's updating and adds one to its own counter. So {A:2} means "a version that includes two changes made at node A", and {A:1, B:1} means "one change made at A, then one made at B".

To compare two versions, compare their clocks counter by counter. Dynamo's rule is: "If the counters on the first object's clock are less-than-or-equal to all of the nodes in the second clock, then the first is an ancestor of the second and can be forgotten. Otherwise, the two changes are considered to be in conflict." Try it on Leo's two carts. {A:2} has a bigger A counter than {A:1, B:1}, but {A:1, B:1} has a bigger B counter. Neither is less than or equal to the other on every node, so neither version descends from the other.

Three timelines A, B and C with messages between them; every event is labelled with a vector clock such as A:2 B:4 C:1, and shaded regions mark one event's causes and effects
Vector clocks on three processes exchanging messages. Each event's clock lists one counter per process; the shaded region on the left holds the events that came before one event on B, the region on the right those that came after it, and the events in neither region are concurrent with it. Leo's two carts are a pair like that.Image: Vector_Clock.svg, CC BY-SA 3.0, via Wikimedia Commons
Predict before you read on

Leo's laptop wrote version {A:2} and his phone wrote version {A:1, B:1}. What does Dynamo do with them?

That's why a Dynamo get can return more than one object. Its interface is get(key), which returns "a single object or a list of objects with conflicting versions along with a context", and put(key, context, object). A context carries the clock information. When the client writes back a merged cart with the context from a read that saw both versions, the coordinator builds the new clock from both, adds one to its own counter, and the new version, {A:3, B:1} in Leo's case, descends from both siblings, which can then be dropped.

Clocks that keep growing are a cost. A clock gains an entry for each node that coordinates a write to the object, and during failures that can be many nodes. Dynamo stored a timestamp with each pair and, once a clock reached a threshold, "say 10", dropped the oldest pair. Its authors admit this could make the system mistake an old version for a conflicting one, and say the problem "has not surfaced in production". Chapter 26 covers vector clocks in general, including this truncation and the dotted version vectors that Riak later adopted to avoid clock growth.

4.3Merging two carts

Dynamo detects the conflict; it doesn't resolve it. Dynamo calls the two kinds of resolution syntactic and semantic. When one clock descends from the other, the system drops the older version itself. When the clocks are concurrent, the application has to merge, because only it knows what the data means. For the cart, the merge was to take the union: every item in either sibling goes into the merged cart.

A union has a consequence that the paper states in one sentence: "Using this reconciliation mechanism, an 'add to cart' operation is never lost. However, deleted items can resurface." Look at Leo's carts again. On the laptop's side the cable is gone, because he removed it. On the phone's side it's still there, because it never heard about the removal. A union can't tell "removed" from "never had", so the cable comes back.

This program builds Leo's two siblings with vector clocks, compares them, merges them by union the way the 2007 cart did, and then tries a better merge. A better merge tags every item with the clock entry of the write that added it, written as (node, counter). When an item is in one sibling but missing from the other, it checks whether the other sibling's clock has seen that add. If it has, the item is missing because it was removed, and the merge drops it. If it hasn't, the other sibling never heard of it, and the merge keeps it.

Leo's cart on two devices: vector clocks, a union merge, and a merge that remembers removals
python
Python
# Leo's cart on two devices, versioned with vector clocks
 
def descends(a, b):
    """True if clock a has seen everything clock b has."""
    return all(a.get(node, 0) >= n for node, n in b.items())
 
def compare(a, b):
    if a == b:
        return "same"
    if descends(a, b):
        return "a is newer"
    if descends(b, a):
        return "b is newer"
    return "concurrent"
 
def bump(clock, node):
    c = dict(clock)
    c[node] = c.get(node, 0) + 1
    return c
 
def merge_clocks(a, b):
    return {n: max(a.get(n, 0), b.get(n, 0)) for n in sorted(a.keys() | b.keys())}
 
# 1. Leo's cart, written through node A
v1 = ({"A": 1}, {"headphones", "cable"})
 
# 2. The network splits. His phone reaches node B and adds a charger.
phone = (bump(v1[0], "B"), v1[1] | {"charger"})
 
# 3. His laptop reaches node A and removes the cable.
laptop = (bump(v1[0], "A"), v1[1] - {"cable"})
 
print("phone :", phone[0], sorted(phone[1]))
print("laptop:", laptop[0], sorted(laptop[1]))
print("compare:", compare(phone[0], laptop[0]))
 
# 4. The next read gets both siblings. Merge by union, as the 2007 cart did.
union = (bump(merge_clocks(phone[0], laptop[0]), "A"), phone[1] | laptop[1])
print("union merge:", union[0], sorted(union[1]))
 
# 5. Better: tag every add with the clock entry that made it, and keep an
#    item only if both siblings have it, or the other sibling never saw it.
def tagged_add(cart, item, clock, node):
    clock = bump(clock, node)
    return clock, cart | {(item, node, clock[node])}
 
def tagged_merge(a, b):
    (ca, ia), (cb, ib) = a, b
    seen_by = lambda clock, tag: clock.get(tag[1], 0) >= tag[2]
    keep = (ia & ib) | {t for t in ia - ib if not seen_by(cb, t)} \
                     | {t for t in ib - ia if not seen_by(ca, t)}
    return merge_clocks(ca, cb), keep
 
c, cart = tagged_add(set(), "headphones", {}, "A")
c, cart = tagged_add(cart, "cable", c, "A")
start = (c, cart)                                     # clock {A: 2}
phone = tagged_add(start[1], "charger", start[0], "B")
laptop = (bump(start[0], "A"), {t for t in start[1] if t[0] != "cable"})
merged = tagged_merge(phone, laptop)
print("tagged merge:", merged[0], sorted(t[0] for t in merged[1]))
output
C++
phone : {'A': 1, 'B': 1} ['cable', 'charger', 'headphones']
laptop: {'A': 2} ['headphones']
compare: concurrent
union merge: {'A': 3, 'B': 1} ['cable', 'charger', 'headphones']
tagged merge: {'A': 3, 'B': 1} ['charger', 'headphones']

Its first three lines are the situation from the diagram: two clocks, neither descending from the other, so compare says concurrent. On the fourth line, the union merge keeps the charger, which is right, and brings back the cable, which is the resurfacing the paper warned about. Last comes the tagged merge. Leo's cable carries the tag (A, 2): it was added by A's second write. The laptop's sibling has clock {A:3}, so it has seen that add, and the cable's absence there means Leo removed it. As for the charger, its tag is (B, 1), and the laptop's clock has no B entry at all, so the laptop never knew about the charger, and it stays. (The tagged version starts from clock {A:2} because adding the headphones and the cable were two separate writes at A.)

Tagging adds is the idea behind the observed-remove set, one of the data types called CRDTs (conflict-free replicated data types), whose merges are designed so that every replica ends up with the same answer whatever order it sees the changes in. Riak, an open-source database modelled on Dynamo, shipped CRDT sets and maps in its 2.0 release in 2014. Dynamo's 2007 cart used the plain union.

Decision

Two siblings of Leo's cart arrive at a read. How should they become one?

Last write wins
Keep the version with the latest timestamp; drop the other.
  • No merge code
  • Never more than one version to show
  • Silently loses one of Leo's changes
  • Trusts clocks on different machines to agree
chosen
Union, in the application
Return both siblings; the cart service takes every item from either.
  • An Add to Cart is never lost
  • Simple to write
  • Removed items can come back
  • Every reader must handle siblings
A CRDT with tagged adds
Tag each add; merge keeps an item unless the other side saw and removed it.
  • Adds kept, removals respected
  • Same answer on every replica
  • More metadata per item
  • Harder to get right; came later

Dynamo offered both of the first two: the paper says the session-state service used last-write-wins, while the cart used "business logic specific reconciliation" by merging. For a cart, a lost Add to Cart costs a sale, and a resurfaced item costs the shopper one click on the checkout page, where they review the cart anyway. That asymmetry made the union the right call in 2007.

4.4How often does this happen?

All this machinery is for a rare event. Dynamo's authors measured the versions returned to the shopping cart service over 24 hours: "99.94% of requests saw exactly one version; 0.00057% of requests saw 2 versions; 0.00047% of requests saw 3 versions and 0.00009% of requests saw 4 versions."

More surprising was the cause. You'd probably expect failures to create most siblings, but the paper reports that the increase in siblings came from "the increase in number of concurrent writers", which was "usually triggered by busy robots (automated client programs) and rarely by humans". Leo on two devices is the human case, and it seems to be rare. A script hammering one cart from many connections is the common one.

05From Dynamo to DynamoDB

5.1Dynamo was hard to run

Dynamo worked. It's also, by Amazon's own later account, a system few teams wanted to operate. In his January 2012 post announcing DynamoDB, Vogels wrote that Dynamo met the reliability, performance and scalability needs of its users but "did nothing to reduce the operational complexity of running large database systems". A 2022 paper on DynamoDB adds the detail: Dynamo "was a single-tenant system and teams were responsible for managing their own Dynamo installations", and that complexity "became a barrier to adoption". Engineers instead picked Amazon's managed services, S3 and SimpleDB, even when Dynamo fitted their problem better. "Ultimately," Vogels wrote, "developers wanted a service."

SimpleDB had its own lesson for consistency. Vogels' 2012 post lists among its problems that SimpleDB's consistency windows were "up to a second in duration", which developers used to traditional databases found hard to adapt to. So when AWS launched DynamoDB as a public service in January 2012, it took Dynamo's incremental scaling and hash partitioning and changed most of the rest. That paper says DynamoDB "shared most of the name of the previous Dynamo system but little of its architecture."

5.2One leader per partition

Its biggest change is where writes go, and the DynamoDB paper, presented at USENIX ATC in July 2022, describes it. A table is split into partitions, each holding a contiguous range of keys. Each partition has three replicas, in three different Availability Zones, which are separate data centres or groups of them within one AWS region. Those three replicas form a replication group that uses Multi-Paxos, a consensus protocol covered in chapter 27, to elect one of them leader.

Only the leader accepts writes and strongly consistent reads. For each write, it creates a record in its write-ahead log and sends the record to the other replicas, and the write is acknowledged "once a quorum of peers persists the log record", which the paper spells out as "two out of the three replicas from different AZs". Any replica can answer an eventually consistent read. A leader keeps its job by renewing a lease, a time-limited claim to be leader, and if the others decide it has failed and elect a new one, the new leader "won't serve any writes or consistent reads until the previous leader's lease expires", so two leaders never accept writes at once.

A sequence diagram: a client asks a proposer, which sends Prepare to three acceptors, collects Promises, sends Accept, and the acceptors send Accepted to the proposer and to learners
One round of Paxos with no failures. A proposer asks three acceptors to promise (phase 1), then asks them to accept a value (phase 2); a value accepted by a majority is chosen. Multi-Paxos keeps one leader so that phase 1 runs once per leadership, and each write needs only phase 2: one round trip to a majority.Image: Iamfromspace, CC BY-SA 4.0, via Wikimedia Commons
Leo's cart write in DynamoDB
ONE PARTITION'S PAXOS GROUPlog recordlog recordCart serviceRequest routerauthenticates, routesMetadata servicekey range → groupLeader replicaAZ 1ReplicaAZ 2ReplicaAZ 3
Step 1. The cart service writes Leo's cart item. The request router looks up which partition holds his key, and which replica leads it, in the metadata service (and caches the answer).
1 / 5

Leaders are what make vector clocks unnecessary. Every write to one key goes through one leader, one after another, so there's always a single latest version, and no siblings to merge. If the leader's zone fails, the group elects another, and the paper describes adding a log replica, which stores only recent log records, to keep a quorum available while a full replica is rebuilt. It reports running "millions of Paxos groups in a Region".

Decision

How should replicas of a partition agree on writes?

Leaderless, sloppy quorums (Dynamo, 2007)
Any node coordinates; versions tracked by vector clocks; siblings merged by the application.
  • Writes accepted even when most home replicas are unreachable
  • No election pause
  • Applications must handle siblings
  • Reads can miss recent writes
  • Hard for teams to reason about
chosen
A Paxos leader per partition (DynamoDB, 2012 on)
Three replicas in three zones; the leader orders writes; two of three must store each one.
  • One version per item, no merge code
  • Strongly consistent reads available
  • Conditional writes and transactions become possible
  • Writes pause briefly during a leader election
  • A partition cut off from its quorum can't take writes

DynamoDB gave up Dynamo's promise to accept a write on any reachable node, and kept a narrower one: a partition accepts writes as long as two of its three Availability Zones are healthy. In return every item has one current version, which is what developers had asked for. DynamoDB's published targets are 99.99% availability for tables in one region and 99.999% for global tables, which are copied across regions.

5.3What this means for Leo's cart

Amazon hasn't published where today's cart is stored, or how. What it has published is that DynamoDB "powers multiple high-traffic Amazon properties and systems including Alexa, the Amazon.com sites, and all Amazon fulfillment centers", and that in October 2019 its consumer business turned off its last Oracle database, after moving about 75 petabytes of data from nearly 7,500 Oracle databases to DynamoDB, Aurora, RDS and Redshift.

So let's reason about how you'd keep Leo's cart in a store with one version per item. Storing the whole cart as one item would make his two devices' writes overwrite each other: the phone reads the cart, the laptop reads the cart, both write back, and one write wipes out the other's change. There are two standard fixes.

  • One row per cart line. Make the partition key Leo's customer ID and the sort key the product. Adding a charger writes the charger's row; removing the cable deletes the cable's row. The two devices touch different rows, so neither can erase the other's change, and the removal stays removed.
  • A version number and a conditional write. Keep a version attribute on the cart and write with a condition, "only if the version is still 7". If another device got there first, the write fails, and the cart service re-reads and retries. That's optimistic concurrency, and DynamoDB's conditional writes do the check on the leader, so it's exact.

Either way, the problem Dynamo solved with siblings is now solved by ordering writes at the leader, and the cart service's merge code disappears. What's been given up is the guarantee that a write succeeds on any reachable node; what's kept is that it succeeds whenever two of three zones are up. Global tables, which copy a table between regions and let each region accept writes, bring the concurrency question back across regions, and by default settle it with last write wins (chapter 28).

A leader also makes something possible that Dynamo couldn't do: a write that succeeds only if a condition holds, checked and applied atomically. That's exactly what Leo's checkout needs next.

06Checkout is a workflow

6.1Four services, no shared transaction

Leo taps Place your order. In version 1 this was one database transaction. Now the work belongs to several services, each with its own data: pricing decides what he pays, inventory decides whether a unit is his, payments asks his bank whether the card is good, and orders records the order. Section 2's rule says no service may reach into another's database, so there's no single database in which to run one transaction across all four.

A textbook way to make several databases commit together is two-phase commit: a coordinator asks every participant to prepare, and only when all say yes does it tell them all to commit. Chapter 31 walks through it and its main weakness: a participant that has prepared must wait, holding its locks, until it hears the outcome, so a coordinator that crashes at the wrong moment leaves everyone stuck. And one of Leo's participants can't take part at all. His bank doesn't offer to "prepare" a payment and wait for Amazon's say-so, so whatever Amazon asks the card network has to be handled on its own terms.

So checkout becomes a saga: a sequence of steps, each a local transaction in one service, where each step that can be undone has a compensation, an action that undoes it if a later step fails. Chapter 31 covers sagas in general. For Leo's order, the steps and their compensations are:

StepServiceIf a later step fails
1. Price the cartPricingNothing to undo: it only reads
2. Claim a unit of each itemInventoryRelease the claim
3. Authorise the paymentPaymentsVoid the authorisation, so the bank releases the hold
4. Create the orderOrdersNot undone: from here on, problems are handled by cancelling the order
A UML activity diagram: book hotel, book rental car, book flight; if any step fails, the completed bookings are cancelled in reverse order
The saga pattern drawn for a trip booking: each step has a compensating step, and a failure runs the compensations for the completed steps in reverse order. Leo's checkout has the same shape, with claim, authorise and create in place of hotel, car and flight.Image: UlrichAAB, CC BY-SA 4.0, via Wikimedia Commons
Leo's checkout as a saga, when the card is declined
CheckoutPricingInventoryPaymentsOrdersprice cartclaim 1 unitauthorise $161.42declinedrelease claimshow error
Step 1. 1. Pricing returns $149 for the headphones, valid because Leo's claim on a deal unit is still live, plus tax and shipping.
1 / 6

Amazon hasn't published how its own checkout orchestrates these steps, or what it calls them. What follows is how you'd order them, and why.

6.2The order of the steps

That order isn't arbitrary. Steps go from cheap and easily undone to expensive and hard to undo. Pricing only reads. Releasing an inventory claim is a single write in Amazon's own systems. Voiding an authorisation means another round trip to the card network, and a shopper whose bank has already put a hold on their money may notice it for days. Creating the order sends a confirmation email to Leo, and an email can't be unsent.

Putting the irreversible step last means most failures happen before anything is visible outside Amazon. If inventory says no, nothing has been asked of the bank. If the bank says no, nothing has been promised to Leo. The point of no return, creating the order, comes only when everything before it has succeeded.

?Why not take the payment first, to be sure of the money?

Because then every out-of-stock deal would be a payment to undo. At midnight on Prime Day, many more people try to claim the headphones than there are units, so claiming the unit first means most of the losers never touch the payment system at all. So the cheap check that usually fails goes first; the expensive step that usually succeeds goes after it.

6.3Which price?

Step 1 looks the simplest and has a trap in it. Leo saw $149 when he added the headphones at midnight. What should he pay if he checks out at 12:20, after the deal has ended?

Amazon's cart answers part of this itself: when the price of an item changes after you've added it, the cart shows a message saying so, with the old and new prices. A cart is a list of things Leo wants, and the prices in it are reminders. The price is worked out again at checkout, from the pricing service, and the price Leo agrees to is the one on the final order-review page, which the order then records.

A Lightning Deal adds a time limit. Amazon's help pages say that once a Lightning Deal is in your cart, "the item will be available in your cart at the Lightning Deal price for 15 minutes". So the deal price belongs to Leo's claim on a unit, which expires, as much as to the product. Pricing has to ask whether Leo's claim is still live before it quotes $149.

Pricing mistakes in a sale are expensive in both directions. Flipkart, India's largest online store, ran its first Big Billion Day sale on 6 October 2014. A day later its founders, Sachin and Binny Bansal, wrote to customers that while preparing the deals, "the pricing of several products got changed to their non-discounted rates for a few hours". An opposite mistake, a deal price that leaks onto items that shouldn't have it, sells stock at a loss. Either way, the price is data that must be changed carefully and checked before the sale begins, and the price used at checkout must be the one the shopper confirmed.

07The deal: 1,000 units and a million taps

7.1When does a unit become Leo's?

Step 2 is where Prime Day is hardest. Say the seller has 1,000 pairs of headphones at the deal price. At midnight, many times that number of people tap Add to Cart within a second or two. The system must hand out at most 1,000 units, and should hand out all 1,000, and should tell everyone else quickly that they've missed out.

First, when does a unit become Leo's? There are three reasonable answers:

  • When the order is placed. Adding to the cart claims nothing; checkout decrements the stock. That's how ordinary items work: an item sitting in a cart can still sell out before you check out. For a deal it's miserable: thousands of people put the headphones in their carts, fill in their details, and most find at the last step that they're gone.
  • When it's added to the cart, for a while. Adding claims a unit and starts a timer. If the shopper checks out in time, the claim becomes an order; if not, the unit goes back.
  • Never exactly. Take more orders than units and cancel the extras later, apologising. Some sellers do this for preorders. For a deal advertised as limited, it breaks the promise the deal makes.

Amazon's Lightning Deals take the second answer, and its help pages describe it: a status bar shows the percentage of the deal already claimed, and units sitting in someone's cart count, since the Join Waitlist button appears "when all Lightning Deals are bought or are in a customer cart"; the deal is "available one per customer"; and once it's in your cart you have 15 minutes to check out. When a claim expires, the unit goes to the next person on the waitlist, who gets an alert and their own time limit.

This is a reservation: a hold on stock that expires unless it's turned into an order. Ticketmaster does something similar for seats, and chapter 67 builds its seat holds in detail. What differs is what's being held. A concert seat is one particular thing, row F seat 7, so Ticketmaster holds rows. A deal unit is one of 1,000 identical things, so Leo doesn't need a particular unit, only a count that goes down by one. That difference turns out to matter a great deal when a million people arrive at once.

7.2Counting units without overselling

An obvious way to hand out units is the one every web developer writes first: read the stock count, and if it's above zero, decrement it. At Prime Day scale there's a second temptation, which is to read the count from a cache, because the deal page is being loaded millions of times and the "% claimed" bar needs the count anyway.

A safe version makes the check and the decrement one atomic step. In DynamoDB that's a conditional write, an update with a condition, like this:

C++
UpdateItem  key: deal#headphones
            SET remaining = remaining - 1
            CONDITION remaining > 0

DynamoDB's partition leader checks the condition and applies the update together, so two claims can never both take the last unit. But it creates a new problem. Every claim for this deal is a write to one item, so every claim goes to one partition's leader. Chapter 29 gives DynamoDB's documented limit for one partition: 1,000 writes a second. A million taps a second against one item is a hot key, and no amount of extra hardware elsewhere helps it.

A usual fix is to split the count. Put the 1,000 units into ten counters of 100, on ten different keys, which DynamoDB will place on different partitions once they're busy. Each claim picks a counter at random and does a conditional decrement on it. If that counter is empty, try another. Ten counters take ten times the writes, and none of them can go below zero.

This program compares three designs for a crowd of 20,000 claims arriving in the first two seconds. Design A checks a cached count that refreshes every 250 milliseconds and then decrements. Design B does a conditional decrement on one counter whose partition applies 1,000 writes a second, one after another, and a shopper whose request would wait more than a second gives up. Design C splits the units across ten such counters, and a claim that lands on an empty counter tries one other.

Claiming 1,000 deal units: a cached count, one counter, ten counters
python
Python
import random
random.seed(7)
 
UNITS = 1_000                    # headphones at the deal price
CLAIMS = 20_000                  # people who tap "Add to cart" in the first 2 s
arrivals = sorted(random.uniform(0, 2.0) for _ in range(CLAIMS))
 
# A. Check a cached stock count, then decrement. The cache refreshes every 250 ms.
stock, granted, cached, next_refresh = UNITS, 0, UNITS, 0.0
for t in arrivals:
    if t >= next_refresh:
        cached, next_refresh = stock, next_refresh + 0.25
    if cached > 0:                       # looks available, so take one
        stock -= 1
        granted += 1
print(f"A cached count:  granted {granted:,}, oversold {granted - UNITS:,}")
 
# B and C. A conditional write per claim: "decrement if stock > 0".
# Each counter lives on one partition that applies 1,000 writes a second,
# one after another. A shopper whose request would wait over 1 s gives up.
def run(counters, rate=1_000, timeout=1.0):
    free_at = [0.0] * len(counters)      # when each partition is next idle
    sold, waits, gave_up_at, sold_out_at = 0, [], [], None
    for t in arrivals:
        i = random.randrange(len(counters))
        if counters[i] == 0:             # that slot is empty: try one other
            i = random.randrange(len(counters))
        start = max(t, free_at[i])
        if start - t > timeout:
            gave_up_at.append(t)
            continue
        free_at[i] = start + 1 / rate
        if counters[i] > 0:
            counters[i] -= 1
            sold += 1
            waits.append(start - t)
            if sum(counters) == 0:
                sold_out_at = start
    early = sum(1 for t in gave_up_at if t < sold_out_at)
    return (f"granted {sold:,}, oversold 0, sold out at {sold_out_at:.2f} s, "
            f"slowest winner waited {max(waits):.2f} s, "
            f"gave up before the sell-out: {early:,}")
 
print("B one counter:  ", run([UNITS]))
print("C ten counters: ", run([UNITS // 10] * 10))
output
C++
A cached count:  granted 2,502, oversold 1,502
B one counter:   granted 1,000, oversold 0, sold out at 1.00 s, slowest winner waited 0.90 s, gave up before the sell-out: 8,089
C ten counters:  granted 1,000, oversold 0, sold out at 0.11 s, slowest winner waited 0.01 s, gave up before the sell-out: 0

Design A sold 2,502 units of a deal that had 1,000, roughly two and a half times its stock. Every claim in the first 250 milliseconds saw "1,000 left" in the cache and went ahead, and so did the claims after each refresh until the cache finally caught up with a negative count. That's oversell, and at Prime Day arrival rates even a quarter-second of staleness is enough to more than double the units sold.

Design B never oversells, because the condition is checked where the count lives. But the single partition takes a full second to work through the first 1,000 claims, the last winner waits 0.9 seconds for an answer, and 8,089 people give up waiting before the deal has even sold out. Those people saw a spinner or an error, not "sold out", and in real life many of them tap again, which adds to the queue they're stuck in.

Design C is also exact, and sells out in roughly a tenth of a second, so almost everyone gets a fast, clear answer. It costs a little complexity: the "% claimed" bar has to add up ten counters, and near the end a claim may land on an empty counter and need a second try, which the program allows for. Real systems probably also put a limit in front of the counters, such as a waitlist or a rate limit, so that the millionth tap is turned away before it reaches storage at all.

7.3Claims that expire

A claim at add-to-cart lasts 15 minutes, so something has to give units back when claims run out. Store each claim as a record with Leo's customer ID, the deal, and an expiry time, written in the same transaction as the decrement. DynamoDB has offered transactions across items since 2018, described in a USENIX ATC paper in 2023, so the decrement and the claim record can commit together or not at all.

Expiry can be lazy, the same way chapter 67 handles seat holds: any check of Leo's claim treats it as gone once its time has passed. A background sweeper then finds expired claims, increments the counter by one for each, and offers the unit to the waitlist. This sweeper has to be careful about one race: Leo might be checking out at the very second his claim expires. So the checkout's confirmation of the claim and the sweeper's release must both be conditional writes on the claim record ("only if it's still live and still Leo's"), and whichever runs first wins.

Flipkart's first Big Billion Day is a reminder of what happens when this goes wrong. In their letter the next day, its founders wrote that stock for many products ran out "within a few minutes (and in some cases, seconds)" and that this "led to some instances of an order getting over-booked for a product that was sold out just a few seconds ago". That's design A from the program above, in production.

Decision

When should a deal unit become Leo's?

At order time
The cart claims nothing; checkout decrements the stock.
  • No claims to expire
  • Stock never tied up in abandoned carts
  • Most shoppers find out they lost at the last step
  • Checkout takes the whole midnight burst
chosen
At add-to-cart, with a time limit
A conditional decrement plus a claim record that expires after 15 minutes; a waitlist takes returned units.
  • Shoppers know at once whether they got one
  • Units return if they don't buy
  • Units sit in carts unbought for up to 15 minutes
  • Needs an expiry sweeper and a waitlist
Oversell and cancel
Accept every order; cancel those beyond the stock.
  • Simplest; no contention at all
  • Breaks the promise of a limited deal
  • Angry customers, as Flipkart found in 2014

Amazon's Lightning Deal rules describe the second option: units in carts count as claimed, the claim lasts 15 minutes, and returned units go to a waitlist. The cost, stock tied up in carts that never check out, is bounded by the time limit, and the one-per-customer rule stops a single account from claiming a pile of units.

08Taking the payment, once

8.1Authorise now, charge later

Step 3 of the saga asks Leo's bank about $161.42. It doesn't take the money. Amazon's help page on authorisations describes the two halves: "When you place an order, Amazon contacts the issuing bank to confirm the validity of the payment method," and the bank reserves the funds; then, when the order ships, "Amazon notifies your bank that the order is concluded and the amount can be charged." If the order is cancelled, Amazon tells the bank the authorisation is no longer needed.

A small card reader with a card inserted, asking for a PIN, next to a phone showing a payment of 13.37 euros in progress
A card payment in progress on a small chip-and-PIN reader. Whether the card is in a reader or typed into a checkout page, the first message to the bank is the same: an authorisation request asking whether this card can pay this amount, which reserves the money without moving it.Photo: Havarhen, CC BY-SA 3.0, via Wikimedia Commons

That split, authorisation now and capture later, is how card payments work in general, and the Stripe case study (chapter 52) explains the card network's side of it. For a store like Amazon it's a necessity. Leo's order may ship from two warehouses on two days, and each shipment is charged when it leaves; Amazon's help page on declined payments mentions the case where "part of your order was charged and dispatched". An item can turn out to be damaged or missing in the warehouse, and then it isn't charged. Capture at shipping means Leo pays only for what's on its way.

It also means the order can be placed without waiting for the payment to settle, and some payments fail after the order exists. An authorisation may be declined on a later retry, or expire before a slow item ships. Amazon's help pages describe what happens then: the order shows that the payment needs attention, the shopper can retry or choose another card, and "if your payment is declined again, we'll cancel the order".

8.2When the tap is retried

Here's the failure that makes payments hard. Leo taps Place your order. His phone's request reaches checkout, which claims, authorises and creates the order. Its response is on its way back when his phone switches from Wi-Fi to mobile data, and the response is lost. His phone shows an error. He taps again.

Without protection, the second tap is a new checkout: a second authorisation and a second order. The fix is an idempotency key: a unique ID for this one attempt to place this one order, created when the order-review page is rendered and sent with every tap and every retry of it. Checkout stores the key with the order. When a request arrives with a key it has already completed, it doesn't run the saga again; it returns the result of the first run.

Chapter 52 builds idempotency keys in detail for Stripe. Amazon's own version of the rules is in the Amazon Builders' Library, in an article by Malcolm Featonby on making retries safe. Amazon's preferred approach, it says, is "to incorporate a unique caller-provided client request identifier into our API contract". Two rules from it matter for Leo:

  • The same key with different parameters is an error. If a request reuses Leo's key but asks for a different cart, it's safest to assume the caller meant something different, and refuse it, instead of returning the old order.
  • A repeat gets a semantically equivalent response. The retry should look like the first success, "Order placed", with the order's current state, not an error such as "order already exists", which would leave Leo's phone unsure whether its own tap created the order.

AWS's own APIs work this way. EC2's RunInstances takes a ClientToken, and DynamoDB's TransactWriteItems takes a ClientRequestToken, which, its documentation says, makes repeat calls with the same token idempotent for 10 minutes after the first one completes.

?Why generate the key on the order-review page, and not when the tap arrives?

Because the server can't tell a retry from a second order by looking at it. Two identical requests to buy headphones might be Leo's phone retrying, or Leo deliberately buying a second pair. Only the client knows which, so the client must name the attempt. A key made when the page is rendered is the same for every tap on that page, and a fresh page, which is what Leo would load to buy again on purpose, gets a fresh key.

An idempotency key also protects each step inside the saga too. Leo's payment authorisation is sent to the payment processor with a key derived from the order's key, so if checkout crashes after asking the bank and retries the step, the bank sees a repeat, not a second authorisation.

09After Place your order: the asynchronous pipeline

9.1What has to happen before Leo sees 'Order placed'

Once the authorisation is approved, step 4 creates the order. Then there's a long list of things that must happen to it: a confirmation email, fraud checks, choosing which fulfilment centre ships it, telling that centre to pick and pack it, booking the delivery, capturing the payment when it ships, and updating Leo's order history. If checkout did all of that while Leo watched a spinner, his order would take minutes, and any slow downstream system would hold up every checkout.

So only a little is done synchronously. Checkout writes the order record, and in the same local transaction writes an event saying "order 114-… created" to an outbox, a table in the same database whose rows are later published to a queue. Chapter 31 explains why this is needed: writing the order and then sending a message are two separate actions, and a crash between them would leave an order nobody downstream hears about, or a message about an order that doesn't exist. With the outbox, both or neither are written. Then Leo sees "Order placed". Everything else reads the event from the queue and does its work later, at its own pace.

Leo's order, from Place your order to the fulfilment centre
one transactionpublishshippedLeo's phoneCheckoutruns the sagaOrders storeorder + outboxOrder eventsqueueNotificationsconfirmation emailFraud checksFulfilment planningwhich centre ships itPaymentscapture on shipping
Step 1. Leo taps Place your order. Checkout prices, confirms his claim and gets the authorisation.
1 / 5

Every consumer here has to be idempotent too, because a queue delivers each event at least once: if a consumer crashes after doing its work but before recording that it did, it will see the same event again. Sending Leo two confirmation emails is a bit embarrassing; picking two pairs of headphones is a real cost. So each consumer records the event IDs it has handled, and skips repeats.

Amazon hasn't published the internals of its order pipeline. Its terms of sale in the UK, and for its Global Store, show the same split from the customer's side: the order confirmation email only acknowledges that the order was received, and the contract of sale is formed when the item is dispatched, with a separate dispatch confirmation for each package.

Yellow-framed machines over a conveyor, scanning and labelling cardboard boxes, with long conveyor lines behind
The SLAM line at the same Swindon fulfilment centre, where each packed box is scanned, labelled and manifested before it ships. This is the far end of the asynchronous pipeline: the order event Leo's checkout wrote becomes, hours later, a label on a box.Photo: Auledas, CC BY 4.0, via Wikimedia Commons

9.2Why the queue helps at midnight

A queue does a second job on Prime Day. At 00:01, orders arrive far faster than warehouses can pick, and faster than the email and fraud systems would like to work. A queue between them absorbs the difference: orders pile up for a while and drain as the warehouses catch up. Checkout only has to be as fast as the order store, not as fast as the slowest thing downstream.

Prime Day posts give a sense of how much Amazon leans on queues. In Prime Day 2025, Amazon SQS, AWS's queue service, peaked at 166 million messages a second for Amazon's systems; in 2026, 213 million. Behind the queues, the warehouses are a distributed system of their own. The 2025 post says AWS Outposts racks in fulfilment centres sent "more than 524 million commands" to over 7,000 robots during the event, peaking at 8 million commands an hour.

A low, round-edged blue robot with two wheels and a black turntable on top, on display on a table
An Amazon Robotics drive unit on display. In the warehouse, units like this slide under shelving pods, lift them on the turntable and carry them to the people who pick items, each move a command from software that reads the orders.Photo: Auledas, CC BY 4.0, via Wikimedia Commons

That deals with the work after the order. What's left is the moment that decides whether any of it happens: the first minutes of Prime Day, when every service in this chapter sees its worst load of the year at once.

10The first minute of Prime Day

10.1Getting ready for a known peak

Prime Day has one advantage over most traffic spikes: everyone knows when it starts. That changes the job. Autoscaling, adding servers when load rises, reacts in minutes, and the first minute of Prime Day doesn't wait. So capacity is added before the event, sized from forecasts and load tests, and the systems are tested by breaking them on purpose.

AWS's posts show the second habit growing quickly. Amazon ran 733 experiments with AWS Fault Injection Service, a tool that injects failures such as stopped instances or slow network links into a running system, before Prime Day 2024; over 6,800 before Prime Day 2025; and over 44,000 before Prime Day 2026. The DynamoDB paper describes the same practice inside one service: "game days (chaos and load tests)" are among the mechanisms it lists for changing the system safely.

Even with all of that, the opening can go wrong. Prime Day 2018 began at 3 p.m. Eastern time on 16 July, and within minutes many shoppers saw error pages instead of deals. Amazon acknowledged on Twitter that "some customers are having difficulty shopping" and said it was working to fix it. It hasn't published what went wrong, so there's no root cause to learn from here, only the fact that preparation doesn't remove the need for a plan for overload.

An aerial view of many large, flat-roofed windowless buildings and electrical substations among roads and housing
Data centres near Ashburn, Virginia, seen from a plane. Northern Virginia is home to AWS's oldest region, us-east-1, and holds one of the densest clusters of data centres in the world. Capacity for a peak like Prime Day has to exist in buildings like these before the peak begins.Photo: Theodore Christopher, CC0, via Wikimedia Commons

10.2Shedding load to stay up

Suppose, despite the preparation, the cart service gets more requests at 00:01 than it can handle. A naive server accepts them all. Each request waits longer in the queue for a thread, latency climbs, and soon requests take longer than the caller's timeout, the time after which the caller gives up and treats the request as failed. From then on the server is doing work for requests whose callers have already gone, and the callers retry, adding more load.

David Yanacek's article in the Amazon Builders' Library on load shedding describes this spiral. It separates throughput, the rate of requests a server handles, from goodput, the part of that throughput "handled without errors" and fast enough to be useful. Under overload, throughput can stay high while goodput falls towards zero. It gives the turning point: once median latency reaches the client's timeout, about half of all requests time out, and "a latency increase transforms a latency problem into an availability problem".

The fix is load shedding: when a server is near its limit, it rejects extra requests immediately with a cheap error, so the requests it does accept still finish in time. Its practices include checking how much of the caller's deadline is left and dropping requests that can't finish in time, bounding how long a request may wait in a queue, and putting a priority on traffic, so the requests that matter most are the last to be dropped.

For Leo's midnight, the priorities follow from section 1. A request to place an order is worth more than a request to load a "customers also bought" widget, and an Add to Cart is worth more than a product-page view from a web crawler. A page that's missing a recommendations panel still lets Leo buy; a page that times out doesn't. Amazon hasn't published how its retail services rank their traffic. Yanacek's own examples are close to this: it prioritises a load balancer's health checks over normal requests, because dropping health checks makes the load balancer remove the server, and prefers human traffic over crawlers.

Load shedding keeps one service alive under too much traffic. It doesn't help with a different kind of trouble: a bad deployment, a poison request that crashes every server it reaches, or one customer whose traffic overwhelms a shared resource. Those don't need less load; they need a wall.

11Containing failure: cells and shuffle sharding

11.1Cells

So far each service has been one big pool of servers. If a bad change reaches that pool, or one request that crashes servers keeps being retried, the whole service fails, for every customer. Ships have the same problem with water, and solve it with bulkheads: walls that divide the hull into watertight compartments, so a hole floods one compartment and the ship stays afloat.

A 1912 page of ship cross-section diagrams showing watertight compartments, with water confined to some compartments and spilling over low bulkheads in others
From a 1912 book on making ships unsinkable: bulkheads that reach high enough keep a breach to a few compartments, and low ones let water spill from one to the next. A cell is a bulkhead for software, and a dependency shared between cells is a bulkhead that's too low.Image: John Bernard Walker, An Unsinkable Titanic (1912), public domain, via Wikimedia Commons

A cell is a bulkhead for a service. Instead of one pool, run several complete, independent copies of the service, each with its own servers and its own data, and give each a share of the customers. A thin cell router in front sends each request to the right cell using a partition key, such as the customer ID. AWS's 2023 whitepaper on cell-based architecture defines the router as "the thinnest possible layer, with the responsibility of routing requests to the right cell, and only that", and puts the effect in numbers: "If a workload uses 10 cells to service 100 requests, when a failure occurs in one cell, 90% of the overall requests would be unaffected by the failure." The same paper says AWS's service teams have used cell-based architecture "for more than a decade".

The cart service, split into cells by customer
INDEPENDENT CELLSLeoOther shoppersCell routercustomer ID → cellCell 1servers + dataCell 2servers + dataCell 3servers + dataControl planecreates, moves cells
Step 1. Leo's request reaches the cell router, which looks up his customer ID and finds his cell. It does nothing else, so there's very little in it to break.
1 / 5

Cells cost something. Each cell must be big enough to stand alone, which probably wastes some capacity compared with one shared pool. Anything that needs data from many customers at once, such as a report across all carts, has to visit every cell. And choosing the partition key matters: the paper says it should match the "grain" of the service, the natural way it divides "with minimal cross-cell interactions". For a cart, the customer is the obvious grain, because a cart belongs to one customer. Amazon hasn't published which of its retail services are built from cells, or how many cells they use.

11.2Shuffle sharding

Cells contain a failure to a fraction of customers, but the fraction is fixed by the number of cells. With four cells, a poison request from one customer still takes down a quarter of everyone. Colm MacCárthaigh's article in the Amazon Builders' Library describes a way to make the fraction far smaller without adding servers, called shuffle sharding.

Take eight servers. Ordinary sharding splits them into four fixed shards of two, and puts each customer on one shard. A customer who sends a request that crashes servers takes down both servers in their shard, and with them every customer on that shard: a quarter of everyone. Shuffle sharding instead gives each customer their own pair of servers, chosen from all eight. There are 28 different pairs, so customers are spread across 28 combinations. When the bad customer's two servers go down, another customer loses both of their servers only if they were given exactly the same pair. Everyone else still has at least one working server, and a client that retries on the other server carries on.

This program assigns 1,000 customers both ways, lets one customer's poison request crash both of its servers, and counts how many customers lose every server they have. It also computes the number from Amazon's DNS service, Route 53, which the article uses as its example.

Plain sharding versus shuffle sharding: who loses every server?
python
Python
import random
from itertools import combinations
from math import comb
random.seed(3)
 
WORKERS = 8
customers = [f"c{i}" for i in range(1_000)]
 
# Plain sharding: 4 fixed shards of 2 workers; each customer gets one shard.
shards = [(0, 1), (2, 3), (4, 5), (6, 7)]
plain = {c: shards[i % 4] for i, c in enumerate(customers)}
 
# Shuffle sharding: each customer gets its own random pair out of the 8 workers.
pairs = list(combinations(range(WORKERS), 2))
shuffle = {c: random.choice(pairs) for c in customers}
 
def hurt(assign, bad):
    """A poison request from `bad` crashes both of its workers.
    A customer is down only if every one of its workers is down."""
    dead = set(assign[bad])
    return sum(1 for c in customers if set(assign[c]) <= dead) / len(customers)
 
print(f"possible pairs of 8 workers: {len(pairs)}")
print(f"plain shards:   {hurt(plain, 'c0'):.0%} of customers lose every worker")
print(f"shuffle shards: {hurt(shuffle, 'c0'):.1%} of customers lose every worker")
 
# Route 53's numbers: 2,048 virtual name servers, 4 per customer domain.
print(f"ways to choose 4 of 2,048: {comb(2048, 4):,}")
output
C++
possible pairs of 8 workers: 28
plain shards:   25% of customers lose every worker
shuffle shards: 3.3% of customers lose every worker
ways to choose 4 of 2,048: 730,862,190,080

With plain shards, the bad customer takes out 25% of customers. With shuffle shards on the same eight servers, 3.3% lose both servers, roughly the 1 in 28 (3.6%) you'd expect, since each customer's random pair matches the bad customer's with probability 1/28. MacCárthaigh's article puts it the same way: one in 28 is "7 times better than regular sharding".

That last line is why the idea scales so well. Route 53 arranges its capacity as 2,048 virtual name servers and gives each customer's domain a shuffle shard of four. There are roughly 730 billion ways to choose four of 2,048, the article's "staggering 730 billion possible shuffle shards", enough that domains can be assigned so that, in the article's words, "no customer domain will ever share more than two virtual name servers with any other customer domain". An attack on one domain overloads its four servers; every other domain keeps at least two healthy ones.

Decision

How should the cart service's servers be shared between customers?

One shared pool
Every server serves every customer.
  • Best use of capacity
  • Simplest to run
  • A poison request or bad deployment can reach every customer
Cells
Several independent copies of the whole service; each customer lives in one.
  • A failure stays in one cell
  • Deploy one cell at a time
  • Some spare capacity per cell
  • Cross-customer work visits every cell
Shuffle shards inside the service
Each customer gets its own small random set of servers from a shared pool.
  • Tiny overlap between any two customers
  • No extra servers needed
  • Clients must retry on another server
  • Assignment must be stored and kept stable

These combine. AWS's whitepaper describes cells as the outer wall, so that a bad deployment or a broken dependency hurts one cell, and MacCárthaigh's article describes shuffle sharding inside services like Route 53, so that one customer's traffic barely touches anyone else's. Which of these Amazon's retail cart uses isn't published; the cart's natural partition key, the customer, suits both.

12The whole system

12.1Every box, and why it's there

Leo's cart and checkout, end to end
ONE CELLclaimLeo's phone + laptopEdge + cell routersheds, routesCart serviceone row per lineCheckoutsaga, idempotency keyDeal inventorysplit counters, claimsPricingPaymentsauthorise, captureOrders + outboxOrder eventsto fulfilmentCard network
Step 1. At 00:01 Leo adds the deal on his laptop. The edge checks for overload and routes him to his cell; the cart service asks deal inventory for a claim, a conditional decrement on one of the deal's split counters.
1 / 6
ComponentWhat it doesAdded because
Services with their own dataEach team owns one part and its storeOne application and one database couldn't scale or change (§2)
Always-writeable cart storeAccepts writes on any reachable node; merges siblingsA refused Add to Cart is a lost sale (§3, §4)
Leader per partition (DynamoDB)Orders writes to each item; conditional writesSiblings were hard to program against; checkout needs exact checks (§5)
Checkout sagaPrice, claim, authorise, create, with compensationsNo transaction spans four services and a bank (§6)
Deal claims on split countersConditional decrements with 15-minute expiryCached counts oversell; one counter is a hot key (§7)
Authorise, then captureHold funds at order, charge at shippingOrders ship in parts and items can turn out missing (§8)
Idempotency keysA retried tap returns the first resultResponses get lost and taps get repeated (§8)
Outbox and queueOrder work happens after Leo sees "Order placed"Downstream work is slow and the peak is spiky (§9)
Load sheddingRefuses extra work fast, by priorityOverload turns latency into outage (§10)
Cells and shuffle shardsWalls between groups of customersOne bad deploy or poison request shouldn't reach everyone (§11)

12.2From top to bottom

LevelThe choiceData structure or algorithm
OrganisationServices own their data; you build it, you run itPublished interfaces; no shared databases
Cart storage (2007)Always writeable; merge on readConsistent hash ring; preference lists; sloppy quorum (3, 2, 2); hinted handoff; Merkle trees
Cart versions (2007)Detect concurrency, let the app mergeVector clocks of (node, counter); siblings; union merge, or tagged adds (an observed-remove set)
Storage (2012 on)One version per itemPartitions of key ranges; Multi-Paxos group of three across zones; leader leases; write-ahead log quorum of two
Deal inventoryClaim at add-to-cart, with expiryConditional decrement; counter split into slots; claim records with expiry; waitlist
CheckoutA saga, irreversible step lastLocal transactions plus compensations; idempotency keys; outbox
OverloadShed by priorityDeadlines, bounded queues, goodput
IsolationCells, then shuffle shardsCell router by customer ID; random k-of-n server sets (C(2048, 4) ≈ 730 billion)

13What goes wrong, and what it costs

13.1Failures this design has to survive

What happensWhat Leo seesWhat the design does
A storage node holding his cart is downNothingDynamo: a stand-in takes the write with a hint. DynamoDB: two of three replicas still form a quorum
His two devices edit the cart during a network splitA removed item may reappear (Dynamo, union merge)Siblings are merged on the next read; tagged adds or one row per line avoid resurfacing
A million people claim 1,000 deal units"Added to cart", or a quick "join the waitlist"Conditional decrements on split counters; claims expire after 15 minutes
His card is declined at checkout"Payment declined", choose another cardThe saga releases his claim; no order exists
The response to Place your order is lost and he taps againOne orderThe idempotency key returns the first result
The authorisation fails after the order is placedA request to update his paymentRetry or change the card; otherwise the order is cancelled
The email service is slow at midnightThe email arrives lateThe order event waits in the queue; checkout isn't held up
The cart service is overloadedA fast "try again" on low-priority parts of the pageLoad shedding by priority and deadline
A bad deployment reaches one cellCustomers in that cell see errors; Leo, in another, doesn'tCells deploy one at a time and fail alone

13.2The tradeoffs, in one table

DecisionChosenGiven upWhy it was worth it
ArchitectureServices with private data (from 2001)Simple cross-team queries and transactionsTeams scale and ship independently
Cart writes (2007)Always writeable, sloppy quorumsReads that always see the latest writeA refused Add to Cart costs a sale
Cart conflicts (2007)Union merge in the applicationRemovals that always stickAdds are never lost; the shopper reviews the cart anyway
Storage (2012)Paxos leader per partitionWrites during a partition that loses its quorumOne version per item; conditional writes and transactions
Deal stockClaim at add-to-cart for 15 minutesUnits tied up in abandoned cartsShoppers learn at once; no oversell
Hot deal counterSplit into slotsA simple single numberMany partitions share the burst
PaymentAuthorise at order, capture at shippingThe certainty of money in hand at checkoutCharges match what ships
After the orderAsynchronous pipelineEverything finished when Leo sees "Order placed"Fast checkout; downstream absorbs the peak
OverloadShed by prioritySome requests refusedThe rest still finish in time
IsolationCells and shuffle shardsSome spare capacity; more routingA failure reaches a fraction of customers

14Summary

  1. A cart is a sale waiting to happen, so refusing a write to it costs money; a deal's stock and an order are promises, so they must be exact. The two kinds of data get different storage.
  2. Amazon split its monolith into services from 2001, each owning its data behind an interface; a page came to depend on over 150 of them, so each is measured at the 99.9th percentile.
  3. Dynamo (2007) made the cart always writeable: keys on a consistent hash ring, N replicas, quorums of (3, 2, 2), and sloppy quorums with hinted handoff so writes succeed even when home nodes are down.
  4. Vector clocks detect concurrent versions: one counter per coordinating node; if neither clock is at or above the other everywhere, both versions are kept as siblings.
  5. The cart merged siblings by union, so no Add to Cart was ever lost and deleted items could resurface; tagging each add with its clock entry lets a merge tell a removal from an item the other side never saw.
  6. Siblings were rare: 99.94% of cart reads in a day saw one version, and most conflicts came from robots, not people.
  7. DynamoDB (2012) kept Dynamo's partitioning and dropped its versioning: each partition is a Multi-Paxos group of three replicas in three zones, a leased leader orders writes, and two of three must store each one.
  8. Checkout is a saga: price, claim, authorise, create, ordered from cheap and reversible to irreversible, with compensations for each step that can be undone.
  9. Deal units are claimed at add-to-cart for 15 minutes, with conditional decrements on split counters, because a cached count oversells and one counter is a hot key.
  10. Payments are authorised at order and captured at shipping, and idempotency keys make a retried tap return the first result.
  11. Work after the order goes through an outbox and a queue, and overload is handled by shedding low-priority work, while cells and shuffle sharding keep any one failure to a fraction of customers.

15Build this

A Prime Day checkout simulator.

  • Write a key-value store with three replicas and simulated network splits. Implement get and put with vector clocks, return siblings, and merge a cart by union. Show a deleted item resurfacing, then replace the union with tagged adds and show that it stays deleted.
  • Add a deal with 1,000 units as ten split counters with conditional decrements, claim records with a 15-minute expiry (speed up time), a sweeper that returns expired units, and a waitlist that gets them.
  • Write a checkout saga with four steps and compensations. Inject failures at each step and check that no claim, authorisation or order is left behind.
  • Add idempotency keys, then drop 10% of responses and have clients retry. Count orders and authorisations per key; both should be exactly one.
  • Fire 50,000 simulated shoppers at it in two seconds. Add load shedding that refuses "recommendations" requests before "place order" requests, and plot goodput with and without it.

16Interview questions

beginnerWhy did Amazon want a shopping cart that never rejects a write?›

Because a rejected Add to Cart is a customer who might not buy. The Dynamo paper says the cart "must allow customers to add and remove items from their shopping cart even amidst network and server failures". Accepting every write means two versions can briefly exist, and the worst case is a removed item showing up again, which the shopper fixes in one click. That's a much smaller cost than a lost sale.

beginnerWhat's the difference between authorising and capturing a payment, and why does Amazon do both?›

Authorising asks the card's bank to confirm the card is good and reserve the amount, without moving money. Capturing tells the bank to take it. Amazon authorises when the order is placed and captures when each shipment leaves, so a shopper is charged only for what ships, and an order split across warehouses is charged in parts. If the order is cancelled, Amazon tells the bank the authorisation isn't needed.

intermediateTwo versions of a cart have vector clocks {A:2} and {A:1, B:1}. What does that tell you, and what happens next?›

Neither clock is at or above the other on every node: the first has more of A's changes, the second has a change from B the first never saw. So the versions are concurrent, and Dynamo keeps both as siblings. The next read returns both with a context; the cart service merges them (in 2007, by union) and writes the result back with a clock that covers both, such as {A:3, B:1}, which lets every node drop the two siblings.

intermediateA Lightning Deal has 1,000 units and a million people tap at once. How do you avoid overselling, and keep it fast?›

Grant units only with an atomic conditional decrement on the stored count, never from a cached number, which oversells badly at these arrival rates. One counter is a hot key, so split the units across several counters on different partitions and have each claim pick one, retrying on another if it's empty. Record each claim with an expiry, 15 minutes for Amazon's Lightning Deals, so unbought units return to a waitlist, and turn away excess taps early with a waitlist or rate limit.

intermediateLeo taps Place your order, the response is lost, and he taps again. How do you make sure he gets one order?›

Give the attempt an idempotency key when the order-review page is rendered, and send it with every tap. Checkout stores the key with the order; a request with a key it has already completed gets the original result back instead of running the saga again. Following Amazon's Builders' Library guidance, a repeat with different parameters is rejected, and the repeat's response is equivalent to the first success, not an "already exists" error. Each downstream call, such as the authorisation, carries a key derived from it.

deepWhy did DynamoDB drop Dynamo's vector clocks and sloppy quorums?›

Developers found siblings and eventual consistency hard to program against, and Dynamo was hard to operate, so few teams adopted it. DynamoDB makes each partition a Multi-Paxos group of three replicas in three Availability Zones with a leased leader that orders every write, so each item has one version and conditional writes and transactions are possible. The price is that a partition that can't reach two of its three replicas, or is between leaders, can't accept writes; DynamoDB judged that a narrower availability promise, backed by 99.99% and 99.999% targets, was worth the simpler programming model.

deepHow do cells and shuffle sharding limit the damage of a failure, and how do they differ?›

Cells are complete, independent copies of a service, each serving a fixed share of customers chosen by a partition key; a bad deployment or broken dependency in one cell affects only that share, so with ten cells 90% of requests are unaffected. Shuffle sharding works inside a shared pool: each customer gets a small random set of servers, so two customers rarely share all of theirs. With 8 servers in pairs, one customer's poison request takes out 1 in 28 customers instead of 1 in 4, and Route 53's 4-of-2,048 shards make full overlap practically impossible.

17Go deeper

check yourself
A Dynamo cart write with N = 3 and W = 2 finds two of its three home nodes unreachable. Does it fail?›

No. Dynamo uses a sloppy quorum: it writes to the first N healthy nodes on the ring, so stand-ins take the missing copies with hints and hand them back when the home nodes return. The write succeeds once two nodes have it.

Why did Dynamo's cart sometimes show items the customer had deleted?›

Concurrent versions were merged by union. A sibling that never heard about the removal still had the item, and a union can't tell "removed" from "never had", so the item came back.

In DynamoDB, when is a write acknowledged?›

When the partition's leader has its log record stored by a quorum of the replication group: two of the three replicas, in different Availability Zones.

Why does checkout claim the deal unit before asking the bank?›

At a sale most claims fail, and failing at inventory costs nothing outside Amazon. Asking the bank first would mean voiding an authorisation, and a hold on the shopper's money, for every loser.

With 8 servers and shuffle shards of 2, what fraction of customers share both servers with a misbehaving one?›

About 1 in 28, since there are 28 possible pairs: roughly 3.6%, against 25% with four fixed shards of two.

DeCandia et al., 'Dynamo: Amazon's Highly Available Key-value Store' (SOSP 2007)

The always-writeable cart, consistent hashing, vector clocks, sloppy quorums, hinted handoff, Merkle trees and the production numbers. PDF

Elhemali et al., 'Amazon DynamoDB: A Scalable, Predictably Performant, and Fully Managed NoSQL Database Service' (USENIX ATC 2022)

Why DynamoDB kept little of Dynamo's design: Multi-Paxos groups, leases, log replicas, admission control, and Prime Day 2021's 89.2 million requests a second.

Idziorek et al., 'Distributed Transactions at Scale in Amazon DynamoDB' (USENIX ATC 2023)

How transactions across items were added to DynamoDB without slowing single-item reads and writes.

Werner Vogels interviewed by Jim Gray (ACM Queue, 2006), and 'Amazon DynamoDB' (All Things Distributed, January 2012)

The move from the Obidos monolith to services, "you build it, you run it", and why Dynamo became a managed service.

Amazon Builders' Library: shuffle sharding, load shedding, and idempotent APIs

Colm MacCárthaigh on shuffle sharding and Route 53; David Yanacek on goodput and load shedding; Malcolm Featonby on client request tokens.

Reducing the Scope of Impact with Cell-Based Architecture (AWS whitepaper, September 2023)

Cells, cell routers, partition keys and the bulkhead idea, from AWS's own practice.

AWS News Blog Prime Day posts (2019, 2023, 2024, 2025, 2026)

Each year's DynamoDB, Aurora, SQS, CloudFront and fault-injection numbers for Amazon's own systems during Prime Day.

Time & Ordering

Lamport clocks, vector clocks, version vectors and Dynamo's clock truncation, in general. Chapter 26.

Replication & Consistency

Quorums, sloppy quorums, hinted handoff and siblings on a different running example. Chapter 28.

Partitioning

Consistent hashing, virtual nodes, hot keys and DynamoDB's per-partition limits. Chapter 29.

Distributed Transactions

Two-phase commit, sagas, the outbox and idempotency keys for one checkout across two systems. Chapter 31.

Designing Stripe

Authorisation, capture and idempotency keys from the payment processor's side. Chapter 52.

Designing Ticketmaster

Holds on unique seats under contention, waiting rooms and bot defences. Chapter 67.

Reliability

Timeouts, retries with backoff, circuit breakers and load shedding. Chapter 40.