skip to content

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

level: middleimportance: should knowfreq 50%

answer

  1. Each shard returns its portion already ordered
  2. The router does a k-way merge, not a full sort
  3. One of the three modifiers cannot be pushed down
  4. A shard cannot know what precedes it globally
  5. Together, the two paging modifiers are sent as their sum

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.

solid answer

~50 s

For a **sort**, `mongos` sends the sort specification to every participating shard. Each shard returns its own results already in order (ideally straight from an index), and the router performs a **merge-sort** over those streams, so it never has to hold the whole result set sorted in memory. For a **limit**, the router pushes the limit down: each shard returns at most that many documents, and `mongos` re-applies the limit to the merged stream. Ten shards with `limit(20)` means up to 200 documents crossing the network to produce 20. A **skip** cannot be pushed down — shard 3 has no idea how many of the globally-preceding documents live on shards 1 and 2 — so the router retrieves unskipped results and discards the first *n* while assembling. When skip and limit appear together, `mongos` asks each shard for `skip + limit` documents, which bounds the transfer but still grows with the offset.

code

javascript · 7 lines
javascript
// Sharded on { orgId: 1 }; this filter has no shard key, so it fans out.
db.events.find({ level: "ERROR" })
         .sort({ ts: -1 })
         .skip(10000)
         .limit(20)
// Each shard is asked for up to 10020 documents in ts order;
// mongos merge-sorts the streams, drops 10000, returns 20.

go deeper

for a junior

Know that with several shards involved the sort happens on the shards and mongos combines the ordered streams, and that a limit does not mean only that many documents move.

for a middle

Explain the three treatments precisely — merge-sort for sort, push-down-then-reapply for limit, router-side application for skip — and why skip is the one that cannot be pushed down.

for a senior

Reason about the cost model out loud: per-shard fetch of skip plus limit, N times the transfer, latency set by the slowest shard, and the fixes you would actually apply to a deep-paging endpoint.

for a principal

Own the design consequence: paging APIs over sharded data should be addressed by position in the key space rather than by numeric offset, and that is an API contract decision, not a query tweak.

## The problem the router faces `sort`, `skip` and `limit` are defined over the **global** result set, but the data lives on N shards that know nothing about each other. The router has to reconstruct a global answer from N local ones without buying a full global sort. The three modifiers get three different treatments. ## Sort: merge-sort of pre-sorted streams The sort specification goes to every participating shard. Each shard produces its portion **already in sort order** — from an index if one supports the sort, otherwise via an in-memory sort on the shard. `mongos` then does a k-way merge: it peeks at the head of each shard's stream, emits the smallest, and pulls the next document from that stream. Two consequences follow. First, the router's memory cost is bounded by the current batch from each shard, not by the size of the result — the merge is streaming. Second, an index that supports the sort matters *more* under sharding than on a single node, because without it every shard independently pays a blocking in-memory sort before the merge can even start, and the query is only as fast as the slowest of those sorts. ## Limit: pushed down, then re-applied A limit is safe to push down, because the global top-K is always a subset of the union of the per-shard top-Ks. `mongos` sends the limit to each shard and re-applies it after merging. The waste is bounded and predictable: with N shards and `limit(L)`, up to `N × L` documents are produced and shipped so that `L` survive. For small L that is fine; for a large limit across many shards it is real network and serialization cost. ## Skip: cannot be pushed down A skip is *not* safe to push down. "Skip the first 100 globally" cannot be turned into "skip the first 100 on each shard" — a shard has no idea how many documents that globally precede its own live elsewhere, and skipping locally would drop documents that belong in the result. So `mongos` retrieves unskipped results from the shards and applies the skip itself while assembling the merged stream. ## Skip and limit together When both appear, the router does have one optimisation available: no document beyond position `skip + limit` in *any* shard's local order can appear in the final answer. So `mongos` passes `skip + limit` as the limit to each shard, and then applies the skip and the limit itself on the merged result. That bounds the transfer, but the bound grows with the offset: `find().sort(...).skip(10000).limit(20)` across 10 shards asks each shard for 10,020 documents so that 20 are returned. ## Reading the cost The practical model for a fanned-out, sorted, paged query is: - **Per-shard work**: index scan (or worse, collection scan plus in-memory sort) for up to `skip + limit` documents. - **Network**: up to `N × (skip + limit)` documents to the router. - **Router work**: a streaming merge, plus discarding the skipped prefix. - **Latency**: the slowest shard's response, not the average — the merge cannot emit its first document until every stream has produced a head element. That last point is why fan-out queries have worse tail latency than their single-node equivalents even when every shard is individually fast: you are exposed to the worst of N samples on every request. ## What to do about it The cheapest fix is to make the query targetable at all — if the filter carries the shard key, one shard does the sort, the skip and the limit locally and there is no merge to speak of. Failing that, ensure an index supports the sort on every shard so the per-shard step is a bounded index walk rather than a blocking sort, and keep the offset small: paging by a range predicate on the last value seen rather than by a growing numeric offset keeps the per-shard fetch at roughly the page size no matter how deep the user pages. ## Interview framing The distinguishing detail here is *why* limit pushes down and skip does not. Anyone can memorise the behaviour; the candidate who can explain that a limit is monotone over the union of per-shard prefixes while a skip depends on cross-shard ordering has actually understood the router's job.

  • Why can a limit be pushed down to the shards but a skip cannot?
    A limit is monotone over the merge: the global first L documents must come from the per-shard first L, so asking each shard for L loses nothing. A skip depends on cross-shard ordering — a shard cannot know how many globally-earlier documents live on its peers, so skipping locally would silently drop documents that belong in the answer.
  • How much does the router hold in memory during a sorted merge?
    Roughly one batch per participating shard, not the whole result. Because each shard returns its portion already in sort order, the merge is a streaming k-way merge that emits documents as it goes. The expensive sorting, if any, happens on the shards.
  • Why does a fanned-out sorted query have worse tail latency than the same query on one node?
    The merge cannot emit anything until every shard has produced its first batch, so the response time is the maximum over N shards rather than a single sample. Any shard that is compacting, failing over, or simply unlucky sets the latency for the whole request.

saying these in an interview costs you the question

  • Says mongos sorts the entire merged result itself
  • Claims skip is pushed down to each shard
  • Thinks limit(20) transfers exactly 20 documents from the cluster
  • Assumes response time is the average shard's latency
  • Believes an index for the sort matters less once sharded

context