KnowSys

Partitioning & Rebalancing

Follow twelve keys as a database outgrows one machine: how each key finds the node that owns it, why adding one node can force three quarters of the data to move, and what real systems do instead, including how they move data while still serving traffic.

⏱ 40 min read◆ BeginnerAssumes: a terminal and Python; chapter 28 (replication) helps
Start reading

You keep user profiles in a database, and the table has outgrown one machine. It might be too big for the machine's disk, or too busy for its CPU. So you buy three machines and split the profiles among them. From now on, every read and write for "user 7" has to go to the one machine that holds user 7, and every client that talks to the database has to know which machine that is.

The obvious rule is to spread users out by number. User n goes to machine n mod 3, meaning the remainder after dividing n by 3. It balances perfectly, and it works until the day you buy a fourth machine. The rule becomes "n mod 4", and of the users numbered 0 to 11, nine now belong to a different machine than before. A cluster that has to copy three quarters of its data every time it grows is a cluster nobody wants to grow.

Deciding who owns what is called partitioning, and this chapter asks one question about it: when data is split across machines, how does each piece of it find its owner, and how much has to move when the cluster changes? We'll follow those twelve users from this first simple rule, through the designs real systems use to make growth cheap, to moving data while the system keeps serving traffic.

01Giving every key an owner

1.1Splitting the data

Let's give the pieces names. Each machine in the cluster is a node. Each record in the database is found by a key, such as the user number, and the record itself is the value. A slice of the data that lives together on one node is a partition. (Some systems call it a shard, and the two words mean the same thing here.)

If you've read the previous chapter, you'll know replication, which puts the same data on several nodes so that one can fail without losing anything. Partitioning does the opposite job: it puts different data on different nodes, so the total can be bigger than any one machine, and the work can be shared. Real systems do both. Each partition is replicated, so a node holds copies of several partitions, and each partition lives on several nodes. For the rest of this chapter we'll look at one copy of each partition, because the question we care about is where partitions go.

Every partitioning scheme has to answer two questions:

  1. Which partition does this key belong to?
  2. Which node holds that partition, and how does a client find out?

For the first few sections we'll treat them as one question, as our three-machine rule does: key to node, in one step. Section 2 pulls them apart, and that turns out to be the most useful idea in the chapter.

1.2Hash, then modulo

Our rule worked on user numbers, but real keys are often strings such as user:42, and you can't take a remainder of a string. So the first step is to turn the key into a number with a hash function. A hash function takes any key and returns a large number. The same key always gives the same number, and different keys give numbers that look random, so even keys like user:41 and user:42 land far apart. The second step is the remainder: hash(key) mod N, where N is the number of nodes. We'll write % for mod, as most programming languages do.

To follow the arithmetic in our heads, we'll pretend the hash of key n is just n itself. A real hash scrambles the numbers, and the later experiments use real ones, but nothing in this section depends on the scrambling.

The appeal of this rule is that it needs no lookup. Any client can compute the owner of any key from the key and N alone, without asking anyone. With three nodes, key 7 goes to node 1, because 7 % 3 = 1.

1.3What happens when a node joins

The rule balances the keys perfectly, and its weakness appears when N changes. Here's our twelve keys on three nodes, followed through the arrival of a fourth:

Twelve keys on three nodes, then a fourth node arrives
Node 0Node 1Node 2Node 3newkey 00 % 3 = 0key 11 % 3 = 1key 22 % 3 = 2key 33 % 3 = 0key 44 % 3 = 1key 55 % 3 = 2key 66 % 3 = 0key 77 % 3 = 1key 88 % 3 = 2key 99 % 3 = 0key 1010 % 3 = 1key 1111 % 3 = 2
Step 1. Twelve keys on three nodes. Each key goes to the node numbered key mod 3, so key 7 is on node 1 (7 ÷ 3 leaves remainder 1). Every node holds four keys.
1 / 4

Look at the last frame and notice that the keys moved around among nodes that were already in the cluster. Key 4 left node 1 for node 0, and nothing about the new node asked for that. It moved only because the divisor changed, and a different divisor gives a different remainder for most keys.

A key stays put only when hash % 3 and hash % 4 happen to agree. For twelve keys that holds for 0, 1 and 2. In general, going from N to N+1 nodes leaves only one key in N+1 where it was, which means a fraction N/(N+1) of the keys moves. Try it before reading on.

Predict before you read on

1 million keys are spread across 10 nodes with hash(key) % 10. You add one node and switch to hash(key) % 11. Roughly what fraction of keys now belong to a different node?

Moving 91% of your data to add 10% capacity is a full reshuffle of the cluster. Every copy crosses the network, and while a key is in flight, requests for it can reach the wrong node. The best any scheme could do is move exactly the new node's fair share, 1/(N+1) of the keys, and move them all to the new node. The rest of this chapter is a search for schemes that get there.

The problem in hash % N is that the owner of a key depends on the number of nodes, which changes. That suggests the first fix: stop computing the owner from the node count at all.

02Fixed partitions and a map

2.1Splitting the two questions apart

Remember the two questions from section 1.1. Our rule answered both at once, so a change in the number of nodes changed every key's answer. Suppose we answer them separately. Pick a number of partitions when we create the system, and make it larger than the number of nodes we'll ever have. The first function maps a key to a partition, hash(key) % P, where P is that fixed number. The second is a small table that says which node holds each partition. We'll call it the partition map.

The first function never changes, because P never changes, so a key's partition is permanent. Adding a node changes only the table: we hand the new node a fair share of whole partitions from the others, copy exactly those partitions, and edit a few table rows.

?Why not just map each key straight to a node?

Because then any change to the set of nodes changes the mapping of keys, and changing the mapping of a key means moving its data. With a layer of partitions in between, adding a node means reassigning a few whole partitions: a small edit to a small map, followed by a bulk copy of exactly those partitions.

Here is the same growth as in section 1, with six partitions of two keys each. A key's partition is its number mod 6, and each node owns two partitions:

Six fixed partitions on three nodes, then a fourth node arrives
Node 1Node 2Node 3Node 4newp0keys 0, 6p1keys 1, 7p2keys 2, 8p3keys 3, 9p4keys 4, 10p5keys 5, 11
Step 1. Twelve keys sit in six partitions. A key's partition is key mod 6, fixed forever, so key 7 is in p1 (7 % 6 = 1) and key 1 is in p1 too. The partition map says node 1 owns p0 and p1.
1 / 5

The new node got four keys where a perfectly fair share would be three, because partitions are lumps of two keys. With thousands of partitions the lumps are small, and the share gets as close to fair as you like. The cost is a bigger map, which is still tiny next to the data.

2.2Redis Cluster's 16,384 slots

Redis Cluster does exactly this. It calls its partitions hash slots, and there are always 16,384 of them. Here is how it picks a key's slot:

C
unsigned int keyHashSlot(char *key, int keylen) {
    int s, e; /* start-end indexes of { and } */
 
    for (s = 0; s < keylen; s++)
        if (key[s] == '{') break;
 
    /* No '{' ? Hash the whole key. This is the base case. */
    if (s == keylen) return crc16(key,keylen) & 0x3FFF;
 
    /* '{' found? Check if we have the corresponding '}'. */
    for (e = s+1; e < keylen; e++)
        if (key[e] == '}') break;
 
    /* No '}' or nothing between {} ? Hash the whole key. */
    if (e == keylen || e == s+1) return crc16(key,keylen) & 0x3FFF;
 
    /* If we are here there is both a { and a } on its right. Hash
     * what is in the middle between { and }. */
    return crc16(key+s+1,e-s-1) & 0x3FFF;
}

crc16 is the hash function. & 0x3FFF keeps the lowest 14 bits of the result, which is the same as % 16384 because 16,384 is a power of two. The slot count never changes, so a key's slot never changes. Only the slot-to-node map does.

The rest of the function handles curly braces. If a key contains {…}, only the text between the braces is hashed. This is a hash tag, and it lets you force related keys into one slot: {user:1000}:cart and {user:1000}:profile both hash only user:1000, so they live in the same slot on the same node, and commands that touch both keys work. Without the tag, the two keys would most likely sit in different slots, and a command that names both fails with a CROSSSLOT error.

Instagram did the same thing in Postgres. Their sharding post describes several thousand logical shards, each a Postgres schema, mapped to far fewer physical servers. Moving a logical shard moves a schema, and no key's shard ever changes.

2.3When the fixed number is too small

The partition count is the one number in this design you can't change later, and some systems make it small. Kafka, which stores streams of messages (chapter 23), splits each stream, called a topic, into partitions, and a message's key picks the partition. Its default partitioner is a plain modulo over the partition count:

clients/src/main/java/org/apache/kafka/clients/producer/internals/BuiltInPartitioner.java
apache/kafka @ 3.7.0 ↗
java
    public static int partitionForKey(final byte[] serializedKey, final int numPartitions) {
        return Utils.toPositive(Utils.murmur2(serializedKey)) % numPartitions;
    }

?Why doesn't Kafka work like Redis Cluster here?

Kafka has the same two layers as Redis: a key picks a partition with the formula above, and a separate map says which server holds each partition. The difference is the size of the fixed number. A topic usually has tens of partitions, where Redis always has 16,384 slots, so a busy topic can outgrow its count. The only way out is to add partitions, which changes numPartitions in the formula, and that brings back the % N problem from section 1. Kafka's docs warn that "this partitioning will potentially be shuffled by adding partitions but Kafka will not attempt to automatically redistribute data in any way", and Kafka can't reduce the partition count at all. Existing messages stay where they are, new ones for the same key go to a different partition, and since Kafka only promises message order within one partition, the per-key order breaks across the change.

Fixed partitions work when you can name a count that is big enough, and when you can afford a map that every client consults. Sometimes you can't do either. A fleet of cache servers, for instance, has no central table to keep, and its nodes come and go all day. For that case we want a way to compute an owner that survives changes to the node set without a table.

03The hash ring

3.1Keys and nodes on one circle

The idea comes from Karger and colleagues' 1997 paper on consistent hashing, written for web caches. Hash the keys onto a circle of values, as a clock face has positions 0 to 11. Then hash the nodes onto the same circle, by hashing each node's name. Each node's spot on the circle is called its token. A key belongs to the first node you reach by walking clockwise from the key's position.

Take our twelve keys with the circle having twelve positions. Suppose hashing the names of three nodes places them at positions 2, 6 and 10:

01234567891011tokenABCownerAAABBBBCCCCA
A twelve-position circle, laid flat. Node A has a token at 2, B at 6 and C at 10. A key belongs to the first token at or after it, wrapping around from 11 to 0.

Key 7 walks clockwise to position 10 and belongs to C. Key 11 wraps around the end of the circle and belongs to A. The stretch of keys between one token and the one before it, such as 7 to 10 for C, is an arc, and everything in an arc has the same owner. A real implementation keeps the tokens in a sorted array, and finds a key's owner by binary search, which halves the list at each step, so a lookup among T tokens takes about log T steps.

Now add a fourth node. Its name hashes to position 8, which cuts C's arc in two:

A fourth node joins the ring
Node Atoken at 2Node Btoken at 6Node Ctoken at 10Node Djoins at 8key 11key 0key 1key 2key 3key 4key 5key 6key 7key 8key 9key 10
Step 1. Three nodes on the circle, tokens at 2, 6 and 10. Each key belongs to the first token clockwise, so A holds 11, 0, 1 and 2, B holds 3 to 6, and C holds 7 to 10.
1 / 4

The property to remember is that when a node joins, every key that changes owner moves to that node. When a node leaves, its keys pass to the next node clockwise and nothing else moves. Compare that with the nine keys that shuffled among old nodes in section 1.

3.2Uneven arcs and virtual nodes

D got two keys where a fair share was three. That was luck, and with one token per node the luck can be much worse. Names hash to random positions, and random points on a circle leave some gaps several times longer than others. A node's share of the keys is the length of its arc, so some nodes get far more than their share.

The fix is to give each node many tokens, called virtual nodes. Each node's share is then the sum of many small arcs, and sums of many random lengths average out. When a node joins, it cuts many small pieces from many neighbours instead of one big piece from one, and when a node leaves, its load spreads across many nodes instead of landing on one.

A ring of eight tokens t1 to t8, coloured by owner, with a token map showing node n1 owning t1 and t5, n2 owning t3 and t7, n3 owning t2 and t6, and n4 owning t4 and t8
Virtual nodes as Cassandra draws them: four nodes with two tokens each, so the ring has eight arcs. Each node's tokens are spread around the circle, so its share is two separate arcs, and a node that leaves hands one arc to one neighbour and the other to a different one. Real clusters use 16 or more tokens per node.Image: Apache Cassandra documentation, Apache License 2.0

A simulation of 200,000 keys on 10 nodes growing to 11 shows the effect. With one token per node, the busiest node in the median run held 2.8 times its fair share. That is the node that fills its disk first, and it fills it nearly three times too early. With 256 tokens per node, the worst of nine layouts was only 16% over. Section 4.2 has the full table.

Balance is half the story. The other half is how many keys move when the cluster grows, and the next subsection lets you count them yourself.

3.3Trying it: count the moves

This script assigns 100,000 keys to five nodes two ways, adds a sixth node, and counts how many keys end up somewhere new. h hashes a string to a number with MD5. mod_owner is the rule from section 1. ring builds the circle, with 100 tokens per node placed by hashing node0#0, node0#1 and so on, and sorts them. ring_owner finds the first token after a key's hash using bisect, Python's binary search, and % len(points) wraps around from the end of the circle to the start.

Assign 100,000 keys to 5 servers by hash % N and by consistent hashing, add a 6th server, count the keys that move
python
Python
import bisect, hashlib
 
def h(s): return int(hashlib.md5(s.encode()).hexdigest(), 16)
 
keys = [f"user:{i}" for i in range(100_000)]
 
def mod_owner(k, n): return h(k) % n
 
def ring(nodes, vnodes=100):
    pts = sorted((h(f"{n}#{v}"), n) for n in nodes for v in range(vnodes))
    return [p for p, _ in pts], [n for _, n in pts]
 
def ring_owner(k, r):
    points, owners = r
    return owners[bisect.bisect(points, h(k)) % len(points)]
 
old_nodes = [f"node{i}" for i in range(5)]
new_nodes = old_nodes + ["node5"]
 
moved_mod = sum(mod_owner(k, 5) != mod_owner(k, 6) for k in keys)
r_old, r_new = ring(old_nodes), ring(new_nodes)
moved_ring = sum(ring_owner(k, r_old) != ring_owner(k, r_new) for k in keys)
 
print("Add a 6th server to a cluster of 5 (100,000 keys)")
print(f"  hash(key) % N        : {moved_mod / len(keys):6.1%} of keys move")
print(f"  consistent hashing   : {moved_ring / len(keys):6.1%} of keys move   (ideal: {1/6:.1%})")
output
C++
Add a 6th server to a cluster of 5 (100,000 keys)
  hash(key) % N        :  83.4% of keys move
  consistent hashing   :  16.8% of keys move   (ideal: 16.7%)

With hash % N, adding one node sent 83.4% of the keys to a different node. That matches the N/(N+1) rule from section 1.3, which gives 5/6 here. On the ring, only 16.8% moved, close to the ideal 1 in 6, which is exactly the share the new node has to receive.

Every moved key is data copied across the network, and while it moves, requests for it may hit the wrong node. A scheme that reshuffles 83% of the data whenever the cluster changes makes adding capacity a major event. A scheme that moves only the new node's share makes it routine.

3.4Tokens in a real database

Cassandra, a database built on a ring, calls the virtual nodes of a machine its tokens, and it used to give every node 256 of them. In version 4.0 it lowered the default num_tokens to 16. Fewer tokens make repair and streaming cheaper, because a node's data then lives in fewer separate pieces. Sixteen random tokens would give uneven shares, so Cassandra paired the change with a token allocation algorithm (allocate_tokens_for_local_replication_factor) that places a new node's tokens where they even out ownership, instead of at random.

The ring lets any client compute an owner from the list of nodes alone, but it still needs many tokens to balance well. Two newer schemes reach the ideal without a ring at all.

04Jump hash, rendezvous hashing and the rest

4.1Two ways to compute an owner

Jump consistent hash (Lamping and Veach, 2014) is a few lines of arithmetic that maps a 64-bit key to a bucket numbered from 0 to N−1. It needs no memory at all, moves the ideal 1/(N+1) of the keys on average when N grows by one, all of them to the new bucket, and balances almost perfectly. Its limit is that buckets are plain numbers, so you can only add or remove at the end. That is a good fit for a storage system whose shards are replicated and never disappear from the middle, and a poor fit for a cache fleet where any node can die.

Rendezvous hashing, also called highest-random-weight (Thaler and Ravishankar, 1998), works differently. For each key it scores every node with hash(key, node) and picks the node with the highest score. Remove any node, and only the keys it owned move, because the next-highest score takes over each one. The price is that every lookup computes one hash per node, which is fine for tens of nodes and slow for thousands.

4.2Side by side

Here are all the schemes on the same 200,000 keys, growing from 10 nodes to 11. The ideal is 9.1% of keys moved, since 1/11 is the new node's fair share, and a busiest-node figure of 1.00 would be perfect balance. The middle column divides the busiest node's key count by the average count. Where a node lands on a ring depends on its name, so the last column reruns the ring nine times with nine different sets of node names and reports the median run, with the worst run in brackets. The exact values depend on the node names and the hash function, so your own run will give different decimals, but the same pattern.

SchemeKeys moved (ideal 9.1%)Busiest node ÷ meanBusiest ÷ mean, median of 9 seeds (worst)
hash % N90.9%1.01—
Ring, 1 token per node13.7%2.312.77 (3.76)
Ring, 16 tokens per node11.4%1.361.31 (1.82)
Ring, 100 tokens per node9.4%1.141.19 (1.27)
Ring, 256 tokens per node9.1%1.071.11 (1.16)
Jump consistent hash9.1%1.01—
Rendezvous (HRW)9.2%1.01—

Read the table from the top. hash % N balances perfectly and moves almost everything. A ring with one token per node moves roughly the right number of keys but leaves the busiest node with more than twice its share. Adding tokens fixes the balance gradually, and jump hash and rendezvous hashing balance almost perfectly with no tokens at all.

Each scheme has its own trade-offs, summarised below with the systems that use them. The "Weights?" column asks whether you can make one node take a bigger share, because its machine is bigger.

SchemeLookup costMemoryRemove any node?Weights?Used by
Fixed slotsO(1) table lookupSlot tableYes, reassign its slotsBy slot countRedis Cluster, Instagram, Couchbase vBuckets
Ring + vnodesO(log T)T tokensYesBy token countCassandra, Riak, Dynamo
Jump hashO(log N) arithmeticNoneOnly the lastNoStorage shard selection
RendezvousO(N) hashesNode listYesYes, with weighted scoresGitHub's GLB director (a variant)
MaglevO(1) table lookupLookup tableYes, with small disruptionYesGoogle's network load balancers

Maglev hashing belongs to load balancing more than storage. It trades a little extra disruption on changes for a fast table lookup and near-perfect balance.

Every scheme in this section decides where a key lives by hashing it, and the property that makes hashing so good at spreading load also has a price. It scatters neighbouring keys. Keys 4, 5, 6 and 7 are neighbours in the key space, and they sit on three different nodes. That matters whenever a query asks for a stretch of keys.

05Keeping neighbours together: range partitioning

5.1Sorted keys, cut into ranges

Suppose the keys in our example are order numbers, and a query asks for orders 4 to 7. Reading a stretch of neighbouring keys in order is called a range scan, and databases do it constantly: all orders from one day, all messages in one conversation. With hashing, those four keys sit on three different nodes, so the query has to ask all three. If the neighbours lived together, one node could answer.

That's the idea of range partitioning. Sort the keys, cut the sorted order into contiguous ranges, and give each range to a node:

01234567891011hash % 3N0N1N2N0N1N2N0N1N2N0N1N2rangesN0N0N0N0N1N1N1N1N2N2N2N2
The same twelve keys under two rules. Under hash mod 3 (top), keys 4 to 7 are spread over all three nodes. Under ranges (bottom), they all sit on N1, so a scan of 4 to 7 touches one node.

Google's Bigtable called these ranges tablets. HBase and TiKV call them regions, and CockroachDB calls them ranges.

5.2Splitting and finding ranges

Nobody chooses the cut points up front. A table starts as one range, and a range splits when it grows too big. In CockroachDB the default is to split at 512 MiB (range_max_bytes) and to merge adjacent ranges that fall below 128 MiB.

Finding a key's range is a lookup in a sorted index of range boundaries. CockroachDB stores that index in the key space itself, as two levels of meta ranges, and every node caches it. The diagram below follows one split. Real keys are usually strings, so it uses keys that start with letters instead of our numbers. [m, z) means every key from m up to, but not including, z:

A range split in a range-partitioned store
●
✎
Writes
keys arrive
▭
Range [m, z)
on node 2
✂
Split
pick a middle key
▭▭
Two ranges
[m, s) and [s, z)
☰
Meta index
boundaries → nodes
⇆
Rebalancer
moves one half
Step 1. Writes for keys between m and z all land in one range, held (with its replicas) by node 2.
1 / 6

Notice that this is the same two-step design as section 2. A split edits a small map, and the data only moves later, when the rebalancer decides to.

A tree with the meta1 range at the top, two meta2 ranges below it covering keys up to m and from m, and four data ranges [a,g), [g,m), [m,s) and [s,max) at the bottom
CockroachDB's two-level map. To find the range holding a key, a node reads meta1, which says which meta2 range covers the key, and that meta2 range says which data range holds it. A split like the one above only edits one meta2 entry.Image: CockroachDB documentation, Cockroach Labs, CC BY 4.0

?Why split by load as well as size?

Because a small range can still be the busiest one in the cluster. CockroachDB's load-based splitting splits a range whose request rate crosses a threshold, choosing the split key from sampled requests so the two halves get similar load. Then the halves can live on different nodes.

CockroachDB Key Visualizer heat map: keyspace on the vertical axis, time on the horizontal axis, mostly black with a large bright red block covering the rows of one table
CockroachDB's Key Visualizer, with the keyspace top to bottom and time left to right, each cell shaded from black (idle) to red (busy). Nearly all the traffic is on one table, the red block at the bottom. Cockroach Labs' docs use this screenshot to show a run of range splits: a range taking a heavy stream of writes keeps splitting, so its pieces can be spread over more nodes.Image: CockroachDB documentation, Cockroach Labs, CC BY 4.0

5.3The hot-spot problem

Range partitioning has one weakness, and it appears when keys arrive in order. Suppose our keys are timestamps, and new keys 12, 13, 14 and so on keep arriving. They are all larger than every existing key, so every one of them lands in the last range, the one on N2. If the key is a timestamp, an auto-increment ID or anything that starts with one, every new write goes to the last range. Your cluster has fifty nodes and one of them takes all the inserts. A node that takes a disproportionate share of the traffic is called a hot spot.

?Why doesn't splitting fix a time-ordered hot spot?

Because the new half gets all the new writes too. Split [t0, now) into two and every insert still lands in the upper half. All that moves is the hot spot, to whichever node owns the newest range.

Every fix breaks the ordering on purpose. Two of them lead the key with a high-cardinality column, meaning one with many distinct values such as a user id, so that writes for different users fall in different places:

FixHowWhat you lose
Prefix with a hashKey = hash(user) % 16 + timestampTime-range scans become 16 scans
Lead with a high-cardinality columnKey = (user_id, timestamp) instead of (timestamp)Global "latest events" queries
Random or hashed IDsA random 128-bit ID (UUIDv4) instead of a sequenceInserts no longer go to one end of the on-disk index (a B-tree), so they touch more pages; see the storage engines chapter
Hash-sharded indexCockroachDB's USING HASH does the prefixing for youOrdered scans on that index

5.4No scheme wins every row

We've now met two families, hashing and ranges, and each one gives up something the other keeps. Here are the goals a partitioning scheme has to balance, and what works against each:

GoalWhy it mattersWhat works against it
Even dataNo node runs out of disk firstUneven key distribution, partitions of different sizes
Even loadNo node saturates while others idlePopular keys, time-ordered writes
Cheap scansRange queries touch few partitionsHashing, which scatters neighbouring keys
Cheap rebalancingAdding a node moves little dataSchemes that tie key placement to node count
Local operationsMulti-key reads and transactions stay on one nodeAny split of related keys

Range partitioning keeps scans cheap and is prone to hot spots. Hash partitioning spreads load and makes range scans touch every partition. Even a perfect hash has a gap, though, because it spreads keys, and the load comes from requests.

06Hot keys and hot shards

6.1How skewed real traffic is

Hashing spreads keys evenly, and it does nothing about requests, because a single popular key always hashes to one place. When one user, product or hashtag gets a large share of traffic, the partition holding it is hot no matter how the keys were spread.

A common model for popularity is a Zipf distribution, where the k-th most popular key gets traffic proportional to 1/k. Simulating that over one million keys, hash-partitioned onto 10 nodes, gives these shares. The last row shows what salting does (section 6.2) when the 100 most popular keys are each split into 8 pieces:

Share of all traffic
Most popular single key6.9%
Fair share per node10.0%
Hottest node16.8%
Coolest node7.4%
Hottest node, top 100 keys each split 8 ways12.4%

One key alone takes 6.9% of all traffic, most of a node's fair share. The hottest node ran at 1.68 times the average (the exact figure shifts with where the hash happens to put the popular keys), and adding nodes can't dilute a single key. At 100 nodes, that one key would still be 6.9% of traffic on one node, about seven times its fair share.

Log-log plot of word frequency against rank for 30 Wikipedia languages; all the curves fall along nearly the same straight line
Zipf's law in real data: how often each word appears in 30 Wikipedias, against its popularity rank, on log scales. On those scales a 1/k law is a straight line sloping down, and all 30 languages fall close to the same line. The top-left corner is the trouble: a handful of words, like a handful of hot keys, carry a large share of all occurrences.Image: SergioJimenez, CC BY-SA 4.0, via Wikimedia Commons

?Why can't the database just give the hot partition more room?

It can, up to the limit of one partition. DynamoDB, Amazon's hosted key-value database, states that every partition can serve at most 3,000 read units and 1,000 write units per second, where a unit is roughly one small read or write. Its adaptive capacity feature can isolate a hot item into a partition of its own, but then that item gets at most those 3,000 reads and 1,000 writes. Past that, only a change to the key design helps.

6.2What to do about a hot key

Since the cure can't be more nodes, every fix either takes load off the key's partition or splits the key itself:

TechniqueHow it worksCost
Cache the readsPut a cache, or an in-process cache, in front of the hot keyStaleness; the cache is now the hot spot
Read from replicasSpread reads over the partition's replicasReplica lag; only helps reads
Split the key (salting)Write to key#0 … key#7 at random; read all eight and combineEvery read fans out 8 ways
Write sharding with aggregationCounters: increment a random sub-counter, sum on read or periodicallyReads are approximate or slower
Give it a dedicated partitionDetect hot keys and move them to their own nodeOperational machinery to detect and move
Bounded-load hashingCap any server at c × average; overflow goes to the next serverSome requests land on a server without their cached data

DynamoDB's own guide recommends write sharding with a random or calculated suffix, which is salting by another name.

Consistent hashing with bounded loads comes from Mirrokni, Thorup and Zadimoghaddam (2016), and HAProxy, a widely used load balancer, ships it as hash-balance-factor. From the HAProxy docs: "if <factor> is 150, then no server will be allowed to have a load more than 1.5 times the average", with 125 to 200 recommended.

So far every lookup has been by key. Real applications also ask for records by other columns, and partitioning by key does nothing for those.

07Queries that don't use the partition key

7.1Local indexes and global indexes

Partitioning by primary key makes lookups by primary key cheap, because the key says which partition to ask. Now suppose a user table is partitioned by user number, and someone asks for the user with a given email address. The user number isn't in the question, so the partitioning can't say where to look. The usual answer is an index, a lookup table from email to user number, and an index on anything other than the partition key is a secondary index. The question is where that index lives, and there are two answers:

Local (per-partition) indexGlobal index
Where it livesEach partition indexes its own rowsThe index is partitioned by the indexed value
Write costOne partition, same transactionA second partition, usually asynchronous
Read by indexed valueAsk every partition (scatter-gather)Ask one index partition, then fetch rows
ConsistencySame as the base tableOften eventually consistent
ExamplesElasticsearch, Cassandra secondary indexes, DynamoDB local secondary indexesDynamoDB global secondary indexes, CockroachDB and Spanner indexes

"Eventually consistent" means a reader may briefly see the old state before the update arrives (chapter 28 covers what that costs you).

?Why is scatter-gather worse than it sounds?

Because a query that asks every partition is as slow as the slowest partition. With 100 partitions, a reply that is slow one time in a hundred turns up in about two queries out of three (1 − 0.99100, roughly 63%). The contention and tail latency chapter shows how fan-out amplifies the slowest few percent of replies, and it applies to every scatter-gather read.

DynamoDB shows the trade-off in its API: global secondary indexes can only be read with eventual consistency, because they're updated asynchronously from the base table.

7.2Keep things that change together together

Operations that touch several partitions need coordination: a distributed transaction, or a saga (a series of local steps with undo steps for failure), or an application that tolerates partial results. The distributed transactions chapter covers those. A cheaper fix is to avoid them by choosing a partition key that keeps related data together.

WorkloadGood partition keyWhy
Multi-tenant SaaStenant_idAlmost every query is inside one tenant
Chatconversation_idMessages are read and written per conversation
Orders and their line itemscustomer_id, with items stored under itAn order and its items commit together
Time-series metrics(series_id, time bucket)Spreads writes across series, keeps one series' range scan local

Redis hash tags (section 2.2), the vindexes of Vitess (a system that partitions MySQL, which we'll meet again in section 9) and the interleaved tables of Spanner (Google's distributed database) are all ways of saying "put these rows on the same shard".

Throughout the chapter we've assumed that a client can find out which node owns a partition. Where that knowledge lives is the next question.

08Routing: who holds the map

8.1Four places to keep the map

Once partitions can move, a client needs a current map from key to node, and the map is a piece of state like any other, with a home. There are four common places to keep it.

Where the map livesHow a request finds its nodeExamples
Smart clientClient caches the map, and servers redirect when it's staleRedis Cluster (MOVED, ASK), Kafka clients, token-aware Cassandra drivers
Proxy tierClients talk to a proxy that routesVitess vtgate, MongoDB mongos, twemproxy, Envoy's Redis proxy
Any node forwardsClient talks to any node, which forwards or coordinatesCassandra coordinators, CockroachDB gateways
Metadata serviceAn authoritative store holds the map; everyone watches itZooKeeper or etcd, Kafka's KRaft controller, CockroachDB meta ranges

A client's cached copy can always be out of date, so every design has to say what happens then. In the smart-client row, the answer is a redirect: the node the client asked says "not mine, try that one", and the client refreshes its map. Section 9.2 shows Redis doing exactly that.

?Why does the map need consensus?

Because two nodes that both believe they own a partition will both accept writes for it, and nothing afterwards can merge those writes safely. The data itself may tolerate being a little stale, but ownership can't have two answers at once. So the map must be agreed on by a group of machines that follow a consensus protocol, and that's why it lives in etcd, ZooKeeper, a Raft group or (for Redis Cluster) behind an epoch number that only increases, so that every change of ownership carries a higher number and nodes ignore claims with an older one. The consensus chapter explains why those are safe, and the failure detection chapter covers what happens when a node wrongly believes it still owns something.

The map has to change whenever data moves, and in a live system data has to move while clients keep using it.

09Moving data without downtime

9.1The general recipe

Adding capacity means moving some partitions to new nodes while clients keep reading and writing them. Let's go back to our six partitions and watch one of them, p1 (keys 1 and 7), move from node 1 to node 4 while traffic continues. The naive way is to stop writes, copy, and restart, which makes the data unavailable for as long as the copy takes. Every system that does this online follows the same shape, with different names for the steps. This is that shape:

Moving partition p1 to a new node while traffic continues
Routing mapwho serves p1Node 1 · sourceChange logwrites after the snapshotNode 4 · targetreads of p1→ node 1writes of p1→ node 1key 1v1key 7v1key 1v1 (copy)key 7v1 (copy)key 7 = v2
Step 1. p1 holds keys 1 and 7 on node 1, and the routing map sends both reads and writes for p1 there. The goal is to move it to node 4 without stopping either.
1 / 8

In a real system the change log is the database's own: MySQL's binlog, Postgres logical replication, or Kafka. The verify step needs a tool built for it. The last step is delayed because, once writes go to node 4, you keep a reverse stream running for a while, sending node 4's changes back to node 1, so that node 1 stays current and you can fail back to it.

Vitess, the MySQL partitioning system from section 7.2, follows this recipe. Its resharding guide maps onto those phases one-to-one: Reshard create copies and streams, VDiff verifies, SwitchTraffic --tablet-types=rdonly,replica moves reads, SwitchTraffic moves writes, ReverseTraffic exists for the way back, and complete cleans up.

?Why switch reads before writes?

Because a read switch is harmless to undo and a write switch isn't. Once writes land on the targets, the source is stale, and failing back needs the reverse stream. Moving reads first lets you find a bad routing rule, a missing index or a cold cache while going back is still free.

9.2Redis Cluster: moving one slot live

Redis Cluster moves data a slot at a time, and it does it without a pause at all, by letting clients be redirected mid-migration. The source is MIGRATING the slot and the target is IMPORTING it. The order of messages matters here, so this is the sequence for slot 1649 moving from node A to node B:

Migrating slot 1649 from node A to node B
ClientNode A (source)Node B (target)SETSLOT IMPORTING / MIGRATINGMIGRATE … KEYSGET k (already moved)-ASK 1649 BASKING + GET kSETSLOT NODE B
Step 1. An operator (or redis-cli --cluster reshard) marks slot 1649 as IMPORTING on B and MIGRATING on A. A still owns the slot.
1 / 6

The decision between those replies is made in a dozen lines of getNodeByQuery():

C
    /* If we don't have all the keys and we are migrating the slot, send
     * an ASK redirection or TRYAGAIN. */
    if (migrating_slot && missing_keys) {
        /* If we have keys but we don't have all keys, we return TRYAGAIN */
        if (existing_keys) {
            if (error_code) *error_code = CLUSTER_REDIR_UNSTABLE;
            return NULL;
        } else {
            if (error_code) *error_code = CLUSTER_REDIR_ASK;
            return server.cluster->migrating_slots_to[slot];
        }
    }

If the command's keys are all missing from this node, the reply is ASK with the target's address. The middle case is a multi-key command whose keys are split between the two nodes. It gets TRYAGAIN, an error, for as long as the migration is in progress, because neither node can run the whole command.

Here's the whole thing on a two-node cluster. cluster keyslot asks which slot a key hashes to. cluster setslot 1649 importing $A and migrating $B mark the two ends of the migration, where $A and $B stand for the two nodes' ids (each node reports its own with cluster myid). migrate … keys moves a batch of keys, and cluster getkeysinslot supplies the list. M=<one of the migrated keys> is a placeholder for any key that the migrate call moved, and -c makes redis-cli follow redirects the way a smart client would.

Migrate half of one slot, then see ASK, MOVED and TRYAGAIN territory
shell
Shell
redis-cli -p 7291 cluster keyslot '{user:1000}:cart'        # 200 keys live in this slot
redis-cli -p 7292 cluster setslot 1649 importing $A
redis-cli -p 7291 cluster setslot 1649 migrating $B
redis-cli -p 7291 migrate 127.0.0.1 7292 "" 0 5000 keys $(redis-cli -p 7291 cluster getkeysinslot 1649 100)
 
M=<one of the migrated keys>
redis-cli -p 7291 get "$M"                   # a moved key, asked at the source
redis-cli -p 7292 get "$M"                   # same key at the target, no ASKING
redis-cli -p 7291 get '{user:1000}:k1'       # a key that hasn't moved yet
redis-cli -p 7291 -c get "$M"                # -c follows redirects
# ... migrate the rest, then SETSLOT 1649 NODE $B on both nodes
redis-cli -p 7291 get '{user:1000}:k1'
output
Output
1649
ASK 1649 127.0.0.1:7292
MOVED 1649 127.0.0.1:7291
v
v
MOVED 1649 127.0.0.1:7292

Read the output line by line, one per get or keyslot command. The first line is the slot, 1649. The moved key, asked at the source, gets ASK: node A owns the slot but no longer has that key. The same key asked at the target without ASKING gets MOVED back to A, because B doesn't officially own the slot yet. The key that hasn't moved is answered by A as usual (v), and -c follows the ASK and ASKING for us, so the fifth line is the value too. After SETSLOT NODE, A answers MOVED for everything in the slot, so the last line points at B.

9.3How fast to move

Moving data competes with serving traffic for disk, network and CPU on both sides. Every mature system throttles it: Cassandra's stream_throughput_outbound, CockroachDB's rebalance snapshot rate settings, Kafka's replication quotas for partition reassignment.

?Why not rebalance automatically the moment a node looks dead?

Because "looks dead" is a guess, and moving its data is expensive. If a node is merely slow and the cluster starts copying its partitions elsewhere, the copying loads the other nodes, which can make them look slow too. That's why Kubernetes, which runs containers across a fleet of machines, waits five minutes by default before moving work off a node it can't reach. It's also why Cassandra never moves data on its own when a node fails. It marks the node down, other nodes save the writes it misses (Cassandra calls these hints), and when the node returns it catches up from the hints and from repair, a background comparison of the replicas. That guess is the subject of the next chapter.

That leaves the practical side: on a running cluster, how do you tell whether any of the problems from this chapter is happening to you?

10Running a partitioned system

10.1What to watch

The commands below are grouped by the question they answer, with the section that raised it. Output depends on your cluster. One needs a word of setup: redis-cli --hotkeys only works when Redis's eviction policy, the rule it uses to choose keys to drop when memory is full, is an LFU (least frequently used) one, because that policy is what makes Redis count how often each key is read.

Shell
# Is one node carrying more than its share of keys? (sections 1 and 3)
redis-cli --cluster check HOST:PORT     # slots and keys per node, and any slots left mid-migration
nodetool status                         # Cassandra: the Owns and Load columns per node
 
# Is one key taking too much traffic? (section 6)
redis-cli --hotkeys                     # needs an LFU maxmemory-policy
 
# Who owns this key, and is a move in progress? (sections 2, 8 and 9)
redis-cli cluster keyslot KEY           # which slot a key hashes to
redis-cli cluster nodes                 # who owns which slots; look for migrating and importing
nodetool netstats                       # Cassandra: data streaming between nodes right now

Beyond commands, these are the quantities worth plotting:

  • Per-partition size and request rate, not just per node. A node average hides one hot partition. Plot the busiest partition against the median.
  • The skew ratio: busiest node ÷ mean, for disk and for requests. Above about 1.5, find out why before you add hardware.
  • Top keys. redis-cli --hotkeys (with the LFU policy above), DynamoDB's CloudWatch Contributor Insights, or sampled request logs.
  • Redirects and retries. A steady rate of MOVED, stale-routing errors or Kafka's NotLeaderOrFollower means clients hold a stale map.
  • Movement in progress. Bytes remaining in rebalancing and its throttle, so a slow migration is visible before it's a surprise.

10.2Rules that hold up

  1. Separate the key-to-partition function from the partition-to-node map, and never change the first (section 2).
  2. Pick a partition count that covers years of growth wherever the count is fixed at creation, as in Kafka and Elasticsearch (section 2.3).
  3. Never use hash % N with N as the number of nodes when nodes come and go (section 1).
  4. Give ring nodes many tokens, or use a token allocator (section 3).
  5. Choose a range-partitioned key by looking at write order. A leading timestamp makes one node take all the inserts (section 5.3).
  6. Salt only keys that measurement shows are hot (section 6.2).
  7. Verify before you switch when you move data online, and switch reads before writes (section 9.1).

10.3What you trade for what

You getYou payWhen the bill arrives
Fixed partitions: cheap, explicit rebalancingA count you can't change, and a map to store and keep correctWhen the cluster outgrows the count
Hashing: even spread of keysRange scans touch every partitionOn the first query that scans
Ranges: cheap scansHot spots on ordered keysAt peak write load
A ring or jump hash: no central mapUneven shares, or the limit of only adding at the endAs one node filling its disk first
Local secondary indexes: cheap writesScatter-gather readsIn the tail latency of every index query
Global secondary indexes: cheap readsSlower, often eventually consistent writesWhen a reader sees a stale index
Online resharding: no downtimeA copy, a log and a verification to runAs migration load on a busy cluster

10.4Symptom, cause, fix

SymptomLikely causeFix
One node's disk fills firstToo few tokens per node, or uneven partition sizesMore vnodes or a token allocator; split big partitions
One node hot on writes; others idleTime-ordered keys in a range-partitioned storeHash-prefix the key, or lead with a high-cardinality column
One partition throttled; table far under capacityA hot keyCache, salt that key, or give it a dedicated partition
Adding a node triggers a huge data transferhash % N placementFixed partitions or consistent hashing
Per-key ordering broke after scaling a Kafka topicAdded partitions to a keyed topicCreate a new topic with more partitions and migrate consumers
Index queries slow, p99 worstScatter-gather over every partitionGlobal index, or a partition key that matches the query
CROSSSLOT or TRYAGAIN errors in RedisMulti-key command over keys in different slots, or mid-migrationHash tags; retry TRYAGAIN with backoff

("p99" is the latency that 99% of requests beat, so it measures the slowest few.)

11Summary

  1. Every scheme answers two questions: key to partition, and partition to node. Keeping them separate is what makes growth cheap, because the first never changes and the second is a small map.
  2. hash % N moves almost everything. Going from 10 to 11 nodes moved 90.9% of keys in the simulation, and in our twelve-key example 9 of 12 moved, where the ideal is only the new node's share.
  3. Many fixed partitions are the simplest fix. Redis's 16,384 slots and Instagram's logical shards never change a key's partition. Kafka and Elasticsearch fix the count at creation, so size them for years of growth.
  4. A ring moves keys only to the new node, and needs virtual nodes. With one token per node the busiest node held 2.8 times its share, and with 256 tokens, about 1.1 times.
  5. Range partitioning keeps scans cheap and invites hot spots. Time-ordered keys send every write to one node, and splitting doesn't help.
  6. Hashing spreads keys, not requests. One Zipf-popular key took 6.9% of all traffic, and no number of nodes dilutes it.
  7. Partitions have a ceiling. In DynamoDB it's 3,000 reads and 1,000 writes a second, even for a partition holding one isolated hot item.
  8. A secondary index is either local or global. Local indexes make reads scatter-gather over every partition, and global indexes make writes slower and often eventually consistent.
  9. The partition map needs consensus. Two nodes that both believe they own a partition will both accept writes.
  10. Online resharding is copy, stream, verify, switch reads, switch writes, clean up. Only the write switch pauses, and verifying is the step people skip.
  11. Redis Cluster redirects mid-migration with ASK, and afterwards with MOVED. A client that confuses them ping-pongs.

12Build this

A key-movement lab.

  • Implement % N, a ring with configurable virtual nodes, jump hash and rendezvous hashing over the same set of keys. Reproduce the table in section 4.2, then add node removal: remove node 3 of 10 and count moved keys. Jump hash can't do it directly; work out what you'd need to add.
  • Add weights: make one node twice as big and check that each scheme gives it twice the keys. Note which schemes make that easy.
  • Start two Redis Cluster nodes and a writer that increments a counter under {user:1} 1,000 times a second. Migrate the slot by hand, as in section 9.2, and count errors the writer sees with redis-cli -c and with your usual client library.

13Interview questions

beginnerWhat's the difference between partitioning and replication?›

Replication keeps copies of the same data on several nodes, for durability, availability and read capacity. Partitioning splits the data so different nodes hold different subsets, for capacity and write throughput. Real systems combine them: each partition is replicated, so every node holds replicas of several partitions.

beginnerWhy is hash(key) % N a bad way to assign keys to nodes?›

Because changing N remaps almost every key. Going from 10 to 11 nodes changes the owner of about 10/11 of the keys, so adding 10% capacity moves roughly 91% of the data. Consistent hashing or a fixed number of partitions moves only the share the new node should own.

intermediateRange or hash partitioning for a table of events keyed by timestamp?›

Hash, or a range scheme with a hashed or high-cardinality prefix. Plain range partitioning on a timestamp sends every new write to the last range, so one node takes all inserts however many you add. If you need time-range scans, lead with something like (device_id, timestamp) or a hash bucket, and accept that a global time scan fans out.

intermediateWhat are virtual nodes and why do consistent-hash rings need them?›

Each physical node owns many positions (tokens) on the ring instead of one. With one token per node, arc lengths vary a lot: in the simulation the busiest of 10 nodes held about 2.8 times its share. Many tokens average out the arcs (about 1.1 times at 256), and when a node leaves, its load spreads across many nodes instead of landing on one neighbour.

intermediateOne celebrity's account overloads a shard. What are your options?›

Cache reads for that key; spread reads across replicas; split the key into N sub-keys (salting) and fan reads out; for counters, increment random sub-counters and sum them; or detect hot keys and give them a dedicated partition. Adding shards doesn't help, since one key always hashes to one place. Salt only the keys that are measurably hot.

deepIn Redis Cluster, what's the difference between MOVED and ASK?›

MOVED means the slot now belongs to another node permanently: update your slot map and send future requests there. ASK means the slot is mid-migration and this particular key has already moved: send this one request to the target, preceded by ASKING, and don't update the map. The target refuses requests for an IMPORTING slot without ASKING, answering MOVED back to the source.

deepWalk through resharding a live MySQL-backed service from 4 shards to 8.›

Create the 8 target shards. Copy a consistent snapshot of each source shard's rows into the targets by the new key ranges, recording the binlog position. Stream binlog changes from that position into the targets until lag is seconds. Verify with a row-level diff. Switch replica reads, then primary reads. Stop writes on the sources briefly, let the stream drain, flip the routing map to the targets and resume. Keep a reverse stream so you can fail back, then drop the old shards. Vitess automates exactly this as Reshard, VDiff, SwitchTraffic and complete.

deepWhy must the partition map be strongly consistent when the data isn't?›

Because two nodes that each think they own a partition will each accept writes, and nothing afterwards can merge them safely. The data may tolerate staleness; ownership can't. So the map lives in a consensus system (etcd, ZooKeeper, a Raft group, KRaft) or behind a monotonic epoch, and data nodes check the epoch on every request so an old owner's writes are rejected.

14Go deeper

check yourself
What does & 0x3FFF do in Redis's keyHashSlot?›

It's modulo 16384, the fixed number of hash slots. It works as a bit mask because 16384 is a power of two.

Which keys move when a node joins a consistent-hash ring?›

Only keys on the arcs the new node's tokens cut off, and they all move to the new node. No key moves between two existing nodes.

Why can't jump consistent hash remove node 3 of 10 directly?›

Its buckets are the numbers 0 to N−1, and it can only grow or shrink N at the end. Removing a middle bucket needs an extra mapping layer on top.

A Redis MULTI touching two keys in a slot that's half migrated. What happens?›

TRYAGAIN, if the source holds some of the keys and not others. Retry after the migration step completes.

Karger et al., Consistent Hashing and Random Trees (STOC 1997)

The original consistent-hashing paper, written for web caching. PDF

Lamping and Veach, A Fast, Minimal Memory, Consistent Hash Algorithm

Jump consistent hash in five lines, with the proof that it's balanced and minimally disruptive. arXiv 1406.2294

Mirrokni, Thorup and Zadimoghaddam, Consistent Hashing with Bounded Loads

Capping each server's load at c × average while staying consistent. arXiv 1608.01350

Redis Cluster specification

Slots, hash tags, MOVED and ASK, and the live resharding procedure, from the people who built it. redis.io

Vitess: Resharding

Online MySQL resharding, phase by phase, with the commands. vitess.io

DeCandia et al., Dynamo (SOSP 2007), section 6.2

Three partitioning strategies Amazon tried in production, ending with fixed equal-sized partitions assigned to nodes. PDF

Replication & Consistency Models

What each partition's replicas promise to readers, and quorum arithmetic. Chapter 28.

Failure Detection & Membership

How the cluster decides a node is gone before it moves that node's partitions. Chapter 30.

Redis Internals

The single-threaded server underneath each Redis Cluster node. Chapter 22.

Distributed Transactions

What to do when an operation has to span partitions after all. Chapter 31.