Sharding gets the headlines; rebalancing is where sharded systems actually hurt. A cluster that was balanced on launch day drifts: one tenant grows ten times faster than the rest, a new node joins with empty disks, a rack is drained for maintenance, and a celebrity key turns one shard into the hottest machine in the fleet. Rebalancing is the process that moves data and load back toward an acceptable distribution while the system keeps serving traffic.
This article treats rebalancing as a control loop with a cost: the strategy families, when to act, a move planner, a worked node addition, one safe move, and the failure modes of the loop itself. Routing and key hashing are covered in Database Sharding Architecture and Hash Sharding; here we assume a router exists and ask how the map behind it should change.
The rebalancing loop at a glance
What balanced means, and what it costs
Start by being precise about what balanced means, because every later decision depends on it. There are at least three different quantities you might balance, and they disagree.
Storage. Bytes per node, measured against disk capacity. This is what fills disks and causes outages, and it changes slowly. Load. Requests per second, CPU, or IO per shard, which changes in minutes and is spiky. Count. Number of shards or replicas per node, which is cheap to compute and is what many systems balance by default because it is stable, even though it is only a proxy for the other two.
A rebalancer must also respect constraints unrelated to balance: replicas of one shard never share a node, and usually not a rack or zone; a node above a disk threshold receives nothing; some shards are pinned to a region. A plan that violates any of these is wrong, not merely suboptimal.
Every move costs bandwidth, disk IO on both ends, cache warm-up, and a window of reduced resilience. So the objective is an acceptable spread reached with as little movement as possible. Perfect balance is never the goal; the last few percent cost as much movement as the first fifty.
The three strategy families
Real systems fall into three families. The family you pick decides what a rebalance even is.
| Family | Unit that moves | Examples | Strength | Weakness |
|---|---|---|---|---|
| Fixed number of partitions | Whole partition | Redis Cluster (16,384 hash slots), Couchbase (1,024 vBuckets), Elasticsearch primary shards, Kafka partitions | Moves are discrete and simple; the key-to-partition map never changes | Partition count is fixed early; a partition cannot be smaller than its hottest key |
| Dynamic split and merge | Key range | HBase regions, Bigtable tablets, CockroachDB ranges (split at a configured maximum size, 512 MiB by default in recent versions) | Adapts to growth automatically; hot ranges can be split | Splits need good split points; many small ranges cost metadata |
| Consistent hashing with virtual nodes | Token ranges | Cassandra, Dynamo-style stores | Adding a node takes a small, spread-out slice from everyone | Balance is statistical; few vnodes give lumpy load, many make streaming heavier |
The fixed-partition family is the one most teams build, because it decouples two problems. The key-to-partition mapping is a pure function that never changes, and the partition-to-node mapping is a small table the rebalancer edits. Choose many more partitions than you will ever have nodes, commonly ten to a hundred per node at the largest planned size, so averages work but each partition still moves in minutes.
Range-based systems rebalance differently. They first split a range that is too big or too hot, which creates new units, and then move some of those units. Range Sharding covers choosing split points; the planner below works for both families once splitting has produced movable units.
When to rebalance: triggers and hysteresis
A rebalancer that runs whenever the spread is non-zero will never stop moving data. You need explicit triggers and a dead band between them.
- Membership change. A node is added, decommissioned, or failed for longer than a grace period. The most common and most predictable trigger.
- Capacity threshold. A node crosses a disk watermark. Elasticsearch uses low, high and flood-stage watermarks: above low it stops allocating shards to the node, above high it moves shards away. Copy that two-level design.
- Imbalance threshold. Most-loaded node over the mean exceeds a bound, measured on a metric smoothed over tens of minutes.
- Operator request. Draining a node or rack, using the same planner with the drained node at zero capacity.
Hysteresis is the key idea: start when imbalance exceeds an upper bound such as 1.25 and stop below a lower one such as 1.10, or a cluster near the threshold flips on noise. Node failure needs the same care, because most disappearances are restarts; Elasticsearch exposes a delayed-allocation timeout for this. Pick a delay longer than a typical restart and shorter than you will tolerate reduced redundancy.
A move planner with a budget
The planner takes the current assignment, per-partition weights and the constraints, and returns a short list of moves. Optimal assignment is a bin-packing problem and NP-hard, but a greedy planner with a movement budget is the common practical choice, and it works well because the starting point is usually close to balanced. The loop is: find the most loaded node, find the least loaded eligible node to receive from it, pick the partition whose move most reduces the squared distance of both nodes from their targets, apply it in memory, repeat until balanced or out of budget.
from dataclasses import dataclass
@dataclass(frozen=True)
class Move:
partition: str
src: str
dst: str
def plan_moves(assign, weight, nodes, capacity, rack, replicas_of,
max_moves=20, target_ratio=1.10):
"""assign: partition -> node; weight: partition -> load units;
capacity: node -> relative capacity (0 means draining);
replicas_of: partition -> set of nodes holding OTHER replicas."""
load = {n: 0.0 for n in nodes}
for part, node in assign.items():
load[node] += weight[part]
total_cap = sum(capacity[n] for n in nodes)
total = sum(load.values())
def util(n):
return load[n] / capacity[n] if capacity[n] else float("inf") if load[n] else 0.0
def target(n):
return total * capacity[n] / total_cap
live = [n for n in nodes if capacity[n] > 0]
mean = total / total_cap
moves = []
while len(moves) < max_moves:
src, low = max(nodes, key=util), min(live, key=util)
if util(src) <= target_ratio * mean and util(low) >= mean / target_ratio:
break # inside the dead band
best = None
for part in (q for q, n in assign.items() if n == src):
w = weight[part]
for dst in sorted(live, key=util):
if dst == src or dst in replicas_of[part]:
continue # never co-locate replicas
if rack[dst] in {rack[r] for r in replicas_of[part]}:
continue # rack anti-affinity
before = (load[src] - target(src)) ** 2 + (load[dst] - target(dst)) ** 2
after = (load[src] - w - target(src)) ** 2 + (load[dst] + w - target(dst)) ** 2
gain = before - after # > 0 only if the move helps
if best is None or gain > best[0]:
best = (gain, part, dst)
break # least-loaded valid dst only
if best is None or best[0] <= 0:
break # nothing helps; stop, do not thrash
_, part, dst = best
assign[part] = dst
load[src] -= weight[part]
load[dst] += weight[part]
moves.append(Move(part, src, dst))
return movesThree details matter more than the search. The planner stops when no move has positive gain, which prevents ping-ponging. The movement budget caps disturbance per round; the next round replans from fresh metrics. And the weight is a choice: bytes, a smoothed request rate, or a weighted sum, with both components logged so operators can see why a move was chosen.
Worked example: adding a fifth node
Take a cluster of four nodes, each with capacity 1, holding 64 partitions of roughly equal size, 16 per node, about 400 GB per node. You add a fifth node. The target is now 64 / 5 = 12.8 partitions per node.
Hash-mod-N would remap about four fifths of all keys. With fixed partitions the planner instead moves whole partitions to the new node: three from each old node, 12 moves in total, ending at 13, 13, 13, 13 and 12, which is inside a 10 percent dead band. That is 12 of 64 partitions, about 19 percent of the data, close to the floor: the new node must end up with roughly a fifth of everything.
Now the operational arithmetic. Each partition is 25 GB. If you throttle moves to 100 MB/s per move and allow two concurrent moves, which is Elasticsearch's default for cluster-wide concurrent rebalances, you move 200 MB/s, and 300 GB takes 25 minutes plus catch-up. Raise concurrency to eight and it takes about six minutes, but the new node ingests 800 MB/s while starting to serve reads. Set the number from foreground latency measured during a test rebalance, not from impatience.
Notice what the planner did not do: shuffle partitions among the original nodes. Every move targeted the new node, the only destination with positive gain, and that property is what makes rebalances predictable.
Executing one move safely
Planning is cheap; executing a move safely is where the engineering goes. A partition move that keeps serving writes follows four phases, shown in the diagram above.
- Copy. The destination streams a consistent snapshot of the partition from the source, or from any replica, at a throttled rate. Reads and writes continue on the source.
- Catch up. The destination replays the log of writes that arrived since the snapshot. Repeat until the lag is small and stable.
- Fence and cut over. Briefly block writes to the partition, or rely on a replication protocol that already orders them, apply the last entries, then update the partition map with a new epoch. Routers and clients that still hold the old epoch get a stale-epoch error and refresh. The source must reject writes stamped with an old epoch, otherwise two owners accept writes and you lose data.
- Clean up. Delete the source copy only after a grace period, so you can roll back by flipping the map back.
Replicated systems often do this more cheaply by adding a replica on the new node, letting it catch up as a follower, transferring leadership, and then removing the replica on the old node. Kafka's partition reassignment works this way: the new replicas join the replica set, catch up, and the old ones are dropped. Its reassignment tool accepts a --throttle in bytes per second, and running it with --verify after the reassignment completes also removes the throttle, which otherwise keeps limiting replication traffic. Kafka Partition Architecture walks through that path in detail.
Failure modes of the control loop
Most rebalancing incidents are caused by the control loop, not by a bug in copying bytes.
- Oscillation. The planner moves a hot partition from A to B, the hotness moves with it, B becomes the hottest, and the next round moves it back. Fix with hysteresis, a minimum residence time per partition (for example no second move within an hour), and gain that must exceed a margin.
- Rebalancing into an overload. The cluster is slow because it is overloaded, the rebalancer sees imbalance, and starts moving data, which adds IO and makes everything slower. Rebalancing should back off when foreground latency or error rate is above its objective, unless a disk is about to fill.
- Flapping nodes. A node restarts repeatedly; each disappearance triggers re-replication and each return triggers a rebalance back. Delay before reacting, and quarantine a node that has failed several times in an hour.
- Hot keys that cannot be moved away. If one key carries 30 percent of traffic, moving its partition just moves the fire. Rebalancing cannot fix this; you need key splitting, caching or request coalescing.
- Stale metrics. The planner uses load data from before the last round's moves finished, so it double-corrects. Wait for moves to complete and metrics to settle before replanning.
- Forgotten throttles and pauses. A throttle or a disabled-rebalance setting left over from an incident silently prevents later recovery. Make these settings visible on the dashboard with their age.
Operating the rebalancer
Operate the rebalancer like any other production service, with its own objectives and dashboards.
- Per-node utilisation for each balanced quantity, plus the max-to-mean ratio over time.
- Moves in flight and queued, bytes per second moved, and the age of the oldest in-flight move; one far past its expected duration is stuck.
- Foreground latency percentiles annotated with move start and end times, for tuning concurrency.
- A planner dry-run that lists the moves it would make and why, for review before large changes.
- A single kill switch that stops new moves, lets in-flight moves finish, and is reported loudly while set.
Trade-offs: balancing by count is stable but blind to skew; by load, responsive but noisy. A robust default is fixed or size-split partitions, a planner balancing bytes with a load term, conservative concurrency, and a manual emergency path. For hash-based placement with bounded loads, see Consistent hashing architecture.
What to do next
- Write down which quantity you balance (bytes, load, count) and the upper and lower thresholds of your dead band.
- Check your partition count: aim for at least ten partitions per node at your largest planned cluster size, each small enough to move in minutes.
- Implement or configure a delay before reacting to node loss that is longer than a normal restart.
- Add a movement budget per planning round and a minimum residence time per partition.
- Run a test rebalance in staging under production-like load, record foreground p99 at several concurrency levels, and set the default from that data.
- Make the epoch check on the write path mandatory, and test a cutover with a client that holds a stale map.
- Build the dashboard and the kill switch before the next node addition, not during it.