Decide whether AtlasMart has a distribution problem before adding a distributed database topology.
Why Shard: Data Capacity, Throughput, Working Set, and Geographic/Operational Needs
Compare scale pressures, simulate range routing, and build a tiny cluster without claiming one-host performance evidence.
Learning objectives
Explain when horizontal sharding addresses a capacity or throughput bottleneck and when it merely adds coordination cost.
Distinguish dataset size, working set, CPU/write throughput, and geographic/operational requirements as separate scaling pressures.
Use a deterministic router simulator to predict targeted versus broadcast work before building a cluster.
Inspect a real compact sharded topology without mistaking a one-host lab for production fault tolerance.
Identify cases where vertical scaling, indexing, caching, archiving, or workload redesign should precede sharding.
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 mandatory
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 router, exposed only on loopback diagnostic
ports 27115–27118. Single-member replica sets satisfy the
mechanism requirement but provide no production redundancy;
production deployments require properly sized multi-member
replica sets and independent failure domains. Authentication and
TLS are disabled only for this isolated local lab. Feature
Compatibility Version (FCV) is observed and never changed.
Reads/writes go through mongos; direct shard
connections are used only for explicitly labeled diagnostics.
This first lesson includes both a deterministic routing
simulator and the compact cluster; the simulator is the
mandatory fallback if the learner cannot comfortably run four
MongoDB containers. 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; expected
invariants are documentation-derived and measured timing/output
must be recorded on the learner machine.
1. Sharding solves distribution problems, not every slow-query problem
AtlasMart has outgrown one database host only after several independent pressures line up: order history is larger than the desired working set, write traffic saturates the host during promotions, and regional teams need operational isolation. Sharding partitions one collection across shards; each shard owns only part of the shard-key value space. The client still sees one logical collection through a mongos query router.
| Pressure | Evidence before sharding | What sharding can change | What it does not automatically fix |
|---|---|---|---|
| Data capacity | Data/indexes exceed practical single-host storage or backup window. | Distributes collection data and index working sets. | Bad schema, oversized documents, or retention policy. |
| CPU / write throughput | Sustained saturation remains after query/index tuning. | Routes different shard-key ranges to different shards. | A single hot shard-key value or globally serialized invariant. |
| Working set | Useful index/data pages exceed RAM and cause I/O pressure. | Each shard may cache a smaller owned subset. | Queries that broadcast to every shard still touch many caches. |
| Geographic/operational | Placement/compliance/ownership requirements are explicit. | Later chapters can map ranges to zones. | WAN latency, regional failover, or compliance policy by itself. |
“The query is slow, therefore shard” is incomplete. First
measure query shape, indexes, data growth, cache pressure,
write saturation, and failure-domain needs. A sharded cluster
can make an un-targetable hot query more expensive because
mongos must scatter it to multiple shards and
merge the results.
2. Model routing before paying topology cost
A range is an inclusive-lower/exclusive-upper interval of shard-key values owned by one shard. This simulator uses the same mental model as ranged sharding: exact tenant keys target one owner, while a predicate without the shard key broadcasts to every shard. It is intentionally not a performance simulator.
ranges = [ (None, "m", "shardA"), ("m", None, "shardB"),]def owns(value, lo, hi): return (lo is None or value >= lo) and (hi is None or value < hi)def route(filter_doc): if "tenantId" not in filter_doc: return sorted({r[2] for r in ranges}) v = filter_doc["tenantId"] return [name for lo, hi, name in ranges if owns(v, lo, hi)]for q in [ {"tenantId": "b", "status": "open"}, {"tenantId": "p", "status": "open"}, {"status": "open"},]: print(q, "=>", route(q))# Expected owner sets: shardA, shardB, then both shards.
The simulator proves only routing logic for the declared ranges. It says nothing about network latency, shard load, storage engines, replica lag, or merge cost.
3. Build a compact observation cluster
docker rm -f atlasmart-ch16-l1-cfg atlasmart-ch16-l1-s1 atlasmart-ch16-l1-s2 atlasmart-ch16-l1-mongos 2>/dev/null || truedocker network rm atlasmart-ch16-l1-net 2>/dev/null || truedocker volume rm atlasmart-ch16-l1-cfg-data atlasmart-ch16-l1-s1-data atlasmart-ch16-l1-s2-data 2>/dev/null || truedocker network create atlasmart-ch16-l1-netdocker run -d --name atlasmart-ch16-l1-cfg --network atlasmart-ch16-l1-net -p 127.0.0.1:27115:27017 -v atlasmart-ch16-l1-cfg-data:/data/db mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim --configsvr --replSet atlasmart-cfg16-l1 --bind_ip_alldocker run -d --name atlasmart-ch16-l1-s1 --network atlasmart-ch16-l1-net -p 127.0.0.1:27116:27017 -v atlasmart-ch16-l1-s1-data:/data/db mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim --shardsvr --replSet atlasmart-shard16-l1-a --bind_ip_alldocker run -d --name atlasmart-ch16-l1-s2 --network atlasmart-ch16-l1-net -p 127.0.0.1:27117:27017 -v atlasmart-ch16-l1-s2-data:/data/db mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim --shardsvr --replSet atlasmart-shard16-l1-b --bind_ip_allfor PORT in 27115 27116 27117; 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:27115/admin?directConnection=true" --quiet --eval 'rs.initiate({_id:"atlasmart-cfg16-l1",configsvr:true,members:[{_id:0,host:"atlasmart-ch16-l1-cfg:27017"}]})'mongosh "mongodb://127.0.0.1:27116/admin?directConnection=true" --quiet --eval 'rs.initiate({_id:"atlasmart-shard16-l1-a",members:[{_id:0,host:"atlasmart-ch16-l1-s1:27017"}]})'mongosh "mongodb://127.0.0.1:27117/admin?directConnection=true" --quiet --eval 'rs.initiate({_id:"atlasmart-shard16-l1-b",members:[{_id:0,host:"atlasmart-ch16-l1-s2:27017"}]})'for PORT in 27115 27116 27117; 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-ch16-l1-mongos --network atlasmart-ch16-l1-net -p 127.0.0.1:27118:27017 --entrypoint mongos mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim --configdb atlasmart-cfg16-l1/atlasmart-ch16-l1-cfg:27017 --bind_ip_all --port 27017until mongosh "mongodb://127.0.0.1:27118/admin" --quiet --eval 'quit(db.runCommand({ping:1}).ok===1?0:1)'; do sleep 1; donemongosh "mongodb://127.0.0.1:27118/admin" --quiet --eval 'printjson(sh.addShard("atlasmart-shard16-l1-a/atlasmart-ch16-l1-s1:27017"));printjson(sh.addShard("atlasmart-shard16-l1-b/atlasmart-ch16-l1-s2:27017"));printjson({hello:db.hello(),shards:db.adminCommand({listShards:1}).shards});'
const admin = db.getSiblingDB("admin");const app = db.getSiblingDB("atlasmart");printjson(sh.enableSharding("atlasmart", "atlasmart-shard16-l1-a"));app.orders.drop();app.orders.createIndex({tenantId:1,orderId:1});printjson(sh.shardCollection("atlasmart.orders", {tenantId:1}));const docs=[];for (const t of ["a","b","c","n","o","p"]) { for (let i=1;i<=4;i++) docs.push({tenantId:t,orderId:`${t}-${i}`,status:i%2?"open":"closed",amount:i*25});}app.orders.insertMany(docs);printjson(sh.splitAt("atlasmart.orders", {tenantId:"m"}));printjson(sh.moveChunk("atlasmart.orders", {tenantId:"z"}, "atlasmart-shard16-l1-b"));sh.status(true);
Run all DDL and application operations through the router at
27118. The hello response from
mongos contains msg: "isdbgrid", which
is a useful sanity check that you did not accidentally connect
to a shard.
printjson(db.hello());sh.status(true);printjson(sh.getShardedDataDistribution());const cfg = db.getSiblingDB("config");const coll = cfg.collections.findOne({_id:"atlasmart.orders"});printjson(coll);printjson(cfg.chunks.find({uuid:coll.uuid},{min:1,max:1,shard:1,history:1}).sort({min:1}).toArray());
The config database is internal cluster metadata.
Reading selected collections is useful for
education/diagnosis, but applications must not write or depend
on its internal schema as an application API.
4. Verify the decision rather than celebrating the topology
On this tiny single-host cluster, sharding adds processes, metadata, routing, and migration behavior but does not demonstrate production throughput gains. A valid production decision needs before/after measurements under the same workload: per-shard CPU and I/O, working-set fit, query targeting ratio, p50/p95/p99 latency, data skew, balancer activity, replication headroom, backup/restore objectives, and operational staffing.
Shard when a measured capacity, throughput, placement, or operational constraint cannot be solved more simply. Sharding does not strengthen write durability, replace replica sets or backups, fix a bad shard key, or make broadcast queries cheap. It adds metadata availability requirements, migration headroom, router health, per-shard indexes, and a larger security surface. Plan a rollback/migration path before data distribution becomes difficult to reverse.
Bridge. Lesson 2 opens the black box and
assigns responsibility to the shards, config server replica set,
and mongos.
docker rm -f atlasmart-ch16-l1-cfg atlasmart-ch16-l1-s1 atlasmart-ch16-l1-s2 atlasmart-ch16-l1-mongos 2>/dev/null || truedocker volume rm atlasmart-ch16-l1-cfg-data atlasmart-ch16-l1-s1-data atlasmart-ch16-l1-s2-data 2>/dev/null || truedocker network rm atlasmart-ch16-l1-net 2>/dev/null || true
Check your understanding
- What is the difference between a shard and mongos?
- Does a large collection automatically justify sharding?
- Why can sharding worsen a bad query?
- What does the one-host lab prove?
Review the answers
1. A shard is a replica set that owns data; mongos is the stateless query router that uses metadata to route client operations.
2. No. Measure storage/working-set, throughput, query/index, retention, and operational constraints first.
3. If mongos cannot target by shard key, it may broadcast to all shards and merge their work.
4. Routing and metadata mechanisms only, not production availability or scale-out performance.
Authoritative references
- MongoDB Sharding — Sharded-cluster purpose, chunks/ranges, targeted operations, and architecture.
- Routing with mongos — Router metadata cache, targeted versus broadcast operations, and aggregation routing evidence.
- Config Servers — Config server replica sets, metadata responsibilities, and config-shard alternatives.
-
config Database
— Internal metadata including
config.collections,config.chunks,config.shards, and settings. - sh.status() — Shards, databases, ranges, sharded data distribution, and migration summaries.
- sh.shardCollection() — Shard a collection and define its shard key.
-
sh.addShard()
— Add replica-set shards through
mongos. - Split Chunks/Ranges — Controlled manual split examples and operational cautions.
-
moveRange
— Explicit range migration through
mongos. - Sharded Cluster Balancer — Range migration procedure, thresholds, cleanup, and resource impact.
- MongoDB 8.3 Release Notes — Current 8.3 behavior including mongos-only DDL on sharded clusters.
- mongosh Release Notes — mongosh 2.10.0 baseline.
- PyMongo Release Notes — PyMongo 4.17 baseline.