Rebalancing Partitions Under Live Traffic
LESSON
Rebalancing Partitions Under Live Traffic
By the end of this lesson, you will be able to...
trace a live partition move from baseline copy through catch-up and ownership cutover;
explain why copying data does not by itself transfer write authority;
choose the evidence and failure signals that make a migration safe to complete or stop.
Idea in one sentence: A rebalance is safe when data catch-up and authority change meet at one explicit boundary, so only one owner can accept the next decisive write.
Core Insight
Harbor Point routes confirmations for allocation bucket B-173 to shard group B. During a market burst, this bucket drives much more traffic than neighboring partitions. Group F has capacity, so the team wants to move B-173 while traders continue to extend holds and confirm reservations.
The tempting plan is short: copy the bucket to F, point routers at F, then delete the copy on B.
It works only if no write arrives during the copy. In the real case, a trader extends hold H-8821 after the snapshot begins. If that extension reaches only B, F is stale. If both groups accept confirmation “for safety,” two authorities can decide the same scarce bucket.
The stronger model is an ownership handoff. The source remains the one writer while the target catches up. The handoff has a named boundary. Before that boundary, B can decide. After it, F can decide. A router or a stale request must not invent a third answer.
The trade-off is explicit: a short, scoped write barrier can delay one partition's commands, but it avoids the much larger correctness cost of two live owners for the same scarce allocation.
The Moving Parts
Use four stable names for the rest of the lesson:
| Part | Job during the move |
|---|---|
Source group B |
Serves the current owner’s reads and writes until cutover |
Target group F |
Receives a baseline and later changes, but is not yet an authority |
| Partition map | States which group owns partition 173 in each map version |
| Router and shard checks | Carry and enforce the version that selected the owner |
At the start, Harbor Point has this map:
map v104
partition 173 -> B
partition 174 -> F
The values are illustrative. A production system might use range metadata, a directory, or a control-plane API rather than this exact map. The invariant is the same: the data plane needs an unambiguous answer to “who may accept a write for partition 173 now?”
The first model to correct is that the target becomes useful as soon as it has a snapshot. A snapshot is a picture of one point in time. It is not a promise that the target has every command the source accepted after that point.
Build the Handoff One State at a Time
Harbor Point moves partition 173 with one source writer. This is a teaching trace, not a universal vendor procedure.
Starting state
map v104: 173 -> B
B is authoritative
F has no copy
1. Establish a copy point
B records a consistent snapshot point: sequence 918443.
B stays authoritative for new commands.
2. Copy the baseline
F receives the state of partition 173 as of sequence 918443.
F may verify and store it, but cannot accept confirmation writes.
3. Catch up
B continues to accept commands and streams later records to F.
F applies 918444, 918445, ... in order.
4. Create a cutover barrier
B briefly stops or drains new writes for partition 173.
The last source record is 918499.
F proves it has applied through 918499.
5. Switch authority
Publish map v105: 173 -> F.
B rejects old-owner writes; F accepts writes for v105.
6. Verify and clean up later
Observe traffic, compare state, retain recovery evidence,
then remove B's old copy under a defined cleanup rule.
The intermediate state in step 3 is the reason the snapshot alone is insufficient. At 10:03:07, after the snapshot, a trader extends H-8821. Group B assigns the extension sequence 918487. The target cannot safely own the bucket until it applies that record and every later record through the barrier.
snapshot point on B: 918443
extension of H-8821 on B: 918487
barrier chosen for cutover: 918499
target F has applied: 918487 -> not ready
target F has applied: 918499 -> eligible for the switch
This is the evidence that earns the correction: F receiving a complete bulk copy says nothing about 918487. A target that has replayed through the barrier has the source history needed for the next owner state.
Some systems make parts of this sequence internal. They may use consensus membership changes, replicated logs, or a migration controller rather than an application-visible sequence number. The teaching model still applies: establish a source state, catch up later changes, change authority once, and verify the result.
Fence the Old Owner
Copy and catch-up protect the target from stale data. They do not stop stale routers from sending new writes to the source after the switch.
Harbor Point includes the partition-map version in a routed command:
def apply_reservation_command(partition, map_version, command):
if map_version != partition_map.current_version(partition):
raise RetryWithFreshMap(partition)
if not partition_map.is_owner(partition, this_shard_group):
raise WrongOwner(partition)
record_and_apply(command)
Before cutover, a command for partition 173 carries v104 and group B accepts it. After cutover, a router that still sends v104 to B receives a predictable rejection. It refreshes its map and retries through F with v105.
This check is called fencing in this lesson: it prevents a former owner from acting on an old authority claim. The exact protocol may use an epoch, lease, generation number, or configuration version. What matters is that a newer authority claim makes the old one unusable for decisive writes.
Why not quietly forward every stale request from B to F? Forwarding can be useful in a particular product, but it must preserve the command identity, ordering and authorization context. A clear wrong-owner result is often easier to test. Harbor Point chooses it for critical confirmations because the retry is narrower and more observable than hidden forwarding behavior.
A Worked Command Through the Cutover
Follow one idempotent confirmation command, K-44, through the difficult boundary.
10:03:00.000 map v104 says partition 173 -> B.
10:03:00.010 B begins snapshot at sequence 918443.
10:03:02.000 B accepts extend H-8821 at sequence 918487.
10:03:04.000 F replays through sequence 918498.
10:03:05.000 migration starts the write barrier for partition 173.
10:03:05.004 B drains accepted commands and records barrier 918499.
10:03:05.010 F reports applied_through=918499.
10:03:05.012 control plane publishes map v105: 173 -> F.
10:03:05.015 stale router sends confirm K-44 to B with v104.
10:03:05.015 B rejects it: RetryWithFreshMap.
10:03:05.020 router refreshes and sends K-44 to F with v105.
10:03:05.022 F records and applies K-44 once.
The command has an idempotency key because retries are normal at this boundary. If the router loses the response after F commits K-44, the next retry must recover the recorded result, not create a second confirmation.
The trace makes two conditions visible:
Fdid not acceptK-44until it had replayed through the barrier and ownedv105.Bdid not acceptK-44afterv105, even though it still held an old copy of the data.
So far, the migration is not “safe because the copy finished.” It is safe for this command because catch-up, fencing, and idempotency agree on one owner at the same moment.
Read Paths, Verification, and Cleanup
Writes are the sharpest ownership test. Reads still need a contract.
For POST /reservations/confirm, Harbor Point routes to the authoritative owner before and after the move. For an advisory search page, the product may allow a derived index or a replica to lag, but it should state the allowed freshness. A healthy target process is not enough evidence that it can serve a particular guarantee.
After the switch, retain the source copy long enough to verify the move and satisfy recovery rules. Useful checks include:
| Check or signal | What it establishes |
|---|---|
| Target applied position versus barrier | Did the target reach the handoff history? |
| Source and target counts or checksums | Does copied state match at the chosen scope? |
| Stale-route rejection rate | Are routers still using the old map? |
| Idempotency-result mismatch count | Did retries cross the handoff incorrectly? |
| Partition error rate and latency before/after | Did the move solve pressure without harming the endpoint? |
| Cleanup age and rollback evidence | Is the old copy retained long enough, but not forever? |
Counts and checksums are evidence, not a replacement for an ownership protocol. Two copies can have matching rows while a stale source still accepts an unsafe write. Conversely, an old source copy can differ after cutover without being a problem if it is explicitly non-authoritative and retained only for recovery.
The cleanup policy is a situated choice. Fast deletion saves capacity. Delayed deletion preserves an audit and rollback path. The team should define which migration evidence, drain period, and recovery point are required before removing the old copy instead of treating cleanup as an automatic last command.
Costs and Boundaries
Live rebalancing improves capacity management, hotspot relief, hardware evacuation, and planned maintenance. It costs control-plane state, data transfer, temporary catch-up work, retry handling, tests, and migration-specific observability.
A brief write barrier is a preference for an invariant-sensitive partition. It creates a small, scoped availability dip. In return, the handoff has one simple ordering point. A system with a different product contract may use a more elaborate online protocol, but it still needs to explain how it prevents two owners from accepting conflicting commands.
Rebalancing does not repair a bad shard key. If the same logical partition becomes hot again immediately after every move, the problem may be the workload boundary, a single hot tenant, or an access pattern that needs a different design. The signal is repeated migration pressure on the same key range, not merely high cluster-average utilization.
Check Your Understanding
Check: Target F has finished its baseline copy, but it has replayed only through 918487 while source B recorded the cutover barrier at 918499. Can F become the write owner?
Think first, then reveal.
Answer: No. F is missing records after 918487. The control plane must wait until F has applied through the chosen barrier, then change authority. Copy completion is not catch-up evidence.
Check: After map v105 names F, a router sends a v104 confirmation to B. Why is rejection better than accepting it and hoping the router refreshes soon?
Think first, then reveal.
Answer: Accepting it gives the former owner a live authority path after the handoff. Rejecting it fences the old owner and forces one retriable route through the map that now names F.
Practice: Design a Safe Hot-Partition Move
A delivery service moves package partition P-8 from group X to group Y during a busy hour. A package-acceptance command has idempotency key A-91. Describe the smallest safe migration plan.
Include:
- a snapshot point and later catch-up marker;
- the moment at which
Ymay first accept the command; - how
Xhandles an old routing version after cutover; - how a retry of
A-91avoids duplicate acceptance; - one signal that would halt cleanup of
X's old copy.
A strong answer keeps X authoritative during baseline copy and catch-up. It pauses or otherwise fences the handoff, proves that Y reached the marker, then publishes a new routing version. X rejects stale-version writes. Y records A-91 with the command result, so a retry returns the existing outcome. A high stale-route rejection rate, a target checksum mismatch, or an idempotency-result mismatch should stop cleanup and trigger investigation.
Connections
- Partitioning Strategies and Shard Keys makes a move possible by separating a stable authority key from changing placement.
- Secondary Indexes Across Shards applies the same ownership question to derived lookup entries when a base partition moves.
- Membership Changes and Replica Set Evolution examines a related but distinct change: who belongs inside the replica group that serves one authority domain.
Resources
- [DOC] Vitess: Resharding — Focus: Follow validation, staged read switching, write switching, and cleanup in an online reshard workflow.
- [DOC] CockroachDB: Replication Layer — Focus: Inspect ordered replication logs, catch-up, membership change, and load-driven rebalancing as built-in mechanisms.
- [DOC] MongoDB: Sharded Cluster Balancer — Focus: Compare automatic placement balancing with application-level ownership claims and operational throttling.
- [BOOK] Designing Data-Intensive Applications — Focus: Connect partition migration, routing, replication, idempotency, and recovery to an explicit client contract.
Key Takeaways
- A baseline copy creates a candidate target; only catch-up through a barrier makes that target eligible to own new writes.
- A map version, epoch, or equivalent fence stops former owners from accepting decisive commands after cutover.
- Idempotency keys turn cutover retries into recovery of one recorded decision rather than a second decision.
- Verification must test both state parity and the authority boundary; matching copies alone do not prove a safe handoff.
- Repeated moves of the same hot partition are evidence to revisit the workload boundary, not just to move bytes faster.