skip to content

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

level: seniorimportance: nice to knowfreq 33%

answer

  1. The pipeline is cut into two halves
  2. One half runs everywhere in parallel
  3. explain names both halves and where the second runs
  4. Some stages pin the second half to a specific shard
  5. A leading key match removes the split entirely

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.

solid answer

~50 s

When an aggregation cannot be targeted to one shard, the planner cuts the pipeline in two. The **shards part** — typically the leading `$match`, `$project` and as much pre-aggregation as is safe — runs in parallel on every participating shard. The **merger part** runs on a single node and finishes the job: combining partial `$group` results, applying the final `$sort`, `$limit` and `$skip`. A sharded `explain()` shows this directly as `splitPipeline` with `shardsPart`, `mergerPart` and a `mergeType`. The merge normally runs on one of the shards; stages that need an unsharded collection — `$lookup` against an unsharded foreign collection, `$out` — force it onto the **primary shard** of the database. `$out` cannot write to a sharded collection at all, while `$merge` can. If the pipeline starts with a `$match` on the shard key, none of this applies: the whole pipeline runs on one shard and `splitPipeline` is null.

code

javascript · 9 lines
javascript
db.orders.aggregate([
  { $match:  { status: "OPEN" } },          // runs on every shard
  { $group:  { _id: "$orgId", total: { $sum: "$amount" } } },
  { $sort:   { total: -1 } },
  { $limit:  10 }
]).explain()
// splitPipeline.shardsPart : $match + partial $group
// splitPipeline.mergerPart : final $group + $sort + $limit
// splitPipeline.mergeType  : where the merger half runs

go deeper

for a junior

Know that an aggregation over a sharded collection runs partly on each shard and is finished off on one node, and that filtering by the shard key first avoids that entirely.

for a middle

Explain which stages typically land in the shards half versus the merger half, and read splitPipeline, shardsPart, mergerPart and mergeType out of a sharded explain.

for a senior

Reason about where the merge runs and what pins it — an unsharded foreign collection or $out forcing the primary shard — and why that makes one node a bottleneck for heavy pipelines.

for a principal

Own the pattern: decide when a cross-shard pipeline belongs on the request path at all, versus materialising it with $merge into a collection whose key makes the serving reads targeted.

## Why a pipeline has to be split An aggregation over a sharded collection reads data that lives on many shards, but stages like `$group`, `$sort` and `$limit` are defined over the **whole** result. So the planner divides the pipeline into two halves and runs them in different places. - The **shards part** runs on every participating shard, in parallel, against that shard's own data. The planner pushes as much down here as it safely can: the leading `$match` (which also determines *which* shards participate at all), `$project` to cut document size early, and a partial form of accumulating stages. - The **merger part** runs on a single node. It combines what the shards produced — finishing a `$group` by merging partial accumulator states, applying the final ordering, then the skip and limit. `$group` is the interesting case: it appears on **both** sides. Each shard groups its own documents locally and emits partial results; the merger combines the partials for the same key. That is what makes cross-shard aggregation scale at all — shipping raw documents to one node would not. ## Reading the split from explain Run `explain()` on the aggregation through `mongos`. The output contains a `splitPipeline` object with: - **`shardsPart`** — the stages executed on each shard. - **`mergerPart`** — the stages executed on the merging node. - **`mergeType`** — where the merger half runs. Alongside it, a `shards` map shows each participating shard's own query plan, so you can check index usage per shard the same way you would for a `find()`. This output is the fastest way to answer the two questions that matter in a review: *did my `$match` actually get pushed to the shards*, and *did the pipeline fan out at all*. ## When there is no split If the leading `$match` contains the shard key, the router targets one shard and ships the **entire** pipeline there. `splitPipeline` is then null and there is no merger half: no partial groups, no cross-shard transfer of intermediate results, no merge-node bottleneck. This is by far the most effective optimisation available for a sharded aggregation, and it is a modelling and query-shape decision, not a tuning knob. ## What pins the merge to a particular shard Normally the merger half can run on any participating shard. Some pipelines cannot, because a stage needs access to an **unsharded** collection, which lives on the database's **primary shard**: - `$lookup` whose foreign collection is unsharded. - `$out`, which writes its results into a collection. Those force the merge onto the primary shard, which turns one node into a serialisation point for the whole pipeline's output. On a heavy pipeline that shows up as one shard pinned at high CPU while the others idle after finishing their halves. ## $out versus $merge `$out` **cannot write to a sharded collection**. `$merge` can, and it can also insert, replace or update matched documents rather than replacing a whole collection. For materialising a cross-shard result at any real scale, `$merge` is the stage you want: it writes into a sharded target, and it lets you refresh a summary incrementally instead of rebuilding it. This is the shape of the standard fix for an expensive cross-shard pipeline on a request path: run it in the background, `$merge` the result into a collection whose shard key matches how the UI reads it, and serve the request with a targeted `find()`. ## $lookup across shards Recent MongoDB versions (5.1 and later) allow the foreign collection of a `$lookup` to be sharded. That removes an old hard restriction, but it does not make the join cheap: there is no co-location guarantee, because the two collections have different shard keys, so resolving the join for a given input document may itself require reaching other shards. Treat a large cross-shard `$lookup` as a fan-out per input document and design around it — by embedding the needed fields, or by materialising the joined shape — rather than assuming the planner will make it local. ## What to take into an interview Three points carry the answer: the pipeline is split into a shards half and a merger half; `explain()` names both halves and where the merger runs; and a leading `$match` on the shard key removes the split entirely. Add the `$out`-versus-`$merge` distinction and the primary-shard pinning, and you have covered what an interviewer is actually probing — whether you have looked at a sharded pipeline's explain output rather than only written pipelines.

  • Why does putting a $match on the shard key first change the whole execution shape?
    Because it makes the pipeline targetable: the router resolves the filter to one shard and ships the entire pipeline there. There is no split, no merger half and no cross-shard transfer of intermediate results — explain shows splitPipeline as null. It is the single most effective optimisation for sharded aggregations.
  • What is the difference between $out and $merge on a sharded cluster?
    $out cannot write to a sharded collection, and its merging half runs on the database's primary shard, which makes it a serialisation point for large results. $merge can write into a sharded collection and can insert, replace or update matched documents, so it is the stage to reach for when materialising results at scale.
  • Does a $lookup in a sharded pipeline join data locally on each shard?
    No. There is no co-location guarantee: the foreign collection has its own shard key, so looking up a document may itself require contacting other shards. Recent versions do allow the foreign collection to be sharded, but the join remains a fan-out per input document rather than a local join, which is why large $lookup pipelines scale poorly across shards.

saying these in an interview costs you the question

  • Says mongos always performs the merge itself
  • Thinks $group runs only on the merging node
  • Claims $out can write into a sharded collection
  • Assumes $lookup joins are co-located on each shard
  • Believes the split happens even for a shard-key-targeted pipeline

context