Distributed systems
Part 4 of 6 · Consistent hashingVnode Rebalancing & Membership: Stream ~K/N Without Split-Brain
Adding a node only steals ~1/N of keys, but streaming those keys still needs throttling, versioned membership, and hinted handoff so clients and replicas do not split-brain.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Overview
The hub lesson stops at ownership: add a node, ~1/N of keys change primary. This page is the ops algorithm. Those keys live in SSTables, WAL segments, or cache entries. You must detect membership change, compute a small set of vnode ranges to move, stream them without saturating disks, forward in-flight writes, and publish a new generation so clients do not keep sending the stolen keys to the old owner.
Skip this layer and a theoretically perfect ring still split-brains during every scale event.
Why controlled vnode rebalance wins
| Approach | What moves | What goes wrong |
|---|---|---|
| Modulo reshape | ~all keys | Origin / disk stampede |
| Jump hash N → N+1 | ~1/(N+1) keys, all to bucket N | No named nodes; cannot remove a middle shard; still must stream |
| HRW add id | ~1/(N+1) keys, whatever the new max wins | No ranges; you still copy objects; O(N) lookup |
| Naive “add node, copy everything it might own” | Far more than ~1/N | Extra disk, extra time, extra split-brain window |
| Vnode steal + throttle (winner) | ~1/N keys, batched ranges | Requires gossip, generations, write-forward |
- 1
hash function → new owner id
Jump / HRW assignment only
Tells you who should store the key. Does not stream bytes, cap NIC use, or stop clients on generation 12 talking to a node that already dropped the range.
- 2
join → disk / NIC melt
Unthrottled ring fill
Even ~K/N of a 100 TB cluster is terabytes. Ten nodes joining at once without a concurrency cap is a self-inflicted outage.
- ?
Winner: vnode steal, versioned membership, throttled streams
New vnodes carve arcs from predecessors. Stream those ranges in batches. Bump a generation. Forward writes until cutover. Maglev tables rebuild in RAM — the wrong tool for stateful range handoff.
Core concepts
| Term | Meaning on this page |
|---|---|
| Vnode | A token (point on the ring) owned by one physical node. Each node holds many; typical teaching V is 100–256. |
| Token ring | Sorted tokens; key hash → clockwise owner. See the hub for lookup. |
| RF | Distinct physical replicas per key. Placement details: topology. |
| Gossip membership | SWIM-style periodic push of alive/suspect/dead, tokens, load, generation. |
| Generation / schema version | Monotonic epoch of the token map. Coordinators reject or retry on mismatch. |
| Controlled rebalance | Pick a minimal set of vnodes to move; cap max concurrent moves and bytes/sec. |
| Range handoff | Stream data for a stolen arc (SSTables / log), write-forward new mutations, then flip ownership. |
| Hinted handoff | Neighbor stores writes for a temporarily down replica and replays on return. Not a substitute for RF. |
Vnodes exist so a join steals many small arcs from many peers (failure fan-out in reverse: parallel bootstrap) instead of one enormous range from one neighbor.
Join, leave, fail
Decisions
- 1
Node joins
- nextGossip: alive + empty token claim
- 2
Gossip: alive + empty token claim
- nextPlanner assigns vnode tokens
- 3
Planner assigns vnode tokens
- nextRebalance needed?
- ?
Rebalance needed?
- YesSelect vnodes to steal from current owners
- NoIdle
- 5
Select vnodes to steal from current owners
- nextStream ranges in batches of maxConcurrentMoves
- 6
Stream ranges in batches of maxConcurrentMoves
- nextWrite-forward in-flight puts to src and dst
- 7
Write-forward in-flight puts to src and dst
- nextFlip token map, bump generation
- 8
Flip token map, bump generation
- nextClients refresh ring
- 9
Clients refresh ring
- 10
Idle
- 11
Node fails
- nextMark suspect then dead
- 12
Mark suspect then dead
- nextHints on neighbors
- 13
Hints on neighbors
- nextRF restorations if the death is permanent
- 14
RF restorations if the death is permanent
Lesson map
Vnode Rebalancing & Membership: Stream ~K/N Without Split-Brain
Adding a node only steals ~1/N of keys, but streaming those keys still needs throttling, versioned membership, and hinted handoff so clients and replicas do not split-brain.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB a["Node joins"] b["Gossip: alive + empty token claim"] c["Planner assigns vnode tokens"] d["Rebalance needed?"] a -->|Node joins to Gossip: alive + empty token claim| b b -->|Gossip: alive + empty token claim| c c -->|Planner assigns vnode tokens| d
Add (scale out)
- New node appears in gossip as alive, generation pending.
- Allocator places V new tokens (random or balanced). Each token steals the open-open arc from the previous clockwise token’s owner — the same rule as the hub.
- Expected key share for the newcomer: ~1/(N+1) if weights are equal. Expected source fan-out: many physical peers, not one.
- Stream only those arcs, in batches of size
maxConcurrentMoves(teaching default 8), with a bytes/sec cap. - During the stream, write-forward: coordinators that still see the old generation send to the old owner; the old owner replicates into the stream or into a log the newcomer tails.
- Cut over: publish generation G+1 with the new tokens. Old generation readers retry.
Remove / decommission
The leaving node’s vnodes are given to remaining nodes (often the clockwise neighbors, or a planner that preserves balance). Stream out, then tombstone the node. This is the inverse of a join, still ~1/N of keys if one of N equal nodes leaves.
Fail
Gossip marks suspect then dead. Lookups walk to live replicas. Hints absorb writes for a short TTL. If the node is gone for good, the planner restores RF by streaming to a new replica — a second ~1/N-class copy, not “the ring magically has RF again.”
Worked numbers
- 1M keys, 10 nodes, V = 128 vnodes each → 1280 tokens. Add 11th: ~91k keys, about 128 new tokens stolen from ~10 peers (some peers donate several arcs).
maxConcurrentMoves = 8: at least 16 batches if you serialize perfectly; wall clock isceil(128/8)stream slots plus tail.- 3 nodes join at once without a global cap: up to ~3/N of data in flight. Cluster-level
max_concurrent_movesexists so bootstrap does not outrun compaction.
Membership is a protocol, not a hashmap
| Mechanism | Role | Failure mode if skipped |
|---|---|---|
| Gossip (SWIM) | Detect alive/dead in O(log N) rounds | Long partitions look like mass death; flapping causes stream storms |
| Generation / epoch | Fence stale coordinators | Split-brain writes on stolen keys |
| Token map in gossip | Everyone computes the same successor | Divergent rings |
| Repair / anti-entropy | Reconcile missed streams | Silent RF holes after a failed handoff |
| Hinted handoff | Short outages | Hints pile up and replay as a thundering write storm — bound the hint window |
Jump hash’s “state is N” still needs a version of N and a side table if buckets map to hosts. HRW still needs a version of the node list. The vnode ring is just the membership that also names ranges so streaming has a unit of work.
Decommission, repair, and SWIM suspects
Decommission (planned leave) is not a crash:
- Mark the node
leavingin gossip so coordinators stop placing new primaries there. - Planner assigns each of its V tokens to remaining nodes (often the clockwise neighbor, or a balancer that preserves disk share).
- Stream out under the same
maxConcurrentMovesbudget as a join. - Only then remove the node and bump the generation. If you flip ownership before the stream finishes, you read empty ranges.
Repair / anti-entropy is how you notice a stream that said done but missed files. Merkle trees per vnode range (Dynamo/Cassandra) compare replica digests and stream differences. After every failed bootstrap, repair the stolen ranges before you trust RF.
SWIM / gossip suspects: a single missed heartbeat is not death. Nodes go alive → suspect → dead with an infection-style round. Too-aggressive timeouts cause flapping: a node is stolen from, then returns, then you stream twice. Too-slow timeouts leave RF-1 in production. Pair suspect timers with the topology question: during suspect, do you still wait for that replica in W?
Write path during a batched stream. Batch 0 owns tokens 0–7 of the newcomer, batch 1 owns 8–15. If you publish one generation at the end, coordinators keep using the old owner for all 32 vnodes until cutover — write-forward has to cover the whole join. If you publish a generation per batch, clients refresh more often but each stolen arc has a shorter dual-write window. Say which you built.
Client refresh
Drivers cache the token map. Options:
- Push: gossip or a control topic announces generation G+1; drivers swap atomically.
- Pull on error: coordinator returns
STALE_GENERATION; driver fetches the map and retries the same request. - Dual-read during a planned window for caches (both old and new owner, prefer new).
A mobile client that sleeps for an hour is the hostile case: it must fail on generation mismatch rather than write into a decommissioned vnode.
Playground — rebalance counts, batches, generation
In-memory ring. No sockets. Logs keys moved, vnode steals, and batch count against ~1/N.
Press Run. Snippets must be self-contained — no network, files, or native modules.
N=3, V=4, add node D. Draw 12 old tokens, 4 new ones, and arrows from predecessor owners. Box a maxConcurrentMoves=2 schedule (2 batches). Stamp the picture gen 7 → gen 8. Show a client still on gen 7 writing a stolen key — and where the write must go until cutover.
Interview Q&A
Adding a node only moves ~1/N keys. Why is rebalancing still hard?
Answer
Because ~1/N of a large store is still terabytes, streams compete with live traffic, and two membership generations can both accept writes for the same key. The hash bound is necessary, not sufficient.
What does a vnode give you during a join that a single-token node does not?
Answer
Many small arcs stolen from many predecessors, so bootstrap is parallel and no single neighbor donates its entire disk. Same fan-out story as failure, in reverse. Cost: bigger token map, more streams to schedule.
1M keys, 10 equal nodes, V=128, add an 11th. What numbers do you quote?
Answer
Keys moved ~91k. New tokens 128. maxConcurrentMoves=8 ⇒ 16 stream slots if each token is one move. Donors: up to all 10 existing nodes. Generation bumps once at cutover (or per batch if you publish incrementally — say which).
What is write-forwarding versus hinted handoff?
Answer
Write-forwarding is for a planned range move: old owner (or coordinator) duplicates mutations to the newcomer until the generation flips. Hinted handoff is for a down replica: a neighbor keeps a timed hint and replays it. Hints are not how you bootstrap a node, and they are not RF.
How do clients avoid split-brain during a rebalance?
Answer
Carry a membership generation. If the coordinator’s generation is behind, refresh and retry. Optionally dual-write during a planned window. Never silently accept a put on a token you already streamed away.
Jump hash has no ring. Do you still need this page?
Answer
Yes for stateful buckets: N → N+1 still moves ~K/(N+1) keys onto bucket N, and those bytes must stream with a version of N. Jump hash also cannot remove a middle bucket, so decommission is a drain — still a membership protocol. Jump lesson.
HRW has no tokens. What is the unit of streaming?
Answer
Keys (or key ranges you invent). The new node wins whatever scores it maxes; you list those keys from a metadata index or scan. You lose “stream this 1/V of the ring” as a neat unit. That is why storage systems keep vnodes even when they know about HRW.
Three nodes join at once. What melts?
Answer
Without a cluster-wide concurrency and bandwidth cap, you launch O(joins · V) streams. Compaction, repair, and the write path stall. Planners serialize bootstraps or share one max_concurrent_moves budget.
A node dies. Where do writes go, and when is RF restored?
Answer
Live replicas (RF walk skips the dead host). Neighbors may take hints. RF is restored only when a new replica is streamed. Until then you are on RF-1 for keys that listed the dead node — topology decides whether the remaining copies share an AZ.
Why gossip instead of a single config file?
Answer
N storage nodes plus client fleets cannot atomically read one file. Gossip spreads alive/dead and tokens in logarithmic rounds. You still want a fence (generation, or a small consensus service for the map) so gossip lag does not dual-own a range.
Pitfalls
- Shipping ~K/N without a throttle — theoretically small, operationally huge.
- No generation — split-brain on stolen keys.
- Counting vnodes as replicas while streaming — RF walk must skip the same physical node.
- One token per node — join/fail dump onto a single neighbor.
- Hint backlog — replay storms; bound TTL and size.
- Parallel joins sharing no global budget.
- Forgetting repair after a failed stream — RF looks fine in gossip, data is not there.
- Hot keys on a range you are streaming — the victim is both the donor disk and the newcomer. Bounded loads.
Deep dive · Cassandra / Dynamo vocabulary
Dynamo described consistent hashing, preference lists, hinted handoff, and anti-entropy. Cassandra’s vnode era (num_tokens) made token assignment automatic so operators stopped picking one token per box. Modern guidance often lowers vnode counts and uses allocation algorithms that place tokens more carefully than “hash node#i.” The interview still wants: many tokens, steal small arcs, stream with a cap, version the map.
Scylla-style shard-aware drivers pin connections to cores that own the vnode — a client optimization on top of the same membership. DynamoDB hides the partition map behind a service; you still feel partition splits as a rebalance, just not one you schedule with nodetool.