skip to content

During a MongoDB chunk migration, what happens on the donor and recipient shards, and how does live traffic feel it?

level: seniorimportance: must knowfreq 55%

answer

  1. clone, then catch up, then commit
  2. only one phase blocks writes
  3. the donor still holds the data afterwards
  4. orphans are removed on a delay
  5. confine it to a quiet window

basics

~20 s

The recipient clones the range's documents and then catches up on changes; the donor enters a short critical section that blocks writes to that range while the config metadata is committed; afterwards the donor deletes the moved documents in the background as orphans.

solid answer

~50 s

A migration runs in phases. The balancer sends `moveRange` to the donor's primary. The recipient creates any missing indexes, then **clones** the documents in the range while the donor keeps serving traffic, and applies the changes that happened during cloning until it is nearly caught up. The donor then enters a brief **critical section**: writes to that range are blocked and queued while the final changes transfer and the new ownership is committed to the config servers. Routers discover the change through stale-config errors and refresh. After the commit the donor still physically holds the documents. They are deleted asynchronously by the range deleter once no cursors reference them and `orphanCleanupDelaySecs` has elapsed (default 900 seconds). The workload feels this as read and write amplification on both shards, extra oplog and replication traffic, cache churn, and a short write stall on the migrating range. Mitigate with a balancer active window, `_secondaryThrottle`, or disabling balancing on a collection during a bulk load.

code

javascript · 12 lines
javascript
// confine migrations to a quiet period
use config
db.settings.updateOne(
  { _id: "balancer" },
  { $set: { activeWindow: { start: "01:00", stop: "05:00" } } },
  { upsert: true }
)

// pause balancing for one collection during a bulk load
sh.disableBalancing("app.events")
// ... load ...
sh.enableBalancing("app.events")

go deeper

for a junior

Know that MongoDB moves data between shards in the background, that it is copied before ownership changes, and that this competes with normal traffic for disk and memory.

for a middle

Be able to name the phases in order — clone, catch up, commit, background cleanup — and say which single phase blocks writes and to what scope.

for a senior

Demonstrate that you would schedule migrations into a window, throttle them if secondaries lag, and diagnose with config.changelog and the range-deletion backlog rather than switching the balancer off.

for a principal

Treat migration bandwidth as a capacity input: the one-migration-per-shard limit sets how fast a cluster can absorb a new shard, and that number belongs in the growth plan, not in an incident postmortem.

## The phases of a range migration 1. The balancer, on the config-server primary, sends `moveRange` (historically `moveChunk`) to the **donor** shard's primary. 2. The donor tells the **recipient** to start receiving. The recipient builds any indexes it is missing for that collection. 3. The recipient **clones** the documents whose shard-key values fall in the range. The donor keeps serving reads and writes normally throughout. 4. Writes that land in the range while cloning is in progress are tracked and shipped to the recipient in a catch-up phase, repeated until the recipient is nearly in sync. 5. The donor enters the **critical section**: incoming writes to that range are blocked and queued, the last changes are transferred, and the ownership change is committed to the config servers. 6. Routers and shards that still believe the old owner holds the range get a stale-config error, refresh their routing table, and retry. 7. The documents remain on the donor as **orphans** until the range deleter removes them — after any open cursors on the range finish, and after `orphanCleanupDelaySecs` (default 900 seconds). Pending deletions are tracked in each shard's own `config.rangeDeletions` collection. ## What the workload actually experiences **A short write stall on that range.** Only step 5 blocks writes, and only for documents in the range being moved. It is normally sub-second, but it is real: on a hot range under heavy write load, or when the donor's primary is struggling, the queued writes show up as a latency spike. **Read and write amplification on both shards.** The donor reads the range out; the recipient writes it in. Those recipient writes go through its oplog and replicate to its secondaries, so a migration inflates oplog volume and can push secondaries behind. The donor's later range deletion is another burst of write work and IOPS. **Cache churn.** Scanning a range on the donor and inserting it on the recipient evicts hot working-set pages on both sides, so query latency can degrade for minutes after a big move even though nothing is blocked. **Metadata refresh chatter.** Every router that had the old ownership cached takes a stale-config error and refreshes. With many `mongos` processes and continuous migrations, this is visible though rarely dominant. **Concurrency limits.** A shard participates in at most one migration at a time, so a cluster of *n* shards runs at most about *n*/2 migrations concurrently. That is a hard ceiling on how fast a cluster can rebalance, and it is why adding a shard to a very large cluster can take days to fill. ## Controlling when and how hard it hits **Balancer active window.** The most common control: confine migrations to a quiet period by upserting an `activeWindow` into the balancer settings document. The balancer will not *start* new migrations outside the window. ```javascript use config db.settings.updateOne( { _id: "balancer" }, { $set: { activeWindow: { start: "01:00", stop: "05:00" } } }, { upsert: true } ) ``` **Per-collection control.** `sh.disableBalancing("app.events")` stops migrations for one collection while leaving the rest of the cluster balanced — useful during a bulk load or a heavy backfill. `sh.enableBalancing()` restores it. **Cluster-wide stop.** `sh.stopBalancer()` and `sh.startBalancer()`, with `sh.isBalancerRunning()` to confirm the current round has finished before you begin maintenance such as a backup. **Secondary throttling.** The `_secondaryThrottle` setting makes each migration write on the recipient wait for a write concern rather than firing as fast as the disks allow, trading migration speed for replication headroom. `_waitForDelete` makes the balancer wait for the donor's range deletion to finish before starting the next migration, which keeps the orphan backlog from growing. ## Diagnosing migration pain Look at `config.changelog` for `moveChunk.start` / `moveChunk.commit` entries to see migration volume and duration. Check each shard's `config.rangeDeletions` for a growing cleanup backlog. Correlate replication lag on the recipient's secondaries with migration timestamps. If latency spikes line up with migrations, the fix is scheduling and throttling — not turning the balancer off permanently, which just lets skew accumulate until the eventual rebalance is far more painful.

  • After a migration commits, can a query routed through mongos see the leftover documents on the donor?
    No. Each shard keeps filtering metadata describing which ranges it owns and excludes documents outside those ranges from results, so orphans are invisible to routed queries. They still consume disk and cache until the range deleter removes them, which is why a large orphan backlog is an operational problem even though it is not a correctness one.
  • How would you keep a large bulk load from fighting the balancer?
    Disable balancing for just that collection with sh.disableBalancing("db.coll") for the duration, and pre-split the range space so the data lands spread out instead of piling onto one shard and then being migrated away. Re-enable afterwards and let the balancer clean up any residual skew inside a maintenance window.
  • What does _secondaryThrottle change about a migration?
    It makes the migration's writes on the recipient wait for a configured write concern instead of proceeding at full speed. Migrations then take longer but generate replication pressure the recipient's secondaries can keep up with, which protects read preferences that target secondaries and keeps lag-based alerts quiet.

saying these in an interview costs you the question

  • Says writes to the whole collection are blocked for the entire migration
  • Believes the donor deletes the documents synchronously at commit
  • Thinks orphaned documents are returned by queries through mongos
  • Claims migrations never touch the oplog or replication
  • Leaves the balancer permanently disabled as the fix for latency spikes

context