Assemble the architecture end-to-end and prove every hop of a targeted and broadcast request.
Build a Small Sharded Lab and Trace a Routed Query from mongos to Target Shards
Rebuild the compact topology, trace requests through mongos, verify physical placement, and document production gaps.
Learning objectives
Build the complete compact sharded topology from blank Docker resources.
Prove the application is connected to mongos and inspect the cluster registry before writing data.
Create two owned ranges and trace targeted versus scatter/gather queries with explain.
Observe per-shard document placement and distinguish diagnostic direct reads from supported application routing.
Validate cleanup, security, resource, and production-topology gaps before carrying the cluster model into shard-key design.
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 27131–27134. 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 capstone repeats the complete setup intentionally so the
lesson stands alone. Budget approximately 4 MongoDB processes
plus Docker overhead; if that is too large, run the Lesson 1
deterministic simulator and study the included expected state
transitions instead. 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. Build the topology from zero
docker rm -f atlasmart-ch16-l5-cfg atlasmart-ch16-l5-s1 atlasmart-ch16-l5-s2 atlasmart-ch16-l5-mongos 2>/dev/null || truedocker network rm atlasmart-ch16-l5-net 2>/dev/null || truedocker volume rm atlasmart-ch16-l5-cfg-data atlasmart-ch16-l5-s1-data atlasmart-ch16-l5-s2-data 2>/dev/null || truedocker network create atlasmart-ch16-l5-netdocker run -d --name atlasmart-ch16-l5-cfg --network atlasmart-ch16-l5-net -p 127.0.0.1:27131:27017 -v atlasmart-ch16-l5-cfg-data:/data/db mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim --configsvr --replSet atlasmart-cfg16-l5 --bind_ip_alldocker run -d --name atlasmart-ch16-l5-s1 --network atlasmart-ch16-l5-net -p 127.0.0.1:27132:27017 -v atlasmart-ch16-l5-s1-data:/data/db mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim --shardsvr --replSet atlasmart-shard16-l5-a --bind_ip_alldocker run -d --name atlasmart-ch16-l5-s2 --network atlasmart-ch16-l5-net -p 127.0.0.1:27133:27017 -v atlasmart-ch16-l5-s2-data:/data/db mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim --shardsvr --replSet atlasmart-shard16-l5-b --bind_ip_allfor PORT in 27131 27132 27133; 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:27131/admin?directConnection=true" --quiet --eval 'rs.initiate({_id:"atlasmart-cfg16-l5",configsvr:true,members:[{_id:0,host:"atlasmart-ch16-l5-cfg:27017"}]})'mongosh "mongodb://127.0.0.1:27132/admin?directConnection=true" --quiet --eval 'rs.initiate({_id:"atlasmart-shard16-l5-a",members:[{_id:0,host:"atlasmart-ch16-l5-s1:27017"}]})'mongosh "mongodb://127.0.0.1:27133/admin?directConnection=true" --quiet --eval 'rs.initiate({_id:"atlasmart-shard16-l5-b",members:[{_id:0,host:"atlasmart-ch16-l5-s2:27017"}]})'for PORT in 27131 27132 27133; 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-l5-mongos --network atlasmart-ch16-l5-net -p 127.0.0.1:27134:27017 --entrypoint mongos mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim --configdb atlasmart-cfg16-l5/atlasmart-ch16-l5-cfg:27017 --bind_ip_all --port 27017until mongosh "mongodb://127.0.0.1:27134/admin" --quiet --eval 'quit(db.runCommand({ping:1}).ok===1?0:1)'; do sleep 1; donemongosh "mongodb://127.0.0.1:27134/admin" --quiet --eval 'printjson(sh.addShard("atlasmart-shard16-l5-a/atlasmart-ch16-l5-s1:27017"));printjson(sh.addShard("atlasmart-shard16-l5-b/atlasmart-ch16-l5-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-l5-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-l5-b"));sh.status(true);
2. Verify router identity and metadata before tracing a request
printjson({hello:db.hello(),shards:db.adminCommand({listShards:1}).shards});sh.status(true);const cfg=db.getSiblingDB("config");const c=cfg.collections.findOne({_id:"atlasmart.orders"});printjson(cfg.chunks.find({uuid:c.uuid},{min:1,max:1,shard:1}).sort({min:1}).toArray());
The router should report msg:"isdbgrid". If it
does not, stop: the client is not connected to mongos.
Starting in MongoDB 8.3, sharded-cluster DDL belongs on
mongos, reinforcing the rule that the router is the
application/control entry point.
3. Trace one targeted request and one broadcast request
function shardNames(explain) { const found = new Set(); function walk(v) { if (!v || typeof v !== "object") return; if (typeof v.shardName === "string") found.add(v.shardName); if (Array.isArray(v.shards)) { for (const s of v.shards) { if (s && typeof s.shardName === "string") found.add(s.shardName); walk(s); } } for (const k of Object.keys(v)) walk(v[k]); } walk(explain); return Array.from(found).sort();}const app=db.getSiblingDB("atlasmart");const targeted=app.orders.find({tenantId:"b",status:"open"}).explain("executionStats");const scatter=app.orders.find({status:"open"}).explain("executionStats");printjson({kind:"targeted",shards:shardNames(targeted),returned:targeted.executionStats?.nReturned});printjson({kind:"scatter",shards:shardNames(scatter),returned:scatter.executionStats?.nReturned});
Conceptually, the request path is: client → mongos → metadata lookup/cache → target shard replica-set primary/eligible member → local query plan → mongos result merge → client. The broadcast case fans the middle of that path out to both shards.
4. Inspect placement directly, but do not turn diagnostics into architecture
echo 'Shard A diagnostic:'mongosh 'mongodb://127.0.0.1:27132/atlasmart?directConnection=true' --quiet --eval 'printjson(db.orders.find({}, {_id:0,tenantId:1,orderId:1}).sort({tenantId:1}).toArray())'echo 'Shard B diagnostic:'mongosh 'mongodb://127.0.0.1:27133/atlasmart?directConnection=true' --quiet --eval 'printjson(db.orders.find({}, {_id:0,tenantId:1,orderId:1}).sort({tenantId:1}).toArray())'echo 'Application query must still use mongos at 27134.'
Direct reads are useful in a disposable lab to prove physical placement. They are not the supported application interface because they ignore router metadata, ownership changes, and sharded-cluster command restrictions.
5. Optional PyMongo trace through mongos
from pymongo import MongoClientclient=MongoClient("mongodb://127.0.0.1:27134/",serverSelectionTimeoutMS=5000)hello=client.admin.command("hello")assert hello.get("msg") == "isdbgrid", helloorders=client.atlasmart.ordersprint("router:", client.address)print("tenant b:", list(orders.find({"tenantId":"b","status":"open"},{"_id":0}).sort("orderId",1)))print("all open count:", orders.count_documents({"status":"open"}))
6. Failure/misuse checklist before production
| Misuse | Concrete consequence | Safer response |
|---|---|---|
| One-member shard/CSRS in production | No replica redundancy; single-process failures remove availability. | Use properly sized multi-member replica sets across failure domains. |
| Hot path omits shard key | Scatter/gather fan-out and tail-latency dependence on all shards. | Redesign request/index/shard key; Chapter 17 formalizes the choice. |
| Balancer/migrations with no headroom | Cache/I/O/network/replication pressure and tail spikes. | Capacity plan, observe migrations, schedule where justified. |
| App connects directly to shard | Bypasses routing/metadata and can hit command restrictions or stale ownership. | Connect application and DDL through mongos. |
| Treat replication as backup | Logical deletes/corruption replicate too. | Independent tested backup/restore with RPO/RTO. |
A production sharded cluster is a distributed system of
replica sets plus routers and metadata control plane. Capacity
planning must include the slowest shard, migration traffic,
CSRS availability, router fleet, backup/restore,
authentication/TLS, tenant authorization, monitoring, and
upgrade/FCV coordination. The compact lab intentionally omits
those redundancies to keep learning free and local. The
maintenance-only directShardOperations role is
not an application escape hatch and should be used only for
narrowly defined maintenance under expert guidance.
Next chapter. Now that routing and range ownership are observable, Chapter 17 asks the harder question: which shard key produces good cardinality, distribution, targeting, write spread, and future evolvability?
docker rm -f atlasmart-ch16-l5-cfg atlasmart-ch16-l5-s1 atlasmart-ch16-l5-s2 atlasmart-ch16-l5-mongos 2>/dev/null || truedocker volume rm atlasmart-ch16-l5-cfg-data atlasmart-ch16-l5-s1-data atlasmart-ch16-l5-s2-data 2>/dev/null || truedocker network rm atlasmart-ch16-l5-net 2>/dev/null || true
Check your understanding
- What confirms that a client is connected to mongos?
- Why does the targeted query touch fewer shards?
- Why are direct shard reads used at all in this lesson?
- What does this four-process lab not demonstrate?
- What is the next design decision after architecture?
Review the answers
1. The hello response includes msg:"isdbgrid".
2. Its exact tenantId maps to one current shard-key range owner.
3. Only as isolated diagnostics to prove placement; applications should use mongos.
4. Production redundancy, independent failure domains, WAN behavior, or scale-out throughput.
5. Choosing a shard key that balances distribution, targeting, cardinality, and write behavior.
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.