As a sharded system grows and needs to add capacity, what does the process of resharding (changing the number of shards or moving data between them) actually involve, and what makes it risky to do on a live production system?
answer
- backfill plus CDC/dual-write plus cutover
- atomic routing flip
- keep old shard as rollback window
- consistent hashing minimizes data moved
- verification/checksum before cutover
basics
~20 sResharding means moving data around to add more machines to a sharded database, similar to redoing a filing system while people are still using the files. It's risky because data can get lost, duplicated, or briefly unreachable while it's being moved, so it has to be done carefully in steps.
solid answer
~60 sResharding is the process of changing shard count or key-to-shard mapping while the system stays live - it requires migrating a subset of data to new/target shards, redirecting reads/writes to the new location at the right moment, and doing so without losing writes that happen mid-migration or serving stale/duplicate reads. The standard pattern is backfill plus dual-write or change-data-capture: start replicating writes for the moving key range to both old and new locations (or capture changes via CDC), backfill historical data to the target, verify the target has caught up, then atomically cut over routing for that range and stop writing to the old location, keeping it around briefly as a rollback path. Consistent hashing minimizes how much data needs to move per rebalance; directory-based or Vitess-style systems automate much of this with online, incremental resharding tooling. The main risks are split traffic during cutover, needing spare capacity to run old and new in parallel during migration, and the sheer duration for very large key ranges, during which the migration must keep up with ongoing production write volume.
go deeper
Understands that adding shards means data has to move somewhere and that this doesn't happen automatically.
Can describe backfill and cutover at a basic level and knows migrations take real time.
Designs a CDC-based migration with a verification step and understands the cutover race-condition risk.
Plans org-scale resharding strategy: picks partitioning schemes upfront that minimize future migration cost, sets standard tooling/runbooks for reshard operations, and budgets spare capacity and rollback policy across the platform.
## What resharding is Resharding is the operation of changing how a sharded system's key space maps onto physical shards - adding shards to handle growth, retiring underused shards, or fixing an unbalanced distribution discovered after the fact - without taking the system offline. It sounds like a bookkeeping change, updating a mapping table or changing a modulus, but it is actually one of the **highest-risk operations** in a sharded system's lifecycle, because it means physically moving a large volume of live, actively-mutating data from one machine to another while the application keeps reading and writing it. ## The mechanical core, step by step The mechanical core of a safe online resharding process is a variant of the same pattern regardless of the underlying technology: **backfill plus dual-write or change-data-capture, plus cutover.** 1. **First,** you identify the subset of data that needs to move - a key range, a hash bucket, or a specific tenant's rows. 2. **Second,** you begin capturing every new write to that data as it happens, either by having the application dual-write to both the old and new location, or, the more robust and now more common approach, by tailing a change-data-capture stream off the source shard's write-ahead log or binlog (tools like Debezium, or built-in mechanisms in Vitess's VReplication) so every committed write is captured without relying on the application to remember to write twice. 3. **Third,** you run a backfill: a bulk copy of the existing historical data for that key range from the old shard to the new one, typically throttled so it doesn't starve production traffic of I/O capacity. 4. **Fourth,** once the backfill completes, the new shard is replaying live changes and catching up to the source in real time - you wait until the replication lag between old and new is effectively zero and stays there. 5. **Fifth,** the actual cutover: atomically flip the routing layer (the directory table, the consistent-hash ring assignment, or the proxy's shard map) so new queries for that key range go to the new shard instead of the old one. 6. **Sixth,** you keep the old shard's copy of the data around, read-only, for a rollback window in case something is wrong with the new shard, before finally decommissioning it. ## Why it stays risky The reason this is risky, even when every step above is followed carefully, comes down to a handful of concrete failure scenarios. - **The cutover moment itself is the sharpest edge.** Routing changes are rarely instantaneous across every client, proxy, and cache simultaneously, so there's a window where some requests are still being routed to the old shard while others have already moved to the new one - a write that lands on the old shard right after cutover can be silently lost if nothing is reading from there anymore, or a read immediately after cutover can miss a write that hadn't finished replicating from old to new. Systems mitigate this by making the cutover itself very short (flip a single routing key/version rather than reconfiguring every client individually) and by briefly pausing or queuing writes to the migrating range during the actual flip, accepting a tiny latency blip in exchange for correctness. - **A second risk is capacity.** During migration, both the old and new shard are live and serving/receiving traffic for the same logical data, and the source shard is additionally paying the cost of the backfill and CDC stream, so real spare I/O and CPU headroom is needed on the shard being migrated away from - attempting a migration on an already-saturated shard risks the backfill never catching up to incoming write volume, or degrading production latency for unrelated data still living on that same physical shard. - **A third risk is scale and duration.** Migrating a shard holding terabytes of data can take hours to days, during which the CDC/replication pipeline must remain reliable - any gap in the change stream, a dropped connection, or a restart that loses buffered events can silently desynchronize the target from the source in a way that isn't caught until data is queried and found missing or stale, which is why migrations are typically paired with a verification/checksum pass comparing row counts or hashes between source and target before cutover, not just trusting the pipeline. ## Why the earlier partitioning choice pays off here This is exactly why the choice of partitioning scheme made much earlier pays off or costs dearly at this stage: **consistent hashing** was specifically designed so adding or removing one shard only requires moving the data adjacent to that one change on the hash ring, rather than reshuffling the entire dataset the way naive `hash(key) mod N` does - this difference alone can be the difference between a resharding operation that takes hours and one that takes weeks. - Systems like **Vitess** automate much of this backfill-CDC-cutover choreography as a first-class 'reshard' workflow so operators don't hand-roll it per migration. - **DynamoDB and Bigtable** hide the mechanism entirely, splitting hot partitions automatically behind the scenes as part of normal operation rather than exposing resharding as a manual operator task at all.
- Why is change-data-capture (CDC) generally considered more robust than application-level dual-writes for resharding migrations?Dual-writes depend on every code path remembering to write to both locations and handling partial failure correctly, which is easy to get wrong across a large codebase. CDC instead taps the source shard's own durable write-ahead log or binlog, so it captures every committed write automatically regardless of application code, eliminating an entire class of missed-write bugs.
- What's the purpose of the 'rollback window' - keeping the old shard around, read-only, after cutover instead of decommissioning it immediately?It gives operators a fast recovery path if the new shard turns out to have a problem, such as missed data or unexpected load behavior, discovered shortly after cutover. They can revert routing back to the old shard, which still has a complete and current copy of the data, rather than reconstructing it from backups or replaying the CDC stream from scratch.
- Why does consistent hashing typically make resharding cheaper than a plain hash(key) mod N scheme, in concrete terms?With mod N, changing the shard count changes the modulus, which changes the target shard for the vast majority of keys system-wide, forcing almost a full data reshuffle. With consistent hashing, shards and keys sit on a ring, and adding or removing one shard only affects the keys immediately adjacent to it on the ring, so only a small, bounded fraction of data needs to move.
Resharding a live database is like renovating a highway from four lanes to six while cars keep driving on it - you build the new lanes alongside the old ones, gradually merge traffic over, and only tear out the old lanes once you're sure every car is safely using the new ones; do the merge too abruptly and cars, meaning writes, get lost or collide.
saying these in an interview costs you the question
- thinks resharding is a config change with no data movement
- proposes a hard cutover with no backfill/replication catch-up step
- doesn't mention keeping a rollback path after cutover
- assumes dual-writes are equivalent in reliability to CDC-based migration
- doesn't account for spare capacity needed during migration