Vnodes flatten the lumpy ring

Many small ring positions per physical server make arc sizes uniform by averaging, and spread a dead node's load across all survivors instead of crushing one neighbour.

Previously

The ring is great in theory, but with one position per server those arcs were wildly unequal — and worse, a dead server dumped its entire arc onto a single clockwise neighbour. We need to chop each server's stake into many small pieces.

Scene 05

Vnodes flatten the lumpy ring

  1. Watch
  2. Try it
  3. Predict
  4. Capture
mode: vnodesload variance: ±101% (ideal 0%)4 tok/nodeS0S1S2S3S4S5S6S7ARC SIZES (deg)S037°S160°S259°S334°S415°S556°S652°S747°ideal: 45° · range: 15–60°4 vnodes per physical server — same 8 servers, 32 total ring positions.
↑ vnode — a virtual ring position (NOT a VM)
← variance shrinks as vnodes/server goes up
What to watch for

Same 8 servers as scene 4 — but each one now holds 4 small wedges scattered around the ring. The variance bar up top is what to watch: arcs are far closer to the ideal 45° share, and no single neighbour is on the hook for a whole server's worth of keys anymore.

Implementation

Highlighted lines are the ones running in the diagram right now.

tokensForNode
deterministic vnode positions per physical server
def tokensForNode(nodeId, vnodeCount):
out = []
for i in range(vnodeCount):
# pure arithmetic, no rng — same answer every render
t = ((nodeId * 1009 + i * 2017) * 137) mod 360
out.append(t)
return sorted(out)
# ring tokens = union of tokensForNode(n, k) for every node n
loadVariance
variance of arc sizes shrinks as vnodeCount grows
def loadVariance(nodes, vnodeCount):
arcs = [] # (owner, size) per ring slice
for token in sorted(all_tokens(nodes, vnodeCount)):
arcs.append((owner(token), arc_size(token)))
per_node_total = sum_by_owner(arcs)
return stddev(per_node_total) / mean(per_node_total)
# law of large numbers: each node's stake is a sum of
# vnodeCount random arcs, so variance ~ 1 / sqrt(vnodeCount).
# 1 -> ~50% ; 4 -> ~20% ; 64 -> ~5%.
absorbDeadNode
where a dead server's load goes
def absorbDeadNode(dead):
for vnode in dead.tokens:
# each greyed wedge hands its keys to its own
# clockwise successor — a different physical node
# than the previous wedge, when vnodeCount is high.
successor = ring.successor(vnode + 1)
successor.adopt(vnode.keys)
# 1 token -> 1 successor crushed
# 64 tokens -> ~64 successors each take 1/64 of the load

Where this sits in Build a wide-column store (Cassandra / DynamoDB family)

Scene 05 of 13, in the The ring act — Consistent hashing + vnodes — adding a server moves only ~1/N of keys.. Give each physical server many small ring positions; arcs become uniform and deaths spread their load.

Up next. Vnodes spread the load evenly, but the dead node's data is still lost — we need copies.

All 13 scenes in Build a wide-column store (Cassandra / DynamoDB family) · Every curriculum

Built with Arqly
Every scene in Build a wide-column store (Cassandra / DynamoDB family) builds on the one before it.All 13 Build a wide-column store (Cassandra / DynamoDB family) scenes