skip to content

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