Prompt 17 · Lesson 05 · Evolving a live distribution
Refine, Reshard, Unshard/Move Collections Where Supported, and Plan Safe Distribution Changes
MongoDB can refine, reshard, unshard, or move collections, but each operation changes a different contract and carries different data-movement, write-blocking, and rollback costs.
Learning objectives
Distinguish refine, reshard, unshard, and moveCollection by the state change each performs and the cases each can solve.
Run a safe shard-key refinement and verify that it appends a suffix rather than replacing existing fields.
Use reshardCollection demoMode on disposable data, inspect progress/current key, and explain abort-versus-commit boundaries.
Plan unshardCollection and moveCollection with their resource/write-blocking constraints instead of treating them as instant rollback buttons.
Create a migration decision record with capacity, zone, Search/encryption, observability, cutover, and rollback/rebuild checks.
This lesson pins MongoDB Community Server
8.3.8 with
mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim, mongosh 2.10.0, and PyMongo
4.17.0 where the driver is used. The local topology
is a disposable single-host sharded cluster with one
single-member config server replica set, two
single-member shard replica sets, and one
mongos, exposed only on loopback diagnostic ports
27151–27154. This is sufficient for routing/distribution
mechanics but is not production high availability.
Authentication and TLS are disabled only for the isolated lab.
Feature Compatibility Version (FCV) is inspected in the lab. FCV
is observed and never changed. Application/DDL operations go
through mongos; direct shard connections are
diagnostics only. The mandatory fast lab executes
refineCollectionShardKey and a tiny
reshardCollection with demoMode:true.
unshardCollection and
moveCollection are taught as optional long-running
procedures because current documentation imposes heavier
resource and write-blocking requirements and a five-minute
minimum for unsharding. Atlas Free/Flex restrictions are called
out separately. Atlas, Search, Vector Search, KMS, and
Enterprise Advanced are not mandatory. Product commands were not
executed in this generation environment because Docker, mongod,
mongos, mongosh, and PyMongo are unavailable here, so all
timing/load outputs must be measured on the learner machine
rather than copied as invented results.
1. Four tools solve four different distribution problems
| Operation | What changes | Redistributes existing data? | Good fit |
|---|---|---|---|
refineCollectionShardKey |
Appends suffix field(s) to existing shard key; existing field order/types stay intact. | Not a full rebuild/redistribution like resharding. | Need more cardinality/split precision while retaining the existing prefix contract. |
reshardCollection |
Replaces the shard key (or redistributes on the same key with forceRedistribution). | Yes, through donor/recipient resharding phases. | Current key/distribution no longer meets workload needs. |
unshardCollection |
Removes sharding and consolidates the collection onto one shard. | Yes, consolidates to a recipient shard. | Collection no longer benefits from sharding and one shard can safely own it. |
moveCollection |
Moves one unsharded collection to another shard. | Yes, changes placement of unsharded collection. | Operational placement change after/without sharding. |
After resharding commits, the old routing table is not an automatic rollback switch. You can start another controlled reshard later, restore from tested backup where appropriate, or use an application migration strategy—but each is another operation with its own risk. Plan rollback before starting.
docker rm -f atlasmart-ch17-l5-cfg atlasmart-ch17-l5-s1 atlasmart-ch17-l5-s2 atlasmart-ch17-l5-mongos 2>/dev/null || truedocker network rm atlasmart-ch17-l5-net 2>/dev/null || truedocker volume rm atlasmart-ch17-l5-cfg-data atlasmart-ch17-l5-s1-data atlasmart-ch17-l5-s2-data 2>/dev/null || truedocker network create atlasmart-ch17-l5-netdocker run -d --name atlasmart-ch17-l5-cfg --network atlasmart-ch17-l5-net -p 127.0.0.1:27151:27017 -v atlasmart-ch17-l5-cfg-data:/data/db mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim --configsvr --replSet atlasmart-cfg17-l5 --bind_ip_alldocker run -d --name atlasmart-ch17-l5-s1 --network atlasmart-ch17-l5-net -p 127.0.0.1:27152:27017 -v atlasmart-ch17-l5-s1-data:/data/db mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim --shardsvr --replSet atlasmart-shard17-l5-a --bind_ip_alldocker run -d --name atlasmart-ch17-l5-s2 --network atlasmart-ch17-l5-net -p 127.0.0.1:27153:27017 -v atlasmart-ch17-l5-s2-data:/data/db mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim --shardsvr --replSet atlasmart-shard17-l5-b --bind_ip_allfor PORT in 27151 27152 27153; do until mongosh "mongodb://127.0.0.1:$PORT/admin?directConnection=true" --quiet --eval 'quit(db.runCommand({ping:1}).ok===1?0:1)'; do sleep 1; donedonemongosh "mongodb://127.0.0.1:27151/admin?directConnection=true" --quiet --eval 'rs.initiate({_id:"atlasmart-cfg17-l5",configsvr:true,members:[{_id:0,host:"atlasmart-ch17-l5-cfg:27017"}]})'mongosh "mongodb://127.0.0.1:27152/admin?directConnection=true" --quiet --eval 'rs.initiate({_id:"atlasmart-shard17-l5-a",members:[{_id:0,host:"atlasmart-ch17-l5-s1:27017"}]})'mongosh "mongodb://127.0.0.1:27153/admin?directConnection=true" --quiet --eval 'rs.initiate({_id:"atlasmart-shard17-l5-b",members:[{_id:0,host:"atlasmart-ch17-l5-s2:27017"}]})'for PORT in 27151 27152 27153; do until mongosh "mongodb://127.0.0.1:$PORT/admin?directConnection=true" --quiet --eval 'quit(db.hello().isWritablePrimary?0:1)'; do sleep 1; donedonedocker run -d --name atlasmart-ch17-l5-mongos --network atlasmart-ch17-l5-net -p 127.0.0.1:27154:27017 --entrypoint mongos mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim --configdb atlasmart-cfg17-l5/atlasmart-ch17-l5-cfg:27017 --bind_ip_all --port 27017until mongosh "mongodb://127.0.0.1:27154/admin" --quiet --eval 'quit(db.runCommand({ping:1}).ok===1?0:1)'; do sleep 1; donemongosh "mongodb://127.0.0.1:27154/admin" --quiet --eval 'printjson(sh.addShard("atlasmart-shard17-l5-a/atlasmart-ch17-l5-s1:27017"));printjson(sh.addShard("atlasmart-shard17-l5-b/atlasmart-ch17-l5-s2:27017"));printjson({hello:db.hello(),fcv:db.runCommand({getParameter:1,featureCompatibilityVersion:1}).featureCompatibilityVersion,shards:db.adminCommand({listShards:1}).shards});'
2. Start with a refine: append a suffix without pretending it redistributes everything
const app=db.getSiblingDB("atlasmart");sh.enableSharding("atlasmart", "atlasmart-shard17-l5-a");app.orders_evolve.drop();app.orders_evolve.createIndex({tenantId:1});printjson(sh.shardCollection("atlasmart.orders_evolve",{tenantId:1}));const docs=[];for (let i=1;i<=500;i++) docs.push({tenantId:i<=250?"mega":`t${i%25}`,orderId:`o-${String(i).padStart(5,"0")}`,status:"paid"});app.orders_evolve.insertMany(docs);// Refine requires a supporting index for the refined key.app.orders_evolve.createIndex({tenantId:1,orderId:1});printjson(db.getSiblingDB("admin").runCommand({ refineCollectionShardKey:"atlasmart.orders_evolve", key:{tenantId:1,orderId:1}}));printjson(db.getSiblingDB("config").collections.findOne({_id:"atlasmart.orders_evolve"},{key:1,uuid:1}));
Refinement preserves the existing tenantId prefix
and adds orderId. It creates more possible split
points inside a frequent tenant value, but it does not magically
redesign routing for queries that never include
tenantId. MongoDB explicitly warns not to change
the range/hashed type of existing key fields during refinement;
use resharding for a genuinely different key.
3. Reshard on disposable data and observe the state transition
const app=db.getSiblingDB("atlasmart");app.orders_evolve.createIndex({orderId:"hashed"});printjson(db.getSiblingDB("admin").runCommand({ reshardCollection:"atlasmart.orders_evolve", key:{orderId:"hashed"}, demoMode:true}));printjson(db.getSiblingDB("config").collections.findOne({_id:"atlasmart.orders_evolve"},{key:1,uuid:1}));printjson(sh.getShardedDataDistribution());
demoMode:true exists for testing/demonstration and
bypasses the normal minimum resharding duration. In production,
resharding goes through coordinator, clone, index, oplog
catch-up, critical-section, and commit phases. Recipient shards
need substantial extra disk/I/O/CPU/oplog headroom. Search
indexes must be considered separately because current
documentation requires rebuilding Search after resharding
completes.
printjson(db.getSiblingDB("admin").aggregate([ {$currentOp:{allUsers:true,localOps:false}}, {$match:{opStatus:{$exists:true}}}, {$project:{opStatus:1,ns:1,desc:1,command:1}}]).toArray());// Before the commit phase, abortReshardCollection can stop an in-progress reshard.// Once the operation reaches commit, MongoDB documents that it can no longer be aborted.
commitReshardCollection deliberately forces the
operation toward the write-blocking completion phase. It is
not a performance shortcut to use without checking the
remaining-time estimate and the application’s latency
tolerance.
4. Unshard and move are supported in MongoDB 8.0+, but they are not cheap reversals
// OPTIONAL after reading the current unshard/move requirements.// Community/self-managed supports both; Atlas Free/Flex does not support these operations.// Consolidate a sharded collection onto one shard and remove its shard key:// printjson(db.getSiblingDB("admin").runCommand({// unshardCollection:"atlasmart.orders_evolve",// toShard:"atlasmart-shard17-l5-a"// }));// After a collection is unsharded, move that entire collection to another shard:// printjson(db.getSiblingDB("admin").runCommand({// moveCollection:"atlasmart.orders_evolve",// toShard:"atlasmart-shard17-l5-b"// }));
Current unshardCollection guidance requires the
application to tolerate a brief write-blocking period and
substantial recipient headroom; it also restricts concurrent
topology/index/DDL operations. It has a five-minute minimum
duration, is not available on Atlas Free/Flex, and cannot be
used for Queryable Encryption collections.
moveCollection has similar resource/concurrency
constraints for unsharded collections. These operations should
have capacity checks, observability, abort procedures, and
post-move index/Search verification.
5. Migration decision record and rollback checkpoints
| Before start | During operation | After commit |
|---|---|---|
| Verify current key, zones, indexes, Search/QE constraints, backups, oplog window, disk/I/O/CPU headroom, app latency budget. | Monitor currentOp/progress, recipient/donor resource use, replication lag, error logs, write latency, and sampled targetability. | Verify new key/placement, query targeting, data counts/invariants, indexes, Search rebuilds, backup coverage, and remove temporary compatibility paths only later. |
| Choose abort/cutover threshold and who can execute it. | Do not add competing index builds/topology changes that current docs restrict. | Keep rollback data/configuration long enough to detect workload regressions; a committed reshard is not instant-reversible. |
Do not schedule a production reshard because a dashboard says one shard is “hot” without proving the cause. The hotspot may be an application key frequency, query scatter, low cardinality, a zone constraint, or unrelated local index/storage pressure. Change only after the evidence supports the mechanism.
6. Production judgment and chapter bridge
Refine when the existing prefix remains useful and suffix cardinality can solve the split problem. Reshard when the key itself must change or when complete redistribution is justified. Unshard only when one shard can safely absorb the collection and the application accepts the operation's constraints. Move an unsharded collection for deliberate placement, not as a substitute for load-balancing a sharded collection. Every path needs security/RBAC, TLS, backup/restore validation, resource headroom, version/FCV checks, and post-change workload measurement.
Next chapter. Distribution solves where data lives. Chapter 18 turns to how applications observe changes over time with Change Streams, resume tokens, and idempotent event-driven integration.
docker rm -f atlasmart-ch17-l5-cfg atlasmart-ch17-l5-s1 atlasmart-ch17-l5-s2 atlasmart-ch17-l5-mongos 2>/dev/null || truedocker volume rm atlasmart-ch17-l5-cfg-data atlasmart-ch17-l5-s1-data atlasmart-ch17-l5-s2-data 2>/dev/null || truedocker network rm atlasmart-ch17-l5-net 2>/dev/null || true
Check your understanding
- What can refineCollectionShardKey change?
- What is the main difference between refine and reshard?
- Can an in-progress reshard always be aborted?
- What does unshardCollection do?
- Why is moveCollection not a sharded load-balancing command?
Review the answers
1. It appends suffix field(s) to the existing shard key; it must not change the existing fields or their ranged/hashed types.
2. Refine extends the current key contract; reshard can replace the key and redistributes data under the new distribution.
3. No. MongoDB documents that once resharding reaches the commit phase it can no longer be aborted.
4. It consolidates a sharded collection onto one shard and removes the shard key.
5. It operates on an unsharded collection and relocates that whole collection to another shard.
Authoritative references
- Choose a Shard Key — Cardinality, frequency, monotonicity, query patterns, and shard-key tradeoffs.
- Troubleshoot Shard Keys — Hot ranges, uneven load, jumbo/indivisible ranges, and scatter/gather symptoms.
- analyzeShardKey — Key-characteristic and sampled read/write-distribution metrics.
- configureQueryAnalyzer — Query sampling for prospective shard-key analysis.
- Hashed Indexes for Sharding — Write distribution and range-query limitations.
- Zones — Zone membership, non-overlapping zone ranges, prefix requirements, and balancing behavior.
- addShardToZone — Associate shards with zone labels.
- Change a Shard Key — Refine versus reshard decision boundary.
- Refine a Shard Key — Append suffix fields without changing existing shard-key field types.
- reshardCollection — Full shard-key change, redistribution phases, demo mode, and abort/commit boundaries.
- unshardCollection — MongoDB 8.0+ consolidation of a sharded collection onto one shard.
- moveCollection — MongoDB 8.0+ relocation of an unsharded collection between shards.
- MongoDB Sharding — Routing, ranges, balancer behavior, and mongos-only client access.
- MongoDB 8.3 Release Notes — Current 8.3 release and sharding/DDL changes.
- mongosh Release Notes — mongosh 2.10.0 baseline.
- PyMongo Release Notes — PyMongo 4.17 baseline.