skip to content

Query Routing: Targeted vs Scatter-Gather

Whether a query hits one shard or all of them depends entirely on whether the shard key is in the filter. This is the payoff question for shard-key design, and the reason a sharded cluster can be slower than a single node.

part ofMongoDBoverview, primer and where to startread it →
on this pageshow

questions

6

Which query filters let mongos target one shard, and which force a broadcast to every shard?

level: middleimportance: must knowfreq 72%

answer

  1. Routing is read off the filter, nothing else
  2. A compound key follows the same prefix idea as an index
  3. Hashing scrambles the order of neighbouring values
  4. Being unique per shard is not the same as being locatable
  5. Between one shard and all shards there is a middle case

basics

~20 s

A filter containing the shard key — or a leading prefix of a compound shard key — lets mongos map values to chunks and contact only the owning shards. Any other filter, including one on _id when _id is not the shard key, is broadcast to all shards.

solid answer

~50 s

Targeting is decided from the filter alone. If the filter has an equality match on the full shard key, `mongos` resolves it to one chunk and therefore one shard. With a **ranged** shard key, a range predicate on the key resolves to the chunks covering that range — a few shards, not all. With a **hashed** shard key, only equality is targetable: neighbouring key values hash to unrelated chunks, so a range predicate on a hashed key broadcasts. For a compound shard key such as `{orgId: 1, createdAt: 1}`, a filter on the **leading prefix** (`orgId`) is targetable; a filter on `createdAt` alone is not. Everything else — no shard key, a filter on `_id` when `_id` is not the shard key, a `$or` whose branches are not all targetable — fans out to every shard owning data for the collection.

code

javascript · 13 lines
javascript
// Shard key: { orgId: 1, createdAt: 1 }

// targeted - leading prefix present
db.orders.find({ orgId: "acme" })

// targeted and narrower - prefix plus range on second field
db.orders.find({ orgId: "acme", createdAt: { $gt: ISODate("2026-01-01") } })

// broadcast - leading field missing
db.orders.find({ createdAt: { $gt: ISODate("2026-01-01") } })

// broadcast - _id is not the shard key here
db.orders.find({ _id: ObjectId("6600000000000000000000aa") })

go deeper

for a junior

Remember the headline: the shard key in the filter means one shard, no shard key means all shards. Be able to point at a filter and say which it is.

for a middle

Explain the prefix rule for compound shard keys, why a hashed key targets equality but not ranges, and why an index on the filtered field changes nothing about routing.

for a senior

Talk about shards-touched-per-query as the metric, verify it with explain rather than by reasoning, and be able to walk a real workload's query shapes one by one.

for a principal

Own the consequence: broadcast-heavy traffic means the cluster's throughput for that traffic does not grow with shard count, which is the argument you take to a capacity or re-architecture decision.

## The rule in one line `mongos` can restrict a query to a subset of shards **only** when it can turn the filter into a set of shard-key ranges. Everything about targeting follows from that. ## Equality on the full shard key This is the ideal case. The router hashes or compares the value, finds the single chunk whose range contains it, and sends the query to that chunk's owner. One shard does the work; the router does no merging worth mentioning. This is what "the shard key is in the query" means when people say a workload shards well. ```js // shard key: { orgId: 1 } db.orders.find({ orgId: "acme", status: "OPEN" }) // targeted ``` The extra `status` predicate is irrelevant to routing; it is evaluated on the shard. ## Ranges, and why hashed keys behave differently With a **ranged** shard key, chunk boundaries are ordered by the key's value, so a predicate like `{createdAt: {$gte: X, $lt: Y}}` on a `createdAt` shard key maps to a contiguous run of chunks. That may be one shard or several, but it is usually far fewer than all of them. With a **hashed** shard key, chunks are ranges of the *hash*, not of the value. Two adjacent values hash to unrelated positions, so no contiguous range of values corresponds to a contiguous range of chunks. The consequence is sharp: **equality on a hashed shard key is targetable; a range on it is not** and produces a full broadcast. Candidates routinely miss this and describe hashed keys as universally better because they spread writes evenly — they do, at the cost of range targeting. ## Compound shard keys and the prefix rule A compound shard key behaves like a compound index for targeting purposes: the router can use a **leading prefix** of it. With `{orgId: 1, createdAt: 1}`: - `{orgId: "acme"}` — targetable, to the shards owning that `orgId` range. - `{orgId: "acme", createdAt: {$gt: T}}` — targetable, and narrower. - `{createdAt: {$gt: T}}` — **not** targetable; the leading field is missing, so the router broadcasts. This mirrors the equality-first intuition from compound indexes, which is why the two are so often confused in interviews. Keep them separate: the index decides how efficiently one shard finds documents; the shard key decides how many shards are asked at all. ## What always broadcasts - A filter with no shard-key field in it. The common case: querying by `email`, by `status`, by a date when the key is `userId`. - A filter on `_id` when `_id` is not the shard key. `_id` is unique per shard, but the router has no idea which shard owns a given `_id`, so it asks all of them. This surprises people constantly. - Most `$or` filters where at least one branch lacks the shard key — the union of possible targets is then every shard. - Negation-shaped predicates on the shard key such as `$ne` or `$nin`, which can match values anywhere in the key space. - An empty filter, a `count` over the whole collection, or a text-style search on a non-key field. ## Targeting is not all-or-nothing Between one shard and all shards there is a middle ground: a range predicate on a ranged shard key may hit three shards out of twenty. That is still an enormous win over a broadcast. When evaluating a workload, think in terms of *shards touched per query*, not a binary. ## How to check rather than guess Run `explain()` through `mongos`. For a sharded collection, the winning plan carries a top-level stage of `SINGLE_SHARD` when one shard was contacted and `SHARD_MERGE` when several were, plus a `shards` array naming them. That output is the ground truth for whether your reasoning about the filter is right; index hints and clever rewrites do not change routing, only the filter's shape does. ## Why this is the payoff question of sharding A broadcast makes every shard do work for every request. The cluster's total capacity for such queries does not grow when you add shards — each new shard just adds another participant to every request, plus one more stream to merge and one more chance to be the slow one. Targeted queries are the opposite: each shard sees roughly 1/N of the traffic, so capacity scales. Sharding pays off exactly to the extent that the hot queries carry the shard key.

  • Why does a filter on _id broadcast when _id is not the shard key?
    Because routing needs a shard-key value. `_id` is unique within the collection, but the router has no mapping from an `_id` to a chunk, so it cannot exclude any shard and must ask all of them. Each shard then does an efficient `_id` index lookup, but you paid a cluster-wide fan-out to find one document.
  • Does adding an index on the filtered field make the query targeted?
    No. Indexes affect how fast a shard finds matching documents once the query reaches it; they have no bearing on how many shards the router contacts. A broadcast query with perfect indexes is still a broadcast — it is just a fast one on every shard.
  • How would you confirm targeting rather than reason about it?
    Run `explain()` through mongos and look at the winning plan's top stage: `SINGLE_SHARD` means one shard was contacted, `SHARD_MERGE` means several, and the accompanying `shards` array names them. Do this for each hot query shape, since routing is per-query, not per-cluster.
  • Is a query that touches three of twenty shards a failure?
    Not at all. Targeting is a spectrum. A range predicate on a ranged shard key often resolves to a handful of adjacent chunks, so the query touches a few shards instead of all twenty. That still scales: adding shards keeps the per-query participant count roughly constant while total capacity grows.

saying these in an interview costs you the question

  • Says any indexed field makes the query targeted
  • Thinks a range on a hashed shard key targets shards
  • Claims a filter on _id always goes to one shard
  • Believes the trailing field of a compound key can target alone
  • Treats routing as all shards or exactly one shard

context

open as a page

After sharding, p99 read latency got worse and every shard is busy — how do you confirm broadcast queries are the cause?

level: seniorimportance: must knowfreq 58%

basics

~20 s

Run explain() through mongos on the hot query shapes: a SHARD_MERGE stage with every shard named in the winning plan means scatter-gather. Compare per-shard keysExamined against nReturned, then fix the filter to carry the shard key rather than adding shards.

open as a page

In a MongoDB sharded cluster, what does mongos do with a find() the driver sends it?

level: juniorimportance: should knowfreq 45%

basics

~20 s

mongos is the query router. It uses the chunk map cached from the config servers to pick which shards can hold matches — one when the filter names the shard key, all shards otherwise — and merges their replies into a single cursor.

open as a page

How does mongos handle sort(), limit() and skip() when a find() fans out to several shards?

level: middleimportance: should knowfreq 50%

basics

~20 s

Each shard sorts and limits its own results; mongos merge-sorts the pre-sorted streams and re-applies the limit. It cannot push skip down — it fetches unskipped results and skips while assembling, passing skip plus limit to the shards when both are present.

open as a page

What must an updateOne or deleteOne on a sharded collection include in its filter, and what changed in MongoDB 7.0?

level: seniorimportance: should knowfreq 38%

basics

~20 s

Before MongoDB 7.0, a single-document update, delete or findAndModify on a sharded collection had to include the shard key or _id in its filter, or it errored. MongoDB 7.0 lifted that: the cluster now locates the document itself, at the cost of a broadcast. Upserts still need the full shard key.

open as a page

How is an aggregation pipeline split between the shards and a merging node in a sharded cluster?

level: seniorimportance: nice to knowfreq 33%

basics

~20 s

MongoDB splits the pipeline into a shards part that runs on every participating shard and a merger part that runs on one node. explain() shows splitPipeline with shardsPart, mergerPart and mergeType; stages needing an unsharded collection, such as $out, force the merge onto the database's primary shard.

open as a page