yield point
Famous Papersdeepupdated 2026-08-26

Dynamo, and what it gives up

R+W>N is a proof, not a guideline. Dynamo keeps writing through a partition by making the set it was a proof about stop existing.

Three pages on this site each describe one clean idea. Consistent hashing decides which nodes a key belongs to, by hashing both keys and nodes into one space and walking clockwise.[3]paperConsistent Hashing and Random TreesKarger, D. et al., STOC 1997, 1997 A quorum decides how many of them have to answer. Vector clocks decide what to do when two answers disagree.

Dynamo is what those become when the requirement is that writes keep succeeding through a partition.[1]paperDynamo: Amazon's Highly Available Key-value StoreDeCandia, G. et al., SOSP 2007, 2007 The paper is worth reading for many reasons, but the one that belongs here is that it is unusually honest about which guarantee it is trading away, and this page is about watching that trade happen.

R=2, W=2, N=3, so R+W>N and no read can miss a write. Press 'Split the cluster' and watch writes start failing. Then turn on 'Sloppy quorum': the writes come back, and the invariant that made R+W>N mean anything goes red.

n0n1n2n3n4n5n6n7n8n9n10n11
Writes accepted
0%
Theory says
97%
Reads missing the write
0%
R + W > N
yes (2 + 2 > 3)
Strict stale chance
0%
Quorum
strict
Hinted writes
0
Values on the key’s replicas
1
tick 0 / 3000
Break it
Do

N: how many nodes the key belongs to, walking clockwise from where it lands.

Accept a write from any W nodes that answer, not only from the key's own N.

Safety properties
  • ✓
    Every read quorum shares a node with the write that preceded it

    held on every tick so far

  • ✓
    Every acknowledged write is on at least W nodes

    held on every tick so far

The proof it starts from

A quorum system with N replicas, R for a read and W for a write, guarantees a read sees the latest write when R + W > N.[2]paperWeighted Voting for Replicated DataGifford, D. K., SOSP 1979, 1979 That is not a rule of thumb. It is pigeonhole: two subsets of one N-element set that are together larger than N have to share a member, and the shared member holds the write.

The widget states the same thing as a number. A write lands on W of the N replicas and a read asks R of them, so the read misses the write entirely only when all R of its nodes come from the N-W that were not written:

P(read misses the write) = C(N-W, R) / C(N, R)

Set Read quorum and Write quorum both to 1 and watch two thirds of reads miss the write, and the overlap invariant go red. That configuration is a perfectly reasonable choice for a cache and a catastrophic one for anything a user expects to read back.

The cost of keeping it

Identical clusters, identical quorum settings, the same split. Press 'Split the cluster' and read the two rows against each other: this is the trade Dynamo makes, and there is no third column where both numbers are good.

Strict quorum

Only the key’s own N replicas will do, so R + W > N stays a proof and writes fail when they cannot be reached.

n0n1n2n3n4n5n6n7n8n9n10n11
Writes accepted
94%
Theory says
97%
Reads missing the write
0%
R + W > N
yes (2 + 2 > 3)
Strict stale chance
0%
Quorum
strict
Hinted writes
0
Values on the key’s replicas
1

Sloppy quorum

Any W nodes that answer will do, so the write lands somewhere and the overlap argument no longer applies.

n0n1n2n3n4n5n6n7n8n9n10n11
Writes accepted
100%
Theory says
100%
Reads missing the write
3%
R + W > N
yes (2 + 2 > 3)
Strict stale chance
0%
Quorum
sloppy
Hinted writes
5
Values on the key’s replicas
1
Measured at tick 120 - both sides, same seed, same inputs
MeasureStrict quorumSloppy quorumGap
Writes accepted0.941.001.1×
Reads that missed the latest write0.000.03∞
tick 120
Break it
Do

N: how many nodes the key belongs to, walking clockwise from where it lands.

Press Split the cluster on both arms.

The strict side stays correct and refuses around half of its writes. Nothing is broken: a client on one side of the split can only reach the replicas on its side, and if fewer than W of the key’s own N are over there, the write has to fail. Eleven healthy nodes sit there able to store the value and are not allowed to.

That is what choosing consistency costs during a partition, and the number is more useful than reciting the letters CAP.[4]talkTowards Robust Distributed Systems (the CAP conjecture)Brewer, E., PODC 2000 keynote, 2000

The trade

Dynamo’s answer is the sloppy quorum. Walk the ring past the nodes that are unreachable and give the write to the first W that answer, with a hint saying who it really belongs to. When the partition heals, the hint is handed back.

Watch the sloppy arm: writes go back to essentially 100%. And the overlap invariant goes red.

Before that, though, notice what the sloppy arm looks like while everything is fine. Its stale rate is zero, its hint count is zero, and its answers are identical to the strict arm’s. A sloppy quorum never walks past the preference list until it has to, so every test you run on a healthy cluster passes, and the guarantee you gave up is one you only needed on the day it was gone.

The reason is worth being precise about, because “sloppy quorums are eventually consistent” skates over it. R + W > N was never a statement about the numbers R, W and N. It was a statement about two subsets of the same N-element set. A sloppy quorum takes the first W nodes that answer, so the write set is drawn from the whole ring while the read set is drawn from somewhere else. The arithmetic is untouched and the proof is gone.

What comes back out

The system now has to hand you two values with no ordering between them, which is what vector clocks are for: they say whether one write knew about the other or whether the two were genuinely concurrent. Dynamo cannot resolve genuine concurrency, so it returns both and makes reconciliation the application’s job.

Watch Values on the key’s replicas while you drive the widget. It reads 1 on a healthy cluster and 1 through a partition under a strict quorum, because only the side holding W of the key’s own replicas is allowed to write and there is never a second branch. Turn on the sloppy quorum and split the cluster and it reads 2: the key now genuinely has two current values, neither descended from the other, and no amount of waiting will decide between them.

That gauge is the honest one to watch, and it is not the same as counting siblings. A sibling is what a single read is handed when it happens to sample both branches, which is a window one write wide. Two live branches on the replicas last as long as the split does - which is what Merkle-tree anti-entropy goes looking for, and why the alert you want is on hint counts rather than on conflicts your clients happened to notice. The shopping cart in the paper merges by union, which is why a deleted item could come back but an added one was never lost - a deliberate choice about which error a customer forgives.

What to do with this

Decide whether you need read-your-writes, then pick the quorum. If you do, a strict quorum with R+W>N is the only thing that gives it, and you are choosing to fail writes during a partition. That is a legitimate choice and it should be a deliberate one.

Treat hinted handoff as a consistency event, not an availability feature. It is doing its job precisely when your guarantee is not there. Alert on hint counts, because they mark the window in which reads may be missing writes.

Write the merge function before you need it. Under a sloppy quorum siblings are not an edge case, they are the design. Deciding at three in the morning what to do with two versions of a cart is worse than deciding it now.

Do not carry R+W>N to a system without checking what N is a set of. The arithmetic travels. The proof does not.

The dial

You gainWrites that keep succeeding through a partition that would stall a strict quorum
You payThe overlap guarantee that made R+W>N worth quoting, and siblings the client has to reconcile

The line - when someone asks

Dynamo composes three things this site covers separately: consistent hashing picks the N nodes a key belongs to, a quorum decides how many must answer, and vector clocks order the answers. Under a partition a strict quorum refuses about half the writes, so Dynamo uses a sloppy one and takes the first W nodes that reply whether or not the key belongs to them. Writes survive. The R+W>N argument does not, because it was a pigeonhole proof about two subsets of one N-element set, and there is no longer one set.

Recall

Loading…

Where are you with this?

Saved on this device. Sign in to keep it across devices.

Sources

Primary

Secondary