skip to content

How would you decide whether a $lookup-heavy pipeline can scale in a sharded cluster?

level: principalimportance: should knowfreq 34%

answer

  1. Ask how many times the join actually runs
  2. Two numbers multiply to give the cost
  3. One shard or all of them makes the difference
  4. Some fixes are not query fixes at all
  5. At some point you stop joining on read

basics

~20 s

Judge it by how many times the join runs and how well each run is targeted: the foreign side is looked up per input document, so shrink the left side first, index the foreign field, and check whether lookups can be routed to one shard. If they cannot, materialize instead.

solid answer

~60 s

Start from the execution shape: `$lookup` resolves the foreign side **once per input document flowing into the stage**, so total cost is roughly (left-side documents) x (cost of one foreign lookup). Three levers follow. First, cut the left side — `$match` and `$limit` before the `$lookup`, never after. Second, make each lookup cheap: index the `foreignField` (or the sub-pipeline's leading predicate) and `$project` inside the join so only needed fields cross. Third, ask where the lookups land. Recent MongoDB supports a sharded `from` collection, but if the join key is not the foreign collection's shard key, each lookup must be broadcast to every shard rather than targeted at one — cost then scales with shard count as well as document count. `$graphLookup` is stricter still: its `from` collection cannot be sharded at all. When those levers are exhausted, the answer stops being a query fix: pre-join with `$merge` into a materialized collection refreshed on a schedule, or accept that the read pattern is telling you the data should be stored together.

code

javascript · 6 lines
javascript
// Costly: joins every order, then filters
db.orders.aggregate([
  { $lookup: { from: "customers", localField: "customerId",
               foreignField: "_id", as: "customer" } },
  { $match: { status: "open", "customer.tier": "gold" } }
])

go deeper

for a junior

Recall the core cost fact: the foreign side is looked up once per document entering the stage, so filtering before the join is what makes it affordable.

for a middle

Explain the concrete levers — $match and $limit before the stage, an index on the foreign field, $project inside the sub-pipeline — and why an unbounded joined array risks the 16 MB document limit.

for a senior

Demonstrate routing judgment in a sharded cluster: whether foreign lookups can target one shard, why an untargeted join gets worse as shards are added, and that $graphLookup cannot use a sharded foreign collection at all.

for a principal

Own the call about where the join belongs: query-time versus materialized versus a different storage shape, argued from read/write ratio, freshness tolerance, growth expectations and workload isolation rather than from micro-optimizing one pipeline.

## Frame the cost before optimizing anything `$lookup` is not a set-based operation you can hand-wave about. The foreign side is resolved **per input document reaching the stage**. That single fact drives every decision: > total join cost ≈ (documents entering the $lookup) x (cost of one foreign lookup) So there are exactly two things to attack — the multiplier and the multiplicand — plus a third question specific to sharded clusters: *where* each lookup has to be executed. ## Lever 1 — shrink the left side Every document you eliminate before the `$lookup` is a foreign lookup you never pay for. In practice this means: - Put the selective `$match` **before** the join, and make sure it is index-supported. A `$match` after the join has already paid for every lookup. - If the query is paginated or top-N, get the `$sort` + `$limit` before the join too, so you join 50 documents rather than 500,000. - Beware the pipeline whose left side is a `$group` output: grouping first can *shrink* the left side dramatically, and joining on the grouped key is often orders of magnitude cheaper than joining raw documents and grouping afterwards. The optimizer does move some `$match` stages earlier on its own, but it cannot do so across stages that change the fields being filtered. Do not rely on it — write the pipeline in the order you want executed. ## Lever 2 — make one lookup cheap - **Index the foreign side.** The equality form needs an index on `foreignField`; the pipeline form needs the sub-pipeline's leading `$match` to be index-eligible. Joining on `_id` is free. Joining on an unindexed field means a scan of the foreign collection per input document — this is the single most common cause of a pipeline that ran fine on test data and collapsed in production. - **Trim what comes back.** The equality form returns whole foreign documents. Switch to the pipeline form and `$project` down to the fields you need; on a wide foreign collection this cuts both memory and network cost. - **Bound the fan-out.** A `$limit` inside the sub-pipeline caps how many documents each parent can pull. Remember the joined array lives inside one output document and is still bound by the 16 MB BSON limit — a hub document joining to hundreds of thousands of children will fail outright, not merely run slowly. ## Lever 3 — routing in a sharded cluster This is where the judgment gets genuinely architectural. Current MongoDB supports a **sharded `from` collection** for `$lookup` (earlier versions required it to be unsharded, which is why older material says joins and sharding do not mix). Support, however, is not the same as efficiency. The question is whether each individual lookup can be **targeted**: - If the field named by `foreignField` **is** (or is a prefix of) the foreign collection's shard key, a lookup can be routed to the one shard that owns that key range. - If it is not, every lookup has to be broadcast to every shard, and each shard does work for keys it does not own. Cost now scales with document count **and** shard count, and adding shards makes the join worse rather than better. The left side matters too: a pipeline whose left side is itself scatter-gather starts from a large document count on every shard, and the merging node becomes a bottleneck. A useful heuristic: if neither side of the join is targeted, the pipeline is doing an all-shards-to-all-shards operation, and no amount of index tuning will change its scaling behaviour. `$graphLookup` is more constrained: its `from` collection **cannot be sharded**. If the hierarchy lives in a sharded collection, recursive traversal is simply off the table there and must be solved another way. ## When to stop optimizing and change the shape Signals that the query is the wrong layer to fix: - The join is on the **hot read path** of the product's main screen, not a report. - The left side cannot be shrunk because the query is inherently "all of X". - The foreign lookups cannot be targeted and the cluster is expected to grow. - The pipeline already runs in the analytics window and still misses it. The realistic responses, in increasing order of commitment: **materialize** the joined result with `$merge` into a collection refreshed on a schedule and read that instead; **route** analytical pipelines to secondaries or a dedicated analytics node so they cannot destabilize the transactional workload; or accept the read pattern as a schema signal and change how the data is stored. That last one is not a query decision, and it is the one a lead should be willing to raise: `$lookup` is best understood as the escape hatch for the joins you did not design for, not as the foundation of a high-traffic read path. ## What to measure Before any of this is more than opinion, get the numbers: how many documents actually enter the `$lookup`, whether the foreign lookups are index-driven, and where the time is spent. Run the pipeline's explain output on representative data volumes — a join whose foreign collection fits in cache on a test box behaves nothing like the same join over a production working set.

  • What is the difference between a $lookup being supported on a sharded foreign collection and being efficient there?
    Support means the operation is legal; efficiency depends on routing. If the `foreignField` is the foreign collection's shard key (or its prefix), each lookup can be sent to the single shard owning that range. Otherwise every lookup is broadcast to all shards, so cost scales with shard count as well as document count — and adding capacity makes the join slower rather than faster.
  • How would you decide between fixing a slow $lookup pipeline and materializing its result with $merge?
    By read/write ratio and freshness tolerance. If the joined result is read far more often than the underlying data changes, and consumers accept data that is minutes old, materializing with `$merge` on a schedule converts many expensive joins into one, plus cheap indexed reads. If results must be transactionally current, or the query shape varies per request, keep joining and spend the effort on left-side reduction and indexing.
  • A $lookup-based dashboard query is destabilizing the transactional workload. What do you change first, before touching the schema?
    Isolate it. Route the analytical pipeline away from the nodes serving user traffic — a dedicated analytics member or secondary reads — so a heavy join can no longer contend with the write path. That buys time to do the real work: measure how many documents enter the stage, confirm the foreign lookups are index-driven, and decide whether the result should be materialized. Schema change is the last resort, not the first response.
  • Why can a $lookup pipeline that passes testing fail outright in production rather than just running slowly?
    Two hard limits. The joined array lives inside a single output document, so a hub record matching an enormous number of foreign documents breaches the 16 MB BSON document limit and errors. And aggregation stages have their own memory ceiling. Test data with small fan-out never reaches either, so the failure mode only appears against real key distributions — which is why representative data volumes matter more than a synthetic row count.

saying these in an interview costs you the question

  • Says $lookup is impossible on sharded collections
  • Assumes the join runs once for the whole collection
  • Puts the selective $match after the join
  • Expects adding shards to speed up an untargeted join
  • Treats an unbounded joined array as safe

context