TECH
Rebalancing a Cluster While the Write Path Stays Open
Rebalancing a live database cluster is a scheduling and throttling problem because the write path stays open throughout. Every node you move risks a stalled quorum, and the migration job must coexist with client traffic rather than run during a quiet hour that never arrives. This piece covers why naive rebalancing stalls, which techniques keep writes flowing, how to choose a throttle rate, what to instrument, and the mistakes that turn a routine migration into an outage.
The write path stays open
Maintenance windows are mostly fiction for systems that take writes around the clock. A payments ledger, a session store, a telemetry pipeline: each one has clients that will retry, buffer, or fail loudly the moment a quorum is unavailable. The rebalancing job has to run inside that reality, not around it. Write availability is the hard constraint, and everything else in the plan is negotiable.
The scheduling question is where to spend the cluster's spare capacity. Moving a shard consumes network bandwidth, disk IOPS, and CPU on both the source and the destination. If that consumption overlaps with a peak write period, latency climbs for unrelated tenants sharing the same nodes. A move that takes twenty minutes at 3 a.m. might take two hours at noon and still degrade the p99 for everyone else.
Treat rebalancing as a background workload with a budget. It gets a fixed slice of bandwidth and disk throughput, and it yields when the foreground write path needs the resources back. That framing changes the design: you stop looking for a quiet hour and start building a controller that can run for days without anyone watching it.
Why naive rebalancing stalls
The simplest approach is a full-stop migration: quiesce writes, copy data, reassign ownership, resume. It works for small datasets and short outages, and it fails badly for anything with a strict availability target. During the freeze, every client retry piles up. When writes resume, the burst can overwhelm the newly assigned nodes before replication catches up.
Bulk copying has its own failure mode. Saturating a network link to move a large shard starves the replication traffic that keeps replicas current. Replication lag grows, and if a node fails during the move, the surviving replicas may not have the latest writes. The migration itself creates the fragility that makes failover dangerous.
A shard handling a disproportionate share of writes will show elevated latency the moment its data starts moving, because the migration competes for the same disk and CPU. If the hot shard is also the one being split, the split point matters: a bad key range leaves one child still hot and the other nearly idle. The next rebalance then repeats the exercise on the same keys.
Failover during migration compounds risk. A node that goes down mid-move forces the controller to choose between aborting the migration and continuing with reduced redundancy. Both choices have costs, and the decision usually has to be made in seconds. Teams that have not rehearsed it tend to pick wrong.
Techniques that keep writes flowing
Consistent hashing limits how many keys get remapped when the ring changes. Adding or removing a node moves only the keys owned by the affected arc, not the whole keyspace. That property is why it shows up in so many distributed stores, and it is the first thing to check when evaluating a rebalancing strategy.
Rendezvous hashing takes a different route. Clients compute the same assignment independently from a shared node list, so no central coordinator has to broadcast a new mapping. The trade-off is that every client needs the current list. When a node is removed, a client still holding the previous generation of the list will keep sending requests to an address that no longer owns those keys. The receiving node can respond with a redirect or a 'wrong node' error, and the client uses that signal to fetch the updated list from the coordination service before retrying. Until the refresh completes, that client's requests either fail or land on a node that forwards them, adding a hop and some latency. Detection usually relies on a version number or generation tag embedded in the list, so a client can tell at a glance whether its cache is stale.
Dual-write windows with read fallback let you move data without a hard cutover. The source keeps serving reads while the destination receives writes, and the controller compares results until they agree. Only then does the destination take over reads. The window costs extra write amplification, so it should be bounded and monitored.
Throttled chunk migration at bounded rates is the workhorse. Move one chunk at a time, cap the rate, and pause when replication lag crosses a threshold. The migration takes longer, but the foreground write path barely notices. This is the technique most production systems converge on, and it is the one that benefits most from good instrumentation.
Choosing the throttle rate
Throttling trades total migration time for write-path stability. A faster migration finishes sooner but competes harder with foreground traffic; a slower one stretches the window during which the cluster is in a mixed state, with some shards moved and others not. The right rate is the fastest one that keeps write latency within its normal band.
Start by measuring spare IOPS and bandwidth on both the source and destination nodes during a representative period. Cap migration throughput at a fraction of that spare capacity, then watch per-shard write latency and replication lag. If either crosses a threshold you set in advance, the controller should pause the current chunk and resume when the metric recovers. A simple proportional controller that halves the rate on a lag spike and creeps back up works well enough for most clusters.
The trade-off is not free. Longer migrations mean more time spent with dual-write windows open, more chances for a node to fail mid-move, and more operational attention. Teams that ignore this end up either migrating too fast and degrading the write path, or too slow and leaving the cluster in a half-moved state for weeks.
Instrumentation before and during moves
Baseline normal write patterns before touching anything. Per-shard write latency percentiles, replication lag, queue depth, and disk utilization all need a known-good reference. Without a baseline, you cannot tell whether a latency spike is caused by the migration or by a client that started sending larger payloads that morning.
Watch replication lag as a leading indicator. It rises before quorum loss and before client-visible errors, which gives the controller time to pause the migration. A related piece on this site about retry logic in service clients makes the same point from the client side: the retry policy determines how much a brief lag spike actually costs.
Alert on quorum loss, not just node down. A node can be reachable and still fail to participate in a quorum because of disk stalls or network partitions. The alert that matters is the one that fires when the cluster can no longer commit writes, and it should page someone who can pause the migration.
Per-shard write latency percentiles catch the hot-shard problem early. If one shard's p99 doubles while the others stay flat, the migration is probably competing with a tenant that happens to live on that shard. Pausing that shard's move and finishing the others first is usually the right call.
What teams actually get wrong
Assuming replication factor equals safety is the most common error. Three replicas protect against a single node failure, not against a migration that leaves two of them behind on the same shard. The redundancy is nominal if the replicas are not current, and the migration is exactly the workload that makes them stale.
Reassigning ownership before the destination has caught up is another frequent mistake. The controller marks the new node as the owner and stops routing reads from the old one, but the destination's replica set may still be missing the most recent writes. Reads then return stale data, and if the old node is decommissioned before the gap closes, those writes are gone. The fix is to gate the ownership change on a verified catch-up: compare a checksum or a write sequence number between source and destination, and only flip the routing table when they match.
Ignoring client-side retry storms turns a brief quorum gap into a sustained overload. Clients with aggressive retry policies and no jitter will synchronize their attempts, and the resulting traffic can keep a recovering node from catching up. This site has argued in a piece on migrating off Lambda that the cost of a transition often lands in places the plan did not budget for; rebalancing has the same shape.
Skipping dry runs on production-like load hides the failure modes that only appear at scale. A migration that works on a test cluster with a tenth of the write volume may stall on the real one because the throttle settings are wrong or the disk IOPS budget is too generous. The dry run should use the same client mix and the same payload sizes.
A practical rebalancing checklist
Measure write throughput and latency baselines first, per shard, over at least a full business cycle. Without those numbers, the throttle settings are guesses and the alerts have no threshold that means anything.
Choose consistent or rendezvous hashing deliberately, based on whether you can broadcast a node list to clients or need a coordinator to own the mapping. Both limit remapping; the operational difference is who has to know the current state.
Throttle migration to a fraction of spare capacity, and pause automatically when replication lag exceeds a threshold you set in advance. A controller that can only be stopped by a human is a controller that will run too long during an incident.
Keep dual-write windows until reads verify consistency, then close them on a schedule rather than on a hunch. The window is the safety net; removing it early is how a successful migration becomes a data-loss incident.
Schedule moves during known low-write periods, not just low-traffic periods.