skip to content

When nodes are added to or removed from a partitioned cluster, walk through how rebalancing actually moves data. What are the main strategies (e.g., a fixed number of partitions per node vs. dynamically splitting partitions by size), and what can go wrong during the process?

level: seniorimportance: should knowfreq 65%

answer

  1. assignment vs data movement are separate problems
  2. fixed partition count vs dynamic split/merge
  3. stream then atomic cutover
  4. throttle concurrent moves (thundering herd)
  5. Cassandra one-at-a-time bootstrap, HBase region split

basics

~20 s

Rebalancing is moving data between nodes so each one holds a fair share after the cluster changes size. One approach creates many more partitions than nodes up front and just reassigns whole partitions to nodes as they join/leave; another splits/merges partitions dynamically as they grow or shrink. Done badly, it can overload the network or serve stale/missing data mid-move.

solid answer

~60 s

Rebalancing has two sub-problems: deciding the new assignment of partitions to nodes, and physically moving the data to match. A common strategy is to fix the number of partitions well above the node count at cluster creation and let partitions move as whole, unsplit units between nodes as membership changes — simple to implement, but requires guessing the right partition count upfront since partitions don't split further, capping how thin data can eventually spread. The alternative is dynamic partitioning: partitions split when they exceed a size/load threshold and merge when they shrink, so partition count grows with data volume — needs a live coordinator to decide when/how to split, more complex, but adapts to unpredictable growth (used by HBase/Bigtable-style range-partitioned systems). Regardless of strategy, moving data involves streaming it from source to destination while the source keeps serving reads/writes, then atomically switching ownership/routing once the destination catches up — get this wrong and you get thundering-herd network/disk saturation if too many partitions move at once, or a window where routing metadata is stale and requests hit the wrong owner.

go deeper

for a junior

Should understand at a basic level that rebalancing means moving data so nodes stay balanced after cluster changes, and that this takes time and resources, without needing to name specific strategies.

for a middle

Should be able to describe at least one concrete rebalancing strategy (fixed partition count or dynamic splitting) and know that data streaming happens before ownership switches.

for a senior

Should contrast both major strategies, explain the stream-then-cutover mechanism and why it's safer than a single-shot copy, and identify the thundering-herd and stale-routing failure modes.

for a principal

Should reason about capacity planning for partition-count choice at system-design time, evaluate throttling/sequencing policies for large-scale rebalances, and cite operational practices from real systems as evidence of how these failure modes are mitigated in production.

## What rebalancing is Rebalancing is the process of moving data between nodes so that each node in a partitioned cluster ends up with a proportionate share after the cluster's membership changes: - a node is added to relieve load, - a node is removed for maintenance, - or a node fails and needs its data reconstructed elsewhere. It has two logically separate parts worth pulling apart: 1. **The assignment decision** — which partitions should live on which nodes now? 2. **The data-movement mechanism** — how do bytes actually get from the old owner to the new owner without breaking availability or correctness along the way? ## The two assignment strategies The assignment decision is where the two dominant strategies diverge. **The fixed-partition-count approach** decides, at cluster creation time, on a partition count well above the number of nodes ever expected — a common rule of thumb is roughly 10x the eventual maximum node count — so each node initially owns many partitions, and as nodes are added, whole partitions get reassigned from existing nodes to the new one. Because partitions are moved as indivisible units and never split further, the assignment algorithm is comparatively simple: essentially a load-balancing/bin-packing problem over a fixed set of items. The real risk is choosing the wrong partition count upfront — too few, and you can't spread data thin enough across a cluster that grows much larger than expected; too many, and each partition carries fixed overhead (metadata, compaction bookkeeping, replication streams) that adds up needlessly at low node counts. This strategy is common in hash-partitioned systems like Cassandra and Riak, both descended from the Dynamo design. **The dynamic-partitioning approach** instead lets the number of partitions grow and shrink with the data itself: a partition exceeding a configured size or request-rate threshold gets split in two, each half potentially migrating to a different node to spread load, and partitions that shrink can be merged back together. This adapts naturally to unpredictable or highly skewed growth without needing to guess a partition count in advance, but requires an active coordinator continuously deciding when and how to split/merge, plus bookkeeping to track changing partition boundaries over time — meaningfully more operational and implementation complexity. This is the model used by range-partitioned, LSM-tree-based systems like HBase and Bigtable, where region splitting (HBase's term) is a well-known, automatic background process. | Strategy | Partition count | Assignment | Seen in | |---|---|---|---| | **fixed-partition-count** | decides at cluster creation time; partitions move as indivisible units | comparatively simple — a load-balancing/bin-packing problem | Cassandra and Riak | | **dynamic-partitioning** | grows and shrinks with the data itself | an active coordinator continuously deciding when and how to split/merge | HBase and Bigtable | ## Stream, then cut over The data-movement mechanism matters independently of which assignment strategy is used. The safe pattern most systems converge on is: 1. the destination node begins streaming a copy of the partition's data from the source (often a snapshot plus a stream of subsequent changes, to avoid losing writes that happen during the copy), 2. the source keeps serving reads and writes throughout, 3. and only once the destination has fully caught up does the system perform an **atomic cutover** — updating routing metadata so new requests go to the destination, while in-flight requests against the source complete. Getting this ordering wrong is the source of most rebalancing correctness bugs: cut over routing before the destination has all the data, and reads against the new owner return incomplete results; forget to keep capturing writes that land on the source during the copy, and those writes are silently lost from the destination's copy. ## Failure modes in production Several concrete failure modes show up in production. - **Thundering-herd rebalancing** happens when a node addition or removal triggers many partitions to move simultaneously, saturating network bandwidth and disk I/O across the cluster and causing latency spikes on totally unrelated partitions that happen to share the same physical links or disks — well-designed systems throttle the rate of concurrent partition moves for exactly this reason, trading a slower rebalance for a stable one. - **Stale-routing windows** happen when clients or coordinators haven't yet learned about a completed cutover and keep sending requests to the old owner, which must either reject, forward, or (worse, if not handled) silently serve now-outdated local data — this is why many systems keep the old owner able to forward or redirect for a grace period after cutover. - **Compounding movement plans.** A related risk during simultaneous multi-node changes is those plans conflicting with each other if the coordinator doesn't serialize or carefully merge them, which some systems handle by allowing only one rebalance operation in flight per partition at a time. ## Where it shows up in real systems A concrete real-world illustration: - **Cassandra's bootstrap process** streams token-range data to a newly joining node using its consistent-hashing-with-vnodes assignment, and operators are explicitly advised to add nodes one at a time and let streaming complete before adding the next, precisely to avoid the thundering-herd and conflicting-rebalance-plan failure modes described above. - **HBase's automatic region splitting** is the dynamic-partitioning counterpart, splitting a region once it exceeds a configured size and letting the master reassign the resulting two regions across region servers as a separate load-balancing step.

  • Why do most systems choose to stream data from source to destination and cut over atomically, rather than just copying the data in one shot and then switching?
    A single-shot copy would require freezing writes to the source for the duration of the copy, which for any nontrivial partition size means unacceptable write downtime. Streaming a snapshot plus ongoing changes lets the source keep serving writes throughout, and the atomic cutover only happens once the destination is fully caught up, minimizing the unavailability window to roughly zero.
  • What operational practice mitigates the thundering-herd problem when scaling a cluster by several nodes at once?
    Adding nodes one at a time (or throttling the number of concurrent partition moves) and letting each node's data streaming complete before starting the next, rather than triggering a full cluster-wide reassignment all at once. This bounds the concurrent network/disk load from rebalancing instead of saturating shared infrastructure across many simultaneous moves.
  • In a dynamic-partitioning system, what's the operational trade-off of splitting partitions too aggressively (a low size threshold)?
    You end up with a very large number of small partitions, each carrying fixed per-partition overhead (metadata, compaction/compression bookkeeping, replication tracking), which adds up to real resource waste and can slow operations that enumerate or coordinate across many partitions. It's the mirror image of the fixed-partition-count problem — too many small partitions instead of too few large ones.

Rebalancing is like a moving company redistributing furniture between apartments when a new roommate moves in. You can either predefine many small labeled boxes from day one and just hand whole boxes to whoever needs them (fixed partitions), or pack and repack boxes on the fly as any one room gets too full (dynamic splitting). Either way, the movers keep the old apartment livable while they carry boxes to the new one, and only update everyone's 'where do I find my stuff' notes once every box has actually arrived — and you don't want every mover in the building doing this at the same time or the elevator jams.

saying these in an interview costs you the question

  • Thinks rebalancing is instantaneous / doesn't involve any data streaming
  • Doesn't distinguish the assignment decision from the physical data-movement mechanism
  • Unaware that adding many nodes simultaneously can cause a thundering-herd of data movement
  • Assumes routing metadata updates propagate instantly to all clients with no staleness window
  • Can't name any concrete rebalancing strategy (fixed partition count or dynamic split/merge)

context