On the Friday, an engineer at the Ardnave Point wind farm was chasing a blade pitch fault and pulled up Tuesday's power curve for turbine AP-0214. The curve had holes in it. Not a clean gap, which would have been a link outage and would have been boring. A dotted line: three readings missing out of four at nine in the morning, then two out of three, then one in two, thinning out until by 09:14 the curve was solid again.
Galewatch, which collects a reading every two seconds from each wind turbine on the farms it monitors and sells the owners dashboards, had deployed a configuration change that morning. Ardnave Point is its largest site, off the west coast, and lesson 014 is the story of how it got there: at two hundred turbines the farm did not fit inside one shard, so Galewatch widened its shard key from farm to the pair (farm, turbine) and split that one farm across two machines. Which half a turbine's readings go to was decided by the shortest line of code anybody writes for this, hash(serial) % 2.
In March the owners switched on a hundred more turbines. The new ones arrive empty, and you size a shard for the year ahead rather than for the morning: three hundred turbines at the five gigabytes a turbine-year lesson 013 worked out is 1,500 gigabytes, and two shards of that is 750 against the 720 gigabyte budget lesson 013 derived from a restore drill and a contractual hour. Thirty gigabytes over. So Ardnave Point needed a third shard, the constant in that line of code changed from 2 to 3, and the config went out the way every config goes out at Galewatch, one ingest box at a time, four boxes, twelve minutes end to end.
For those twelve minutes some of the ingest boxes thought a turbine's readings belonged on shard 1 and the rest thought shard 3.
Nothing failed. Every write succeeded, on whichever machine the box that handled it believed in. Lesson 013 named this exact shape when it warned that a routing map is production state: a wrong answer here is not an error, it is zero rows, because asking the wrong Postgres for a turbine's Tuesday morning returns an empty result that looks exactly like a turbine that was switched off. The dashboards drew the hole and nobody argued with them.
Two thirds of the turbines had changed address, each of them sending thirty readings a minute, and roughly half of those readings went through a box that had not been told yet. Call it twenty four thousand readings on the wrong machine, and nobody can give you the exact figure, because it depends on which of the four boxes took which request. They were all still there. Finding them meant writing a script that asked all three shards for every turbine and compared, which took two days, and the only reason anybody knew to write it was that one turbine in three had no gap at all. Nothing in a network or a disk is fussy enough to skip every third turbine.
Twelve minutes to cause, three days to notice, two more to be sure it was all of it.
The line nobody writes down
shard = hash(key) % n.
It is the correct first thing to write and I would write it too. It needs no table, no lookup, no deploy to update, nothing to back up and nothing to agree on. Lesson 013 spent a section on why that matters: the moment placement becomes data rather than arithmetic, you own a new piece of production state that has to be stored, distributed and right everywhere at once. A formula is the one placement scheme with no state at all. Every service that can compute the hash can answer the question alone, forever, with no network call.
The bill arrives on the day n changes.
Work it for Ardnave Point. Two hundred turbines, keys hashed, % 2 becoming % 3. A key stays where it was only if h % 2 and h % 3 happen to agree, and over any run of six consecutive hash values that happens twice. One third stay. Two thirds move.
That is 133 of the 200 turbines, and at five gigabytes of history each, 665 gigabytes have to be copied from the machine they are on to the machine the new arithmetic says they belong on. Borrowing lesson 013's two hundred megabytes a second, which it measured on a restore rather than on a live copy, so read it as the optimistic end: 665 gigabytes is 3,325 seconds. Fifty five minutes of copying to make room for a hundred turbines that have not sent a single reading yet.
The remainder is worse at scale, which is the wrong way round
Generalise it. Going from n machines to n + 1, a key keeps its home only when h % n equals h % (n + 1), and the fraction for which that is true is one in n + 1.
| Machines | Keys that stay | Keys that move |
|---|---|---|
| 2 to 3 | 1 in 3 | 67 percent |
| 4 to 5 | 1 in 5 | 80 percent |
| 18 to 19 | 1 in 19 | 95 percent |
Galewatch is not a two machine company. After Ardnave Point's third shard it runs nineteen Postgres machines and twelve hundred turbines, with 5.5 terabytes on disk: the 4.5 terabytes lesson 013 sized for nine hundred turbines across nineteen farms, plus the thousand gigabytes Ardnave Point's first two hundred have written in their year. March's hundred have nothing yet. Suppose Galewatch had taken the by-turbine cut lesson 013 priced and the company turned down, so that a reading's home were hash(turbine) % 18, and suppose you add the nineteenth machine.
Eighteen nineteenths of 5.5 terabytes is 5,210 gigabytes. At two hundred megabytes a second that is 26,050 seconds, which is seven hours and fourteen minutes of copying, during which every key's address is in flux and the honest options are to freeze writes or to write to both places and reconcile afterwards. To add one machine to eighteen. The machine itself costs an afternoon of someone's attention and arrives with a warranty.
And notice which way the curve runs. At two machines the formula wastes half of the copy. At eighteen it wastes seventeen eighteenths of it. The bigger your fleet, the worse the deal, which is precisely backwards, because a big fleet is where you add machines most often.
Put the machines in the same space as the keys
Here is the move, and it is older than most of the systems that use it. Hash the keys into some fixed space, say the whole range a 32 bit hash can produce, zero to about four billion. Then hash the machines into the same space, using their names. Now both live on one line, which you treat as a circle by joining the right hand end back to the left. A key belongs to the first machine you meet walking to the right from where the key landed, which on the circle is clockwise, wrapping round the end if you run out of line.
0 2^32 wraps to 0
|----B------D---------------A------C------------|
^k1 ^k2 ^k3 ^k4
k1 landed before B, so B owns it
k2 landed before D, so D owns it
k3 landed before A, so A owns it
k4 ran off the end, wrapped round, so B owns it too
The lookup is a sorted array and a binary search, and it is about ten lines:
ring = sorted((h(f"{node}#{i}"), node)
for node in nodes for i in range(200))
def owner(key):
k = h(key)
lo, hi = 0, len(ring)
while lo < hi: # first point at or after k
mid = (lo + hi) // 2
if ring[mid][0] < k: lo = mid + 1
else: hi = mid
return ring[lo % len(ring)][1] # ran off the end, so wrap
Ignore the range(200) for a moment; one point per machine is the version to understand first.
Now add machine E. It hashes somewhere, say between D and A. Every key that lands between D and E used to walk past E's position and stop at A; now it stops at E. Every other key in the ring is untouched, because nothing else about the walk changed.
before: |----B------D---------------A------C------|
^k5 ^k3
after: |----B------D------E--------A------C------|
^k5 ^k3
k5 moved from A to E. k3 did not move. Nothing else
anywhere in the ring moved at all.
One arc changed hands, and the expected size of a new node's arc on a circle with n + 1 points is one n + 1th of the circle. So adding a machine moves one in nineteen keys where the remainder moved eighteen in nineteen, and the ratio between the two is exactly eighteen, which is the number of machines you already had.
The waste factor of the remainder is the fleet you already have. That is the whole pitch, and the arithmetic on that nineteenth machine, on the by-turbine cut Galewatch did not take, is seven hours and fourteen minutes of copying against twenty four.
Removal works the same way in reverse. Take a machine out and its arc merges into the arc of its neighbour on the right. Only that machine's keys move, and they move to exactly one place. A machine that dies takes one nineteenth of the fleet's keys with it and hands them to a single survivor.
Sit with that last sentence, because it is a problem and not a feature.
Nineteen darts do not make nineteen equal slices
Throw nineteen points at a circle at random and you do not get nineteen arcs of 5.26 percent each. You get a mess. The arcs between uniformly random points are not even approximately equal, and the arithmetic is unforgiving. Stand on any one of the nineteen points and look back the way the keys came: the stretch it owns is longer than twice the average only if none of the other eighteen points landed in it, and each of them misses with probability 17/19. So (17/19) to the eighteenth power, which is 13.5 percent. Nineteen machines times 13.5 percent is between two and three of them holding double their fair share of your data. About one in twenty exceeds three times the average, so expect one machine at triple.
That is lesson 014's skew, arriving from a completely different direction. Lesson 014 got it from the world being lumpy, Kilmore Sands having a hundred and forty turbines when lesson 013's nineteen farms averaged forty seven. This version has nothing to do with the world. It is the hash function's own randomness, and it would happen if every key were identical in size.
The fix is the range(200) in that code block. Give each machine two hundred points on the ring instead of one, by hashing machine#0 through machine#199. Nineteen machines now scatter 3,800 points, each machine owns two hundred small arcs, and its share is the sum of two hundred independent draws rather than one. Sums of many small random things cluster hard around their mean.
How hard is worth knowing, because it is the number that tells you what to set. The relative spread of a machine's share falls off as one over the square root of the points per machine. One point each: the spread is about as large as the share itself, which is how two or three machines in nineteen end up holding double. Two hundred points each: one over the square root of two hundred is seven percent, so machines land within a few percent of their fair share and the busiest of nineteen is maybe fifteen percent heavy rather than the triple we just computed.
One over the square root of the points per machine. To halve your imbalance, quadruple the points. That is the only dial here and it is a cheap one, since 3,800 pairs of integers is a rounding error in memory and a binary search over them is twelve comparisons.
The better argument for virtual nodes is the one that shows up when something breaks. Without them, a dead machine hands all of its keys to one neighbour, which instantly carries double. With two hundred points each, a dead machine's two hundred arcs have two hundred independently chosen right hand neighbours, so its load spreads across all eighteen survivors and each one picks up an eighteenth of what the dead machine held, which is six percent more work than it had a minute ago.
A dead machine either doubles one survivor or adds six percent to eighteen of them. Lesson 020 called the recovering service the one you break, and lesson 027 priced what a box does once it is asked for more than its ceiling. Six percent is a shrug. Double is how you lose the second machine ninety seconds after the first.
What the points cost
The scatter that saves you on failure is a nuisance when the data is heavy and replicated. If a machine owns two hundred arcs, then the set of machines holding copies of its data is effectively the whole cluster, so rebuilding a dead node streams from everywhere, and any repair job that compares ranges has two hundred times as many ranges to compare. Worse, with replicas placed by walking the ring, a wider scatter means a given set of three replicas is more likely to be a set that some pair of simultaneous failures takes out together. Cassandra shipped with 256 tokens per node for years because of the balance, and lowered the default to 16 in version 4.0, once the repair and streaming bill was costing more than the evenness was buying. The number that makes your disks even is not the number that makes your repairs cheap.
The other cost is that you have bought placement you do not control. Redis Cluster takes the opposite position: it defines 16,384 fixed hash slots, maps a key to a slot with CRC16 modulo 16,384, and then keeps an explicit map from slot to node, which an operator can edit. If one slot is hot, you move that slot. On a ring you cannot move anything, because the ring is a function and the only input you have is the set of machine names.
Which is the same trade lesson 014 made when it introduced bucket indirection, key to one of 4,096 fixed buckets and buckets to machines, and handed the algorithm here. 4,096 buckets and 16,384 slots are the same design. The bucket count never changes, so the key's bucket never changes, and all the movement happens in a small table that a person can read.
A formula or a table
Ring against buckets is the argument people have in the room. The axis underneath it is this:
A formula needs no agreement and gives you no control. A table gives you control and needs agreement.
With a ring, every service computes placement alone, and the only thing that must be shared is the list of machine names. With slots or buckets, you can pin a hot key, drain a machine gently, move eleven slots and leave the rest, and in exchange you have a few thousand rows that have to be identical in every process that routes, including the one that started thirty seconds ago and the one that has been up for six weeks. Lesson 048 owns distributing that, and lesson 029 already sent it the same shape when it looked at keeping a revocation list current at every checker.
So which should Galewatch have? At the farm layer, neither, and lesson 014 was right. Twenty farms are twenty keys, not twenty million, and the ring's balance arithmetic above assumed keys so numerous and so small that an arc's length is a good proxy for its load. Place twenty unequal lumps by hash and you get a random assignment of boulders, which is worse than the one a person would make in an afternoon. Lesson 014 said the same thing about buckets, in its own words: they buy you everything when the key is gravel and nothing at all when it is boulders. A ring is a bucket scheme with the table replaced by arithmetic, so it inherits the limit exactly.
At the turbine layer inside a split farm, the ring is right and the remainder is wrong. Ardnave Point's three hundred turbines across three machines is gravel: a hundred pieces per machine, roughly equal in size, and a count that changes whenever a farm grows. That is the layer the hook was about, and it is the layer the course keeps arriving at from different directions. Lesson 047 owns the case where one of those pieces is hot enough to matter on its own.
The part the ring does not fix
Go back to Tuesday morning.
Consistent hashing would have helped and it would not have saved them. Under a ring, adding Ardnave Point's third machine moves one turbine in three instead of two in three, so about half as many readings land in the wrong place during the rollout. Half the damage is worth having. The dotted line is still there.
Because the failure was never the arithmetic. It was that four ingest boxes disagreed for twelve minutes about which machines exist, and a ring built from a different list of names is a different ring. Each box computes confidently, none of them errors, and the key goes to two places depending on who asked. Consistent hashing makes a membership change cheap. It does nothing to make a membership disagreement safe. Somebody still has to decide who is in the cluster and make that decision land everywhere, which is lesson 032 and lesson 048, and it is a harder problem than today's by a wide margin.
What you can do today, cheaply, is stop the disagreement being silent. Redis Cluster does the obvious thing: ask a node for a key whose slot it does not own and it refuses, replying MOVED with the address of the node that does. The client learns the map was stale from the only party that knows. Notice where that check lives. The node does not take the client's word for which slot this is; it works out ownership itself, and that is the only reason a stale client cannot talk it into agreeing. Galewatch's version is a constraint on each shard saying that a row's serial has to hash to this shard, with the three databases getting the new divisor before the four ingest boxes do. A box that has not been told yet then collects an error from Postgres instead of a clean write. Keep the shard id in the row as well, and the handful that still slips through is a query rather than a two day script.
Nobody builds that, because on every ordinary day it catches nothing and costs a hash on every insert. It is worth it anyway, and the reason is in the hook. Twelve minutes of writes landing on the wrong machine is an incident, and twelve minutes of writes being rejected by the wrong machine is a log line at 09:02 and a fixed config by 09:05. The same bug. One check's difference, and three days of nobody knowing.
Recap
The waste factor of the remainder is the fleet you already have. Hashing modulo the machine count moves n out of every n + 1 keys when you add the n + 1th machine, and a ring moves one. The ratio is exactly n, so the scheme gets worse precisely as the fleet gets big enough to need it. Galewatch's nineteenth machine, had it taken the by-turbine cut lesson 013 priced and nobody chose, would be seven hours and fourteen minutes of copying against twenty four.
Put the machines in the same space as the keys. Hash machine names into the same ring as the keys and give each key to the first machine clockwise. Adding a machine steals one arc from one neighbour and leaves the rest of the ring untouched. That is the entire idea; the rest is managing its randomness.
One over the square root of the points per machine. Nineteen random points make arcs so uneven that two or three machines carry double. Give each machine two hundred points and the spread falls to about seven percent. Quadruple the points to halve the imbalance.
A dead machine either doubles one survivor or adds six percent to eighteen of them. Virtual nodes earn their keep at the failure, not at rest, and that is the argument to make when somebody asks why the ring has 3,800 entries in it.
A formula needs no agreement and gives you no control; a table gives you control and needs agreement. A ring is the formula. Fixed slots or buckets, lesson 014's 4,096 and Redis Cluster's 16,384, are the table. You pick which problem you would rather have, and distributing the table is lesson 048's.
Consistent hashing makes a membership change cheap, not a membership disagreement safe. Two processes with different machine lists compute two different rings and neither of them errors. Make the wrong owner loud: a node asked for a key it does not own should say so, which turns three days of silence into a log line.
Check your understanding
A colleague is adding the fifth node to a four node memcached fleet and says the reshuffle does not matter because it is only a cache and the misses will refill. Lesson 008 priced what Marlow Books, the four person online bookshop in this course, is offered without its cache, and lesson 020 collected the shape of a recovering service. Decide whether they are right, and name the one measurement that settles it.
You inherit a service placing user sessions with
hash(user_id) % 12. You need to grow to sixteen nodes. Give two migration plans, one that keeps the remainder and one that moves to a ring, and say what each costs in downtime, in code and in the weeks after.Somebody proposes four thousand virtual nodes per machine, reasoning that if two hundred is good, more is better. Using the square root rule, say what the extra 3,800 points per machine buy, then give the two costs that make you refuse.
Galewatch has nineteen Postgres machines and a routing table. Argue for replacing the table with a ring keyed on
(farm, turbine)for the whole fleet, as strongly as you can, then give the property of Galewatch's data that kills it and the one farm you would point at first.Your ring has twelve nodes and one point each, and node 7 dies at the busiest hour. Work out what happens to node 8, say what you would do in the first five minutes, and say what you would change before the next hour like it.
A reviewer says that making every shard check it owns the rows it is handed is paranoid, since the routing code is five lines and has never been wrong. Give the strongest version of their case, then the version of the hook's incident you would describe to change their mind, and the cost of the check in writes per second.
Next lesson
032 Leader Election: Why Someone Has to Be in Charge. Today ended on four ingest boxes holding different ideas of which machines exist and none of them being wrong about it; next lesson is about how a group of machines agrees on a single answer to a question like that, and what it costs them while they are deciding.