MongoDB Sharding¶
Sharding splits data across multiple replica sets (shards) rather than copying the whole data set to every node — the MongoDB equivalent of what Citus does for PostgreSQL. A mongos router sits in front and directs each query to the shard(s) that actually hold the relevant data, coordinated by a set of config servers holding cluster metadata.
Building a Test Sharded Cluster with mtools¶
mtools' mlaunch is a convenience wrapper for standing up a full sharded + replicated topology locally, useful for testing before building the real thing by hand.
sudo pip install mtools
mlaunch init --sharded 3 --replicaset --nodes 3 --config 3 --auth
This launches 3 shards (each a 3-node replica set), 3 config servers, and a mongos router — 13 mongod/mongos processes in total — and wires them together automatically. With --auth, it also creates an initial admin user and prints the generated credentials to the terminal; capture those immediately, since they aren't shown again.
mongo admin -u <printed_user> -p <printed_password>
mongos> show dbs
admin 0.000GB
config 0.000GB
The Localhost Exception on Individual Shards¶
Cluster-wide user accounts live on the config servers and are accessed through mongos. But each shard is still an independent mongod/replica set underneath, and — same as any standalone instance — allows the localhost exception for its own first user, separate from the cluster's admin user.
mongo --port <shard_member_port>
use admin
db.createUser({ user: "shardLocalAdmin", pwd: "a_strong_password", roles: ["root"] })
db.auth("shardLocalAdmin", "a_strong_password")
show dbs
Useful for direct shard-level maintenance (e.g. inspecting a shard's own oplog) without routing through mongos — but remember this account only has authority on that one shard, not the cluster as a whole.
Restoring a Sharded Cluster From Backup¶
A simplified walkthrough for standing up a temporary cluster from a config-server backup and a shard backup — useful for testing a restore procedure or recovering an isolated environment.
1. Start a Temporary Config Server¶
mongod --configsvr --dbpath <path_to_config_backup> --port 26050 --fork --logpath <log_path> --logappend
2. Restore the Config Server Data¶
mongorestore -h localhost:26050 --db config --dir <path_to_config_backup>
3. Verify and Fix Shard Host References¶
The config server's shards collection stores each shard's replica-set connection string, pointing at the original hostnames — which won't resolve in a temporary/isolated restore environment.
mongo localhost:26050/config
db.shards.find()
{ "_id" : "s1", "host" : "s1/shard-host-1:27501,shard-host-2:27502,shard-host-2:27503" }
{ "_id" : "s2", "host" : "s2/shard-host-4:27601,shard-host-5:27602,shard-host-5:27603" }
Repoint each shard entry at wherever you're actually restoring it in the temporary environment:
db.shards.replaceOne(
{ "_id": "s1", "host": "s1/shard-host-1:27501,shard-host-2:27502,shard-host-2:27503" },
{ "_id": "s1", "host": "localhost:27501" }
)
db.shards.replaceOne(
{ "_id": "s2", "host": "s2/shard-host-4:27601,shard-host-5:27602,shard-host-5:27603" },
{ "_id": "s2", "host": "localhost:27601" }
)
4. Restart the Config Server¶
kill <config_server_pid>
mongod --configsvr --dbpath <path_to_config_backup> --port 26050 --fork --logpath <log_path> --logappend
5. Start the Shard Servers¶
mongod --shardsvr --dbpath /tmp/restore/s1 --logpath /tmp/restore/shardsvr1.log --port 27501 --fork --logappend
mongod --shardsvr --dbpath /tmp/restore/s2 --logpath /tmp/restore/shardsvr2.log --port 27601 --fork --logappend
6. Restore Each Shard's Data¶
mongorestore -h localhost:27501 --dir /tmp/backup/s1
mongorestore -h localhost:27601 --dir /tmp/backup/s2
7. Start a mongos Router¶
Point it at the (now-repointed) config server, and the cluster should come up recognizing both restored shards.
This restore path is genuinely fiddly — the host-repointing step in particular is easy to get subtly wrong (matching on the exact original string). For anything beyond a one-off test restore, prefer mongorestore against a cluster where the shards were provisioned with the hostnames the backup already expects, avoiding the manual shards collection surgery entirely.
Related¶
- Replication — each shard is itself a replica set
- Security — the localhost exception behavior referenced above applies identically here