Resharding & Rebalancing — Dual-Write, Online Moves & Consistent Hash Rings
If you cannot reshard, you cannot shard safely. Online moves are dual-write plus backfill plus cutover, engine range splits, or a consistent-hash vnode remap. Modulo-N moves most keys when N changes. A ring moves a slice. This is a data-plane move, not a schema migration.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
The shard count has to change
Prefer
An online move
Dual-write and backfill, a range split or merge, or a vnode remap with a copy. Reads flip on a flag. The old copy stays until lag is zero and shadow checks are quiet.
- You can abort by sending reads back.
- Only the keys whose owner changed need to move on a ring.
- Whale isolation can be scheduled before the shard melts.
Alternative
Big-bang downtime
Stop writes, copy, restart. Simple, and it trains the team never to reshard because the copy window only grows.
- The SLA is the copy time.
- Imbalance gets worse while you wait for a window.
- There is no shadow read. You discover divergence after the cut.
Overview
If you cannot reshard, you cannot shard safely. Shard count changes. A whale needs isolation. A range goes hot. The interview signal is whether you can move live traffic without a maintenance window.
Three families:
| Family | Mechanism | Pros | Cons |
|---|---|---|---|
| Dual-write plus backfill | App writes old and new, copies history, switches reads | Works with ordinary storage | Divergence risk; a long overlap |
| Range split or merge | Engine moves key ranges (Vitess, TiDB, CockroachDB) | Less application logic | Needs a substrate that can move ranges |
| Consistent-hash vnode remap | Remap virtual nodes onto physical databases | Fewer keys move than modulo-N | You still copy the remapped vnodes |
This is data-plane reshard dual-write. It is not schema expand/contract. Column adds, online DDL, and migration dual-write are Zero-downtime database migrations. Clockwise ring walks, replica selection, and vnode theory are Consistent hashing. Here, each vnode is a set of keys that currently lives on a physical Postgres or MySQL shard.
Flow
- 1
Phase 1: dual-write old and new
- nextBackfill: read old, upsert new
- 2
Backfill: read old, upsert new
- nextPhase 2: shadow-read and compare
- 3
Phase 2: shadow-read and compare
- nextPhase 3: reads and writes go to new
- 4
Phase 3: reads and writes go to new
- nextPhase 4: drop old after lag is 0
- 5
Phase 4: drop old after lag is 0
Lesson map
Resharding & Rebalancing — Dual-Write, Online Moves & Consistent Hash Rings
If you cannot reshard, you cannot shard safely. Online moves are dual-write plus backfill plus cutover, engine range splits, or a consistent-hash vnode remap. Modulo-N moves most keys when N changes. A ring moves a slice. This is a data-plane move, not a schema migration.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB p1["Phase 1: dual-write old and new"] p2["Backfill: read old, upsert new"] p3["Phase 2: shadow-read and compare"] p4["Phase 3: reads and writes go to new"] p1 -->|Phase 1: dual-write old and new| p2 p2 -->|Backfill: read old, upsert new| p3 p3 -->|Phase 2: shadow-read and compare| p4
Dual-write discipline for a reshard
Same primary key on both sides. Old write, then new write, until the flag says otherwise.
- 1
Idempotent keys
Upsert on the same primary key. A retry must not double-apply. - 2
Ordering
Prefer write old, then new. Decide in advance what happens when the new write fails: fail the request, or queue a retry. - 3
Backfill
Bounded workers, checkpointed, safe to replay. - 4
Shadow reads
Compare old and new. Alert on the mismatch rate. Do not cut over on a row count that happens to match. - 5
Cutover
Flip read routing. Keep dual-write until you trust the new side. Then single-write the new shard. - 6
Abort
A feature flag sends reads back to the old shard. Do not drop the old copy in the same step as the flip.
Rings, from the database side
Rebalance means: a vnode's key set moves to another physical shard, you copy the data, you flip routing, you drop the source copy after the lag is zero. Versus modulo-N: changing N moves most keys. A ring with virtual nodes moves about one Nth, and vnodes smooth that further. The primer on Toptal is enough public background if you have not read the hashing study. Do not rebuild the ring math in the reshard runbook.
Vitess can reshard without the application dual-writing every statement. VReplication streams rows while routing rules change. You still need a careful cutover. TiDB's placement driver splits and moves regions. CockroachDB rebalances ranges in the distribution layer. AWS DMS is a homogeneous copy tool when you are moving engines, not a substitute for a shard-routing flag.
Toy cutover
The Python router writes old then new while dual-write is on. After the historical copy, shadow reads can compare. Flipping read_new changes the primary read. Turning dual-write off after that flip single-writes the new store. If you clear dual-write but forget to retarget writes, the new shard freezes and the old shard moves on. That bug is the next page's divergence case. This sketch retargets on purpose.
The TypeScript function lists keys whose vnode owner changed. Only those keys need a copy. Keys that stay on the same physical shard are not part of the move.
Press Run. Snippets must be self-contained — no network, files, or native modules.
Press Run. Snippets must be self-contained — no network, files, or native modules.
"a" has length 1, so its vnode moves from s0 to s2. The other three sample keys stay. That is the whole point of a remap versus hash % N.
Interview Q&A
Why dual-write before cutover?
Answer
The new shard stays warm and you can verify it under live traffic. A backfill alone races with writes that landed after the copy started.
How do you detect divergence?
Answer
Shadow reads, checksum jobs, row counts per key range, and Merkle trees. A matching row count is not a matching payload.
What does consistent hashing buy a database?
Answer
Fewer keys move when you add or remove a physical shard. The ring and vnode placement are the hashing study. This page only needs the consequence: copy the vnodes whose owner changed, then flip routing.
What if the write to the new shard fails?
Answer
Pick consciously. Queue and retry with an idempotency key, or fail the user request. Do not succeed on the old shard, drop the error, and cut over anyway.
How is a schema migration different from a reshard?
Answer
A schema migration expands and contracts columns and code. A reshard moves rows between nodes. The overlap looks similar and the invariants differ. Do not reuse a column-migration runbook as a shard move.
Can Vitess reshard without application dual-write?
Answer
Yes. VReplication streams rows while routing rules update. The cutover is still careful. You did not escape verification, lag, or an abort plan.
When do you isolate a whale?
Answer
When one key dominates a shared shard, schedule a dedicated target early. Online move of that key range beats waiting for a full rehash.
What is the abort plan?
Answer
A flag that sends reads back to the old shard, while the old shard is still being written or is still caught up. Decommission only after lag is zero and you have decided not to roll back.
Pitfalls
- Turning off dual-write while writes still target only the old shard.
- Dropping the old copy in the same change as the read flip.
- Using modulo-N and discovering the reshard moves almost every row.
- Treating a DMS task as the routing layer. The app still has to know which shard is primary.
- Copying vnode-placement math into this runbook instead of linking the hashing study.
- Applying expand/contract column steps and calling that a reshard.
N is changing from 4 to 8 under modulo. A single tenant is 40 percent of one shard. You run Vitess. Say which move you use for each, what you verify before reads flip, and what the abort flag does.
Go deeper
Vitess resharding and VReplication are the MySQL online move. TiDB scheduling is region split and placement. CockroachDB's distribution layer is range rebalancing. The Toptal consistent-hashing primer is public background for why a ring moves fewer keys than modulo. AWS DMS introduces homogeneous copy when the engine itself is what you are moving.
Next: Failure modes.