Designing Data-Intensive Applications
Ch. 6

Rebalancing Partitions

When nodes join or leave the cluster, partitions must move — without excessive downtime.

Adding a machine should improve capacity, not invalidate your entire partition map. Naive hash mod N remaps almost every key when N changes. Better schemes use fixed partitions or consistent hashing to limit data movement.

In practice

Cassandra uses virtual nodes (vnodes) so adding a host only moves a fraction of data. Kafka's kafka-reassign-partitions tool migrates partitions between brokers with throttling. On Kubernetes, scaling a StatefulSet adds pods but your app still needs a safe rebalance plan — automatic resharding without ops oversight has caused production outages.

Netflix at scale

Adding Cassandra nodes during peak streaming hours uses consistent hashing with vnodes so only a fraction of keys move — avoiding the full remap disaster of hash mod N.

typescript — Consistent hashing on the ring
// Netflix / Cassandra — consistent hashing limits key movement on rebalance
function hashRing(key: string, nodes: string[]): string {
  const ring = nodes.flatMap((n) =>
    Array.from({ length: 128 }, (_, i) => ({ hash: hash(`${n}:${i}`), node: n }))
  ).sort((a, b) => a.hash - b.hash);
  const h = hash(key);
  return ring.find((e) => e.hash >= h)?.node ?? ring[0].node;
}
// Adding one node moves only ~1/N of keys — not hash(key) % N
Key Takeaways
  • Hash mod N breaks when N changes — most keys remap to new nodes.
  • Consistent hashing minimizes key movement when nodes are added or removed.
  • Fixed number of partitions with virtual nodes simplifies rebalancing.
  • Automatic rebalancing risks cascading failures; manual triggers are safer.
  • Rebalancing should move minimal data at bounded throughput.
  • Cassandra vnodes, Kafka partition reassignment, and Kubernetes StatefulSet scaling all rebalance data.
rebalancingCassandraKafkaKubernetesconsistent hashingvirtual node