Turn the shard key from a schema field into a routing instruction whose absence can multiply work.
Query Targeting vs Scatter/Gather: How Shard Keys Determine Routing Cost
Create deterministic ranges and use explain to show exactly which shards received each AtlasMart query.
Learning objectives
Predict whether a query is targeted or broadcast from its shard-key predicate.
Use mongos explain output to identify which shards actually received a query.
Explain why a correct result can still be operationally expensive when it scatters.
Connect compound shard-key prefixes to targeting without pre-teaching Chapter 17 shard-key selection.
Quantify fan-out and merge work without using planner labels as a latency guarantee.
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 27123–27126. 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.
The lab creates two deterministic ranges split at
tenantId:"m" so exact tenant predicates have an
obvious owner. This is a routing demonstration, not a shard-key
recommendation for real production. 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. mongos routes from predicates plus metadata
A targeted operation is sent only to the shard or subset of shards whose ranges can satisfy the shard-key predicate. A broadcast or scatter/gather operation is sent to all relevant shards because the router cannot rule them out. The final result can be identical while the operational cost is radically different.
| Query shape | Routing expectation | Why |
|---|---|---|
{tenantId:"b"} |
One shard | Exact shard-key value maps into one owned range. |
{tenantId:{$gte:"a",$lt:"m"}} |
Subset/one shard in this fixture | Range overlaps only the lower range. |
{status:"open"} |
Both shards | Predicate has no shard-key information. |
{tenantId:{$in:["b","p"]}} |
Both shards | Values live in different ranges. |
docker rm -f atlasmart-ch16-l3-cfg atlasmart-ch16-l3-s1 atlasmart-ch16-l3-s2 atlasmart-ch16-l3-mongos 2>/dev/null || truedocker network rm atlasmart-ch16-l3-net 2>/dev/null || truedocker volume rm atlasmart-ch16-l3-cfg-data atlasmart-ch16-l3-s1-data atlasmart-ch16-l3-s2-data 2>/dev/null || truedocker network create atlasmart-ch16-l3-netdocker run -d --name atlasmart-ch16-l3-cfg --network atlasmart-ch16-l3-net -p 127.0.0.1:27123:27017 -v atlasmart-ch16-l3-cfg-data:/data/db mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim --configsvr --replSet atlasmart-cfg16-l3 --bind_ip_alldocker run -d --name atlasmart-ch16-l3-s1 --network atlasmart-ch16-l3-net -p 127.0.0.1:27124:27017 -v atlasmart-ch16-l3-s1-data:/data/db mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim --shardsvr --replSet atlasmart-shard16-l3-a --bind_ip_alldocker run -d --name atlasmart-ch16-l3-s2 --network atlasmart-ch16-l3-net -p 127.0.0.1:27125:27017 -v atlasmart-ch16-l3-s2-data:/data/db mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim --shardsvr --replSet atlasmart-shard16-l3-b --bind_ip_allfor PORT in 27123 27124 27125; 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:27123/admin?directConnection=true" --quiet --eval 'rs.initiate({_id:"atlasmart-cfg16-l3",configsvr:true,members:[{_id:0,host:"atlasmart-ch16-l3-cfg:27017"}]})'mongosh "mongodb://127.0.0.1:27124/admin?directConnection=true" --quiet --eval 'rs.initiate({_id:"atlasmart-shard16-l3-a",members:[{_id:0,host:"atlasmart-ch16-l3-s1:27017"}]})'mongosh "mongodb://127.0.0.1:27125/admin?directConnection=true" --quiet --eval 'rs.initiate({_id:"atlasmart-shard16-l3-b",members:[{_id:0,host:"atlasmart-ch16-l3-s2:27017"}]})'for PORT in 27123 27124 27125; 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-l3-mongos --network atlasmart-ch16-l3-net -p 127.0.0.1:27126:27017 --entrypoint mongos mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim --configdb atlasmart-cfg16-l3/atlasmart-ch16-l3-cfg:27017 --bind_ip_all --port 27017until mongosh "mongodb://127.0.0.1:27126/admin" --quiet --eval 'quit(db.runCommand({ping:1}).ok===1?0:1)'; do sleep 1; donemongosh "mongodb://127.0.0.1:27126/admin" --quiet --eval 'printjson(sh.addShard("atlasmart-shard16-l3-a/atlasmart-ch16-l3-s1:27017"));printjson(sh.addShard("atlasmart-shard16-l3-b/atlasmart-ch16-l3-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-l3-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-l3-b"));sh.status(true);
2. Ask explain which shards were actually touched
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");for (const [name,query] of [ ["target-a",{tenantId:"b",status:"open"}], ["target-b",{tenantId:"p",status:"open"}], ["scatter",{status:"open"}], ["multi-range",{tenantId:{$in:["b","p"]},status:"open"}],]) { const ex=app.orders.find(query).explain("executionStats"); printjson({name,query,shards:shardNames(ex),nReturned:ex.executionStats?.nReturned});}
Do not hard-code the full explain tree into application logic;
execution engines and wrapper stages can evolve. For diagnosis,
inspect the shards work reported by mongos and
compare it with the expected range owners.
3. Correct but expensive: measure scatter fan-out
const app=db.getSiblingDB("atlasmart");const targeted=app.orders.find({tenantId:"b",status:"open"}).sort({orderId:1}).toArray();const scatter=app.orders.find({status:"open"}).sort({tenantId:1,orderId:1}).toArray();printjson({targetedCount:targeted.length,scatterCount:scatter.length,targeted});printjson(scatter.slice(0,6));
A broadcast query is not automatically “wrong.” Analytics and administrative workloads sometimes need all shards. The mistake is putting an unavoidable scatter/gather shape on a latency-sensitive hot path without budgeting per-shard work, merge cost, tail latency, and failure behavior.
4. Prefix targeting for compound shard keys
MongoDB can target operations that contain a shard key or a
prefix of a compound shard key. For a future shard key
{tenantId:1, orderId:1},
{tenantId:"b"} can still narrow routing. A
predicate on orderId alone cannot use that prefix.
Chapter 17 will decide whether such a key is actually a good
distribution choice.
Do not add shard-key fields to every query mechanically. A predicate must be semantically correct for the business request, and targeting alone does not guarantee a fast local plan. Each targeted shard still needs appropriate indexes and enough capacity.
5. Production judgment
Track routing fan-out as a first-class service metric. A query that targets one shard today can target more after range movement or with broader predicates. Measure per-shard execution, router merge behavior, network bytes, and p95/p99 latency. Security also matters: the shard key often carries tenant or business identity, but its presence in a document is not authorization.
Bridge. Lesson 4 follows those ranges while they move between shards and shows why redistribution needs operational headroom.
docker rm -f atlasmart-ch16-l3-cfg atlasmart-ch16-l3-s1 atlasmart-ch16-l3-s2 atlasmart-ch16-l3-mongos 2>/dev/null || truedocker volume rm atlasmart-ch16-l3-cfg-data atlasmart-ch16-l3-s1-data atlasmart-ch16-l3-s2-data 2>/dev/null || truedocker network rm atlasmart-ch16-l3-net 2>/dev/null || true
Check your understanding
- What makes a query targeted?
- Can a correct query still be operationally poor?
- Does targeted routing guarantee an index scan on the shard?
- Can a compound shard-key prefix help targeting?
Review the answers
1. The mongos can map its shard-key predicate to only the ranges/shards that can satisfy it.
2. Yes. A scatter/gather can return correct results while multiplying shard, network, and merge work.
3. No. Each shard performs normal local planning and may still choose a collection scan.
4. Yes; mongos can often target from the full shard key or an appropriate leading prefix.
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.