Virtual Nodes

With one position per shard, a shard’s share of the codes depends on the length of its arc. We do not choose the positions. They come from hashing the shards’ names.

Uneven shares

With A at 20, B at 55, and C at 85, the arcs are close in length. A holds 35 of the 100 hash values, B holds 35, and C holds 30.

Now add D at 40. D takes 20 of B’s 35 values:

Shard Before adding D After adding D
A 35 35
B 35 15
C 30 30
D 20

D took its whole share from B. A and C hold as much as before. If A was nearly full, adding D did not help A. B now holds less than half of what A holds.

Removing a shard has the same problem. Go back to the three shards A, B, and C, and remove B. All of B’s 35 values go to C. C then holds 65 of the 100 values, and A holds 35.

Several positions per shard

Instead of one position per shard, we can give each shard several. We hash a different name for each position, such as A-1 and A-2. Suppose each shard gets two positions. A-1, B-1, and C-1 land where A, B, and C were before. A-2 lands at 38, B-2 at 5, and C-2 at 45.

Each code still goes to the first position clockwise. That position’s shard holds the code. Here are the positions and the codes, going clockwise from 0:

Code or shard name Position Stored on
B-2 5
x7Kp2Qa 17 A
A-1 20
Qz81fWc 30 A
A-2 38
C-2 45
b3Rt9Lm 52 B
B-1 55
Hn4vR8e 71 C
C-1 85
pL2wJ5s 96 B, at B-2 after passing 99

Each shard now holds two arcs. A holds the codes from 6 to 38, which is 33 values. B holds the codes from 86 to 5, after passing 99, and from 46 to 55. That is 30 values. C holds the codes from 39 to 45 and from 56 to 85. That is 37 values.

Now remove B again. The codes on B’s arc that ends at 55 go to the next position clockwise. That position is C-1 at 85. The codes on B’s arc that ends at 5 go to A-1 at 20. So b3Rt9Lm moves to C, and pL2wJ5s moves to A. A ends up with 53 of the 100 values and C with 47. When each shard had one position, removing B left C with 65.

Adding a shard works the same way. Suppose D gets two positions, at 32 and 75. D at 32 takes the codes from 21 to 32 from A, including Qz81fWc. D at 75 takes the codes from 56 to 75 from C, including Hn4vR8e. So D takes codes from two shards instead of one.

A server in a distributed system is often called a node. On the ring, a shard with several positions appears as several nodes. Each of those positions is called a virtual node.

Two virtual nodes per shard are still too few to make the shares even. Real systems use more. Apache Cassandra, for example, gives each server 16 virtual nodes by default. With many virtual nodes spread around the ring, each shard’s arcs add up to about the same share. A new shard typically takes codes from several other shards. A removed shard’s codes are typically spread across several remaining shards.

Shards of different sizes

Virtual nodes also let us give a shard a share of the codes that matches the storage on its servers. Suppose we add shard E on servers with twice the storage of the other shards’ servers. We give E twice as many virtual nodes as each of the other shards. With many virtual nodes, E then holds about twice as many codes as each of the others.