skip to content

Walk through how you move a slice of data from one shard to another while the application keeps reading and writing, with no downtime. What are the steps and where does correctness break?

level: seniorimportance: must knowfreq 45%

answer

  1. Move buckets/chunks, never the whole shard at once
  2. Snapshot at a log position, stream from that exact position
  3. Backfill must upsert-if-newer or it clobbers live writes
  4. Fence on the source: stale routers must fail, not write
  5. Rollback window = source kept read-only, expires once target diverges

basics

~20 s

Copy the slice's existing rows to the target, keep applying ongoing changes until lag is near zero, briefly fence writes for that slice, verify with checksums, flip the routing directory to the new shard, then keep the source read-only as a rollback path before deleting it.

solid answer

~60 s

Per bucket or chunk, not for the whole shard at once: 1. **Bootstrap** — snapshot-copy existing rows to the target in bounded batches, throttled. 2. **Catch up** — stream ongoing changes (logical replication or change data capture) into the target until replication lag is milliseconds. Alternatively dual-write from the application, but that has no atomicity across the two databases. 3. **Verify** — row counts and range checksums per key range, source versus target, while still live. 4. **Fence and cut over** — reject or briefly queue writes for *that slice only* (a per-bucket status such as `MIGRATING`), drain in-flight writes, confirm zero lag, bump the routing directory to a new version with the bucket now pointing at the target. 5. **Settle** — routers with a stale directory version get fenced by the source, which now refuses writes for the moved bucket. Keep the source copy read-only for a rollback window, then delete. Breaks: backfill overwriting newer live writes (apply upsert-if-newer with a version), stale router caches double-writing, non-atomic dual-writes, id collisions on the target, and unique constraints that were only per-shard.

code

sql · 7 lines
sql
INSERT INTO orders (tenant_id, order_id, status, row_version, updated_at)
VALUES (:tenant_id, :order_id, :status, :row_version, :updated_at)
ON CONFLICT (tenant_id, order_id) DO UPDATE
   SET status      = EXCLUDED.status,
       row_version = EXCLUDED.row_version,
       updated_at  = EXCLUDED.updated_at
 WHERE orders.row_version < EXCLUDED.row_version;

go deeper

for a junior

Know the shape: copy the data, keep the copy up to date with ongoing changes, verify, briefly pause writes for that slice, switch the routing, then clean up.

for a middle

Explain snapshot-plus-change-stream anchored at a log position, per-bucket cutover, and why the backfill needs a version-guarded upsert.

for a senior

Own the failure modes and the runbook: fencing on the source, versioned directory, checksum verification, throttling against production load, id and unique-constraint hazards, and a defined rollback window with an explicit expiry.

for a principal

Judge the whole programme — that migrations are resumable, pausable and observable, that a bucket is the unit of blast radius, and that the organization can run this continuously rather than as a one-off heroic event.

## Framing Resharding is not one big copy. It is a loop over small units — a **bucket** in a hash-plus-directory scheme, a **chunk** in a range-sharded system — each moved independently, verified independently, and individually roll-backable. Doing it per unit is what makes the operation resumable, throttleable, and survivable. ## Step 0 — Preconditions Every access must go through a routing layer that reads a **versioned directory**; if application code can reach a shard by hard-coded connection, you cannot cut over safely. Freeze schema changes for the tables involved. Confirm you can throttle the copy so it does not starve live traffic, and decide the rollback criterion before starting. ## Step 1 — Bootstrap copy Copy the slice's existing rows to the target in ordered, bounded batches (key ranges, not `OFFSET`), with a sleep or a token-bucket between batches so replicas do not fall behind and the source's I/O headroom survives. Record progress so a crashed copier resumes rather than restarts. Do not index-build on the target before the copy if you can build after — bulk load first, indexes second, is usually much faster. ## Step 2 — Catch up on live changes The copy is stale the moment it starts, so ongoing mutations must reach the target. Two mechanisms: **Change stream (preferred).** Logical replication or change data capture reads the source's write-ahead log and applies the changes to the target. It is ordered, complete, and needs no application change. Ordering against the bootstrap matters: take the snapshot at a known log position and start streaming from that exact position, otherwise you either lose changes or replay them out of order. **Dual-write from the application.** The application writes to both shards. It is simple to explain and treacherous in practice: there is no atomicity across two databases, so a crash between the two writes leaves them diverged; concurrent writers can apply in different orders on each side; and every write path in the codebase must be found and changed. If you use it, make every write idempotent and carry a monotonic version so the target can reject stale applications, and run the reconciler continuously rather than once. Either way the backfill and the live stream will race on the same rows. The rule that saves you: **apply-if-newer**. Every row carries a version or updated-at, and both the backfill and the stream write with an upsert that refuses to overwrite a higher version. Without it, a slow backfill batch happily stamps a stale copy over a fresh live write, and nothing in the system notices. ## Step 3 — Verify before you trust While both sides are live, run comparisons per key range: row counts, then checksums of the concatenated row images (excluding volatile columns), then targeted row-level diffs on any mismatching range. Re-run until stable. Mismatches early in the process are normal, because of lag; mismatches that persist after lag reaches zero are bugs, and they are the reason to abort. ## Step 4 — Cut over This is the only moment with user-visible impact, and it should last a fraction of a second per bucket. 1. Mark the bucket `MIGRATING` in the directory. Routers begin queueing or rejecting-with-retry writes for **that bucket only** — every other bucket is unaffected, which is the entire reason for moving small units. 2. Drain: wait for in-flight transactions on the source to finish. 3. Confirm the change stream has applied everything up to the source's current log position, i.e. lag is zero. 4. Publish a new directory version with the bucket owned by the target, and mark it `ACTIVE`. 5. Routers refresh. Crucially, the **source must fence**: it rejects writes for a bucket it no longer owns, so a router running a stale directory version fails loudly instead of silently writing to a dead copy. Fencing on the owner, not just on the client, is what makes cutover safe — client caches always lag. ## Step 5 — Settle and clean up Keep the source's copy intact but read-only for a defined window. If a defect surfaces, rollback is a directory flip back — which is only valid while the source is still a faithful copy, so rollback stops being available the moment the target has taken writes that were never replicated back. Be explicit about that expiry. Then delete the source rows, in throttled batches, and reclaim the space. ## Where correctness actually breaks - **Backfill clobbering live writes.** The single most common bug. Fixed by version-guarded upserts, never by "the backfill is fast enough". - **Stale router caches.** Solved by owner-side fencing plus a monotonically increasing directory version, not by hoping caches expire. - **Non-atomic dual-writes.** One side committed, the other not. Reconciliation is mandatory if you take this path. - **Identifier collisions.** Database-local sequences on the target can generate ids that already exist in the migrated rows. Use globally unique ids, or bump the target's sequence past the imported maximum before it takes traffic. - **Constraints that were only per-shard.** A unique index enforced within the source shard does not stop the target from already holding a conflicting row. Check globals before the move, not during. - **Foreign keys and co-location.** Related tables must move as a set, or a co-located join silently becomes a cross-shard one. - **Lag-induced stale reads.** If reads are switched before writes, or vice versa, you get read-your-writes violations. Switch both at the same directory version. - **Copy starving production.** Throttle, watch replica lag and lock waits, and be willing to pause the migration for days. ## What to monitor Replication lag per bucket, copy throughput, verification mismatch counts, cutover fence duration, error rate on the fenced bucket, and the age of the rollback window. A reshard that cannot be paused, resumed, and rolled back is not a plan — it is an outage with a schedule.

  • Why is dual-writing from the application riskier than streaming changes from the write-ahead log?
    Dual-write has no atomicity across the two databases, so a crash or a timeout between the two writes leaves them permanently diverged, and concurrent writers can land in different orders on each side. It also requires finding and changing every write path in the codebase, including background jobs and manual fixes. A log-based change stream is ordered, complete by construction, needs no application change, and can be anchored to the exact snapshot position so nothing is lost or replayed out of order.
  • A router still holding the old directory version sends a write to the source after cutover. What stops the data from being lost?
    The source fences: it knows the current directory version and refuses any write for a bucket it no longer owns, returning an error that tells the client to refresh. The router reloads the directory and retries against the new owner. Relying on cache expiry instead would silently accept the write into a copy nobody reads again, so enforcement has to live on the owner side, keyed to a monotonically increasing version.
  • How long can you keep the option to roll back, and what ends it?
    Rollback is a directory flip back to the source, and it is valid only while the source remains a faithful copy of the data. The moment the target accepts writes that are not replicated back to the source, flipping back would lose them. So either you keep reverse replication running for the rollback window, or you state explicitly that rollback expires at cutover plus a fixed period and treat anything after that as roll-forward only.

It is moving one filing cabinet at a time: photocopy the folders, keep forwarding new mail to the new location, lock that one cabinet for a few seconds, update the office directory, and keep the old cabinet locked-but-intact for a week in case something was missed.

saying these in an interview costs you the question

  • Planning one big copy of an entire shard instead of moving small, independently verifiable units
  • Assuming the backfill is fast enough that it cannot overwrite a newer live write
  • Relying on client-side directory cache expiry rather than fencing writes at the source
  • Cutting over reads and writes at different times, which breaks read-your-writes
  • Forgetting that local sequences on the target can collide with imported identifiers
  • Deleting the source data immediately after cutover, destroying the rollback path

context