skip to content

A hash-sharded cluster running 8 shards needs to grow to 16 because storage per shard is filling up. What makes resharding — moving data to the new shard count — operationally risky, and how do teams typically do it with minimal downtime?

level: seniorimportance: should knowfreq 55%

answer

  1. naive mod-N remaps almost everything on resize
  2. dual-write, backfill, verify, cutover
  3. split a shard vs. rehash the whole cluster
  4. split-brain and read-your-write risk during cutover
  5. Vitess VReplication / Mongo chunk balancer as real tooling

basics

~20 s

Resharding means moving live data between servers while the system keeps running, so the main risks are reads getting stale or wrong-shard answers mid-move and overloading production during the copy. Teams usually write to both old and new locations, copy the history over in the background, verify it matches, then switch reads over.

solid answer

~60 s

A naive modulo scheme (shard = hash(key) % N) is the classic trap: changing N from 8 to 16 remaps almost every key to a different shard, since key % 8 and key % 16 rarely agree, forcing a near-total data migration instead of an incremental one. Production resharding typically follows a dual-write/backfill/cutover pattern: the new topology starts receiving writes for the affected key ranges alongside the old one (via dual writes or a change-data-capture stream), a background job backfills historical data into the new shards, reads gradually cut over per key range once the target shard is verified consistent with the source, and the old location's data is cleaned up afterward. Splitting an existing shard into two (keeping most keys in place, moving only the split-off range) is far less disruptive than a full N-to-2N rehash and is the preferred approach when the sharding scheme supports it. Risks include split-brain if both locations accept writes without proper ordering, prolonged double-write overhead, read-your-write inconsistency during the cutover window, and needing a routing layer that can correctly serve both the old and new topology simultaneously.

go deeper

for a junior

Should understand that resharding involves physically moving data, not just updating a config value.

for a middle

Should be able to describe the basic idea of writing to both locations temporarily while migrating.

for a senior

Should describe the staged dual-write/backfill/verify/cutover pattern and name at least one concrete failure mode like read-your-write inconsistency.

for a principal

Should discuss choosing a sharding scheme up front to minimize future resharding pain, tooling investment, and how to sequence a large-scale resharding project with rollback checkpoints.

## What resharding is **Resharding** is the operation of changing a live sharded system's data-to-shard mapping — typically to add capacity — without taking the system offline. Mechanically, it starts with deciding the new topology (e.g., 8 shards becoming 16, or one overloaded shard being split into two), then moving the affected data from its old location to its new one while both the application and the routing layer stay aware of which rows currently live where. The naive approach — recomputing `shard = hash(key) % new_shard_count` for every row — is deceptively dangerous: because the modulus changed, the vast majority of keys resolve to a different shard number than before, even though most of them didn't need to move at all conceptually. A key that was on shard 3 out of 8 might now belong on shard 11 out of 16, with no relationship to its old placement. This means a naive resize forces something close to a full-cluster data migration rather than an incremental one, which is exactly why production systems avoid plain modulo hashing for anything that needs to grow, or why they lean on schemes (like consistent hashing, whose deeper mechanics belong to distributed-systems theory) or on shard-splitting instead. ## The staged migration The reason resharding exists as an operational discipline, rather than a one-line config change, is that a sharded system's whole purpose is to keep serving traffic continuously while its capacity grows — stopping the world to copy data would defeat much of the benefit of having sharded in the first place. So the standard approach decomposes the migration into stages that each individually preserve correctness. 1. **First, dual writes** (or an equivalent change-data-capture stream): once a target shard for a given key range is decided, every new write to that range is applied to both the old and new locations, ensuring the new location never falls further behind than "real time" from that point forward. 2. **Second, backfill**: a background job copies the pre-existing historical data for that range from old to new, typically throttled so it doesn't compete too aggressively with production traffic for IO and CPU. 3. **Third, verification**: before cutting reads over, the system (or an operator) checks that the new location's data actually matches the old location's — row counts, checksums, or spot-checks — because silently cutting over to an incomplete or diverged copy is worse than not migrating at all. 4. **Fourth, cutover**: reads for that key range are switched to the new location, usually gradually (a percentage of traffic, or one range at a time) rather than all at once, so problems surface on a small blast radius. 5. **Finally, cleanup**: once the cutover is confirmed stable, the old location's copy of that data is deleted and dual-writing stops. ## The trade-off The trade-offs of this approach are mostly about time and cost versus safety. It's slower and more operationally involved than a single atomic switch — a large resharding project can run for days or weeks — but it avoids downtime and gives multiple checkpoints where the team can pause or roll back if something looks wrong. Splitting an existing shard into two, when the sharding scheme supports it (e.g., range sharding where a shard just owns a sub-range of its old range going forward), is generally far cheaper than a full rehash, because only the data that's moving to the new shard has to move at all — the rest of the cluster's mapping is untouched. ## Failure modes The failure modes are concrete and well known operationally. - **Split-brain writes** happen if both the old and new locations accept independent writes for the same key without one being clearly authoritative — a client might read stale data from one location while another client's write only landed on the other. - **Read-your-write inconsistency** shows up during the cutover window: a user writes a row, and an immediate read is routed to whichever location hasn't yet received it, producing a confusing "my data disappeared" bug report. - **Prolonged dual-write windows** add sustained extra load and latency to every write in the migrating range, which can itself become a capacity problem if the migration runs too long. - **A router that isn't correctly aware** of exactly which key ranges have and haven't cut over yet is a common source of subtle, hard-to-reproduce bugs during the transition. ## Where it shows up Real systems build dedicated tooling for exactly this workflow rather than treating it as a one-off script. - **Vitess** (used for sharding MySQL at YouTube and Slack) has a formal "resharding" workflow built on VReplication that automates the dual-write/backfill/cutover sequence with built-in consistency checks. - **MongoDB's** sharded clusters run a background balancer that migrates data "chunks" between shards to keep load even, using a similar copy-then-cutover mechanism at the chunk level. - **Citus** (a sharding extension for PostgreSQL) provides shard rebalancing commands that move individual shards between nodes using logical replication to avoid downtime. The common thread across all of them is that resharding is treated as a first-class, tooled operation with explicit safety stages, precisely because doing it by hand with an ad hoc script is one of the highest-risk operations you can run against a production sharded database.

  • Why is naive modulo hashing (key % N) especially bad for resharding compared to schemes designed to minimize key movement?
    Because changing N remaps nearly every key — key % 8 and key % 16 agree for only a small fraction of values — so a resize forces close to a full-cluster data migration instead of an incremental one. The deeper mechanics of how to minimize key movement on resize (consistent hashing) belong to distributed-systems partitioning theory, but operationally the takeaway is: plan for a large migration, or choose a scheme up front that avoids this, rather than assuming a shard-count change is a cheap operation.
  • What's the difference between splitting an existing shard in two versus rehashing the whole keyspace to a new shard count?
    Splitting takes one shard's existing key range and divides it into two, moving only that shard's data to a new location while the rest of the cluster's mapping is untouched. A full rehash to a new modulus typically reassigns keys across the entire cluster at once. Splitting is the far less disruptive, operationally preferred approach for incremental growth when the sharding scheme supports it.

It's like renovating a building floor by floor while tenants keep living in it, instead of evacuating everyone and rebuilding overnight — you set up the new floor, move belongings over in the background, double-check nothing was left behind, then redirect mail to the new address, all before tearing down the old floor.

saying these in an interview costs you the question

  • Thinks resharding is just a config change with no data movement
  • Doesn't mention any staged process like dual-write/backfill/cutover
  • Assumes zero downtime happens automatically without operational work
  • Unaware that naive mod-N hashing forces a near-total key remap on resize
  • No mention of verification/consistency checking before cutover

context