Coordinate shard data with config metadata, balancing, write quiescence, and destination sharding so cluster recovery preserves routing correctness.
Sharded-Cluster Backup Consistency, Balancer/Metadata Concerns, and Restore Ordering
Teach self-managed sharded backup and restore as a coordinated metadata, balancer, lock, shard-key, and restore-order problem rather than independent shard copies.
Learning objectives
Explain why sharded backup consistency spans shard data, config metadata, balancing, schema operations, and cross-shard transactions.
Use the current mongos fsync-lock workflow
safely and understand its blast radius.
Capture and inspect shard-key metadata needed to reconstruct destination sharding before data restore.
Restore through mongos with the correct
ordering and avoid destroying pre-sharded metadata with
--drop.
Know when a coordinated managed backup service is required to preserve cross-shard transactional guarantees while writes continue.
This chapter pins
MongoDB Community Server 8.3.8 with
mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim, mongosh 2.10.0,
PyMongo 4.17.0 where application checks are useful,
and MongoDB Database Tools 100.18.0 for
mongodump, mongorestore, and
bsondump. Lesson 4’s mandatory path is a
deterministic metadata/restore-order simulator because a
config-server replica set plus multiple shard replica sets can
exceed a learner laptop’s resources. Host exposure is
loopback-only on
simulator only by default; optional lab may reuse the Chapter
16 sharded topology. Authentication and TLS are disabled only for disposable local
labs; production security remains the Chapter 22 prerequisite.
Default read/write concern and primary read preference are used
unless a step states otherwise.
FCV is observed and never changed. Atlas,
Enterprise Advanced, Search, Vector Search, and KMS are optional
unless a lesson explicitly labels them. If the optional real lab
is run, all components are disposable and isolated. Do not run
cluster-wide fsyncLock against a valuable
deployment. Runtime backup/restore labs were not executed in the
generation environment, so artifact sizes, restore durations,
RPO/RTO measurements, snapshot times, and checksums must be
recorded locally rather than copied as invented output.
1. AtlasMart problem: shard files without routing metadata are not a cluster backup
A sharded cluster’s data is distributed across shard replica sets, while the config server replica set stores the metadata that maps ranges to shards. During balancing, chunks/ranges can move. During resharding or DDL, collection metadata changes. A backup that captures shard A before a migration and shard B after the same migration can contain duplicate or missing logical data.
MongoDB’s documented self-managed dump procedure therefore
establishes a quiescent window: stop the balancer, stop
writes/schema transformations, lock via mongos,
take the dump, unlock, and restart balancing. Coordinated
services such as Atlas, Cloud Manager, or Ops Manager are the
supported alternatives when backups must preserve cross-shard
transactional atomicity while normal writes continue.
2. Treat the backup as a topology manifest
manifest={ "cluster":"atlasmart-sharded", "balancerStopped":True, "writesStopped":True, "schemaTransformsStopped":True, "fsyncLocked":True, "collections":{ "atlasmart.orders":{"key":{"tenantId":1,"orderId":1}}, "atlasmart.inventory":{"key":{"sku":"hashed"}}, }, "shards":["shardA","shardB"], "configMetadataCaptured":True,}required=["balancerStopped","writesStopped","schemaTransformsStopped","fsyncLocked","configMetadataCaptured"]missing=[k for k in required if not manifest.get(k)]assert not missing, f"backup preconditions failed: {missing}"for ns,meta in manifest["collections"].items(): assert meta.get("key"), f"missing shard key for {ns}"print("manifest accepted", manifest["collections"])
This simulator does not replace cluster backup; it forces the learner to name the exact consistency preconditions and the shard-key metadata needed during restore.
3. Optional real backup sequence through mongos
Starting with the documented supported server lines,
fsync/fsyncUnlock can run from
mongos. The following commands are shown only for
an isolated sharded lab. Verify the balancer actually stops and
no application writes/schema transformations remain before
locking.
sh.stopBalancer();while (sh.isBalancerRunning().mode !== "off") { sleep(1000); }printjson({balancer:sh.getBalancerState()});printjson(db.getSiblingDB("config").collections.find( {_id:/^atlasmart\./}, {_id:1,key:1,uuid:1}).toArray());printjson(db.getSiblingDB("admin").fsyncLock());printjson(db.getSiblingDB("admin").aggregate([ {$currentOp:{}}, {$facet:{ locked:[{$match:{fsyncLock:{$exists:true}}}], unlocked:[{$match:{fsyncLock:{$exists:false}}}] }}, {$project:{locked:{$gt:[{$size:"$locked"},0]},unlocked:{$gt:[{$size:"$unlocked"},0]}}}]).toArray());
rm -rf /tmp/atlasmart-ch25-l4-sharded-dumpmongodump \ --uri="mongodb://127.0.0.1:<MONGOS_PORT>/" \ --out=/tmp/atlasmart-ch25-l4-sharded-dumpfind /tmp/atlasmart-ch25-l4-sharded-dump -type f -printf '%P %s bytes\n' | sort | head -100
printjson(db.getSiblingDB("admin").fsyncUnlock());sh.startBalancer();printjson({balancer:sh.getBalancerState()});
Production automation needs a finally-style
unlock/start-balancer path and an operator runbook. A script
that exits after locking the cluster can create an outage even
if the backup files themselves are good.
4. Restore ordering: rebuild sharding metadata, then restore data through mongos
The current documented database-dump restore flow does
not blindly overwrite the destination
config database. Stop the destination balancer,
read the source shard keys from live/captured
config.collections metadata, create/pre-shard the
destination collections with the same keys, and then run
mongorestore through mongos while
excluding config.*.
sh.stopBalancer();sh.enableSharding("atlasmart");sh.shardCollection("atlasmart.orders", {tenantId:1,orderId:1});sh.shardCollection("atlasmart.inventory", {sku:"hashed"});
mongorestore \ --uri="mongodb://127.0.0.1:<TARGET_MONGOS_PORT>/" \ --nsExclude='config.*' \ /tmp/atlasmart-ch25-l4-sharded-dump
--drop casually.
The documented sharded restore warns that dropping a collection after it has been pre-sharded removes its sharding configuration. Restore ordering is part of correctness, not cosmetic procedure.
5. Validate routing and business state before reopening traffic
sh.status();const d=db.getSiblingDB("atlasmart");printjson({orders:d.orders.countDocuments(),inventory:d.inventory.countDocuments()});printjson(d.orders.explain("executionStats").find({tenantId:"tenant-a",orderId:"O-000001"}).finish());d.restore_canary_ch25.insertOne({_id:"canary",createdAt:new Date()});printjson(d.restore_canary_ch25.findOne({_id:"canary"}));
Validation must prove all shards are reachable, sharded collections have the intended key/range metadata, targeted queries route as expected, counts/checksums/invariants match, and a temporary post-restore write succeeds. Only then restart the balancer and plan application cutover.
Check your understanding
- Why must the balancer be stopped for the documented self-managed dump workflow?
- Where is shard-key/range metadata stored?
- Why pre-shard destination collections before restoring?
-
Why exclude
config.*during the documented data restore? - When is a coordinated managed backup preferable?
Review the answers
1. A migrating range can otherwise appear on the wrong combination of backup components and create missing/duplicate logical state.
2. In the config server replica set, principally the config database used by mongos for routing.
3. The restore procedure reconstructs the destination sharding layout from captured metadata before loading documents.
4. To preserve the destination cluster’s own configuration after it has been intentionally created/pre-sharded.
5. When writes/transactions must continue and cross-shard transactional atomicity must be preserved without a quiescent self-managed backup window.
6. Production judgment
Sharded recovery is a cluster reconstruction problem, not a set of independent shard copies. Version the config metadata, shard keys, zones, balancer state, restore ordering, security configuration, and application invariants with the backup. If the business cannot tolerate the documented self-managed quiescent window, use a coordinated backup product designed for live sharded clusters. The next lesson combines these mechanisms into a timed restore drill with explicit RPO/RTO and cutover evidence.
Authoritative references
Backup behavior is topology-, tool-, storage-, and service-version sensitive. Re-check the current server, Database Tools, Atlas, encryption, and restore documentation before adopting a production procedure.
- Backup Methods for Self-Managed Deployments
- Back Up and Restore with MongoDB Tools
- mongodump 100.18.0
- mongorestore 100.18.0
- mongodump Examples
- mongorestore Behavior, Access, and Usage
- Database Tools Release Notes
- Filesystem Snapshots
- db.fsyncLock()
- Replication Oplog
- Backup and Restore Self-Managed Sharded Clusters
- Back Up Sharded Cluster with Database Dumps
- Restore Sharded Cluster from Database Dumps
- Config Servers
- Atlas Backup Architecture Guidance
- Atlas Continuous Cloud Backup Restore
- Atlas Backup Policy
- Atlas Disaster Recovery Guidance
- Encryption at Rest
- Queryable Encryption Key Management
- MongoDB 8.3 Release Notes
- mongosh Changelog
- PyMongo Release Notes