Andrew Mercer
on this page

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.

  • Replication — each shard is itself a replica set
  • Security — the localhost exception behavior referenced above applies identically here