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.

Intermediate120–190 minutesSharded-cluster routing/operations labMongoDB 8.3.8 · mongosh 2.10.0 · PyMongo 4.17.0Last reviewed: September 2026

Learning objectives

01

Predict whether a query is targeted or broadcast from its shard-key predicate.

02

Use mongos explain output to identify which shards actually received a query.

03

Explain why a correct result can still be operationally expensive when it scatters.

04

Connect compound shard-key prefixes to targeting without pre-teaching Chapter 17 shard-key selection.

05

Quantify fan-out and merge work without using planner labels as a latency guarantee.

Reproducible lab baseline

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.
start compact sharded topology (l3)
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});' 
create deterministic AtlasMart ranges through mongos
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

extract shard names from mongos explain
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

compare result equality with routing 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));
Evidence interpretation

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.

Wrong approach

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.

cleanup / full reset
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

  1. What makes a query targeted?
  2. Can a correct query still be operationally poor?
  3. Does targeted routing guarantee an index scan on the shard?
  4. 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

Keep knowledge open

Help the academy stay free and grow.

If these tutorials save you time, a small donation supports new lessons, technical review, diagrams, examples, and long-term maintenance.

ETHEthereum / ERC-20 only
0x716c4Ab160C4B66F31a28AE2448BfF68fc3a2ef0

Send only Ethereum or ERC-20 compatible assets to this address.