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 Trees 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 Store 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.
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 Data 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
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.
Sloppy quorum
Any W nodes that answer will do, so the write lands somewhere and the overlap argument no longer applies.
| Measure | Strict quorum | Sloppy quorum | Gap |
|---|---|---|---|
| Writes accepted | 0.94 | 1.00 | 1.1× |
| Reads that missed the latest write | 0.00 | 0.03 | ∞ |
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)
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
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
- [1]
- [2]
- [3]
Secondary
- [4]
- [5]