A $group over 300 million documents spills to disk and runs for an hour — how would you speed it up?
answer
- separate scan cost from group memory
- memory tracks distinct keys, not input size
- only the leading filter reaches an index
- array-building accumulators grow per group
- recomputing history every run is the real bug
basics
~20 sCut what reaches the $group: an index-backed $match at the head of the pipeline, fewer carried fields, and a coarser grouping key. $group memory scales with distinct keys and accumulator state, so if the key is genuinely huge, split the work by time window and materialize partial results incrementally.
solid answer
~50 sStart by locating the cost. `$group` holds one accumulator entry per **distinct key**, so its memory is driven by key cardinality and accumulator size — not by input volume, which drives scan time instead. Then attack in order. First, reduce the input: a selective leading `$match` that an index can serve, so the scan and the group both see less. Second, reduce the payload: keep only the fields the accumulators need, and avoid `$push`/`$addToSet` that materialise arrays per group. Third, reduce the key: grouping by user *and* minute may have a hundred times the cardinality of grouping by user *and* day. If the aggregate genuinely spans hundreds of millions of documents every run, stop recomputing it — run it incrementally over one window at a time and `$merge` the partials into a summary collection, then aggregate the summary. Enabling `allowDiskUse` is not a fix; it is what is already keeping the query alive.
code
javascript · 13 lines// Recomputes all history every run, groups on a very high-cardinality key
db.events.aggregate([
{ $group: { _id: { user: "$userId", minute: "$ts" }, hits: { $sum: 1 } } },
{ $match: { hits: { $gt: 5 } } }
])
// One bounded window, coarser key, result merged into a summary collection
db.events.aggregate([
{ $match: { ts: { $gte: start, $lt: end } } },
{ $project: { userId: 1, day: { $dateTrunc: { date: "$ts", unit: "day" } } } },
{ $group: { _id: { user: "$userId", day: "$day" }, hits: { $sum: 1 } } },
{ $merge: { into: "daily_user_hits", on: "_id", whenMatched: "replace" } }
])go deeper
You are not expected to lead this diagnosis, but know the first move: filter early with an index-backed $match so far fewer documents reach the $group at all.
Explain why $group memory depends on distinct-key count and accumulator type rather than input size, and why $push and $addToSet are the accumulators that blow up.
Show a diagnostic order rather than a grab-bag: read the plan, decide scan-bound versus memory-bound, then cut input, payload and key cardinality in that order, and reach for incremental materialization when the same history is recomputed every run.
Own the workload split: which computations get precomputed on a schedule, where heavy analytical passes run so they do not contend with latency-sensitive traffic, and what staleness the business will accept in exchange.
## Read the shape of the cost first Two different costs hide behind "the aggregation is slow", and they have different fixes. **Scan cost** is proportional to the number of documents and bytes the pipeline reads. If a `COLLSCAN` feeds 300 million documents into the pipeline, you pay for all of them before any grouping happens. **Group memory** is proportional to the number of *distinct* grouping keys times the size of each accumulator's state. Ten billion input documents grouping into fifty keys is cheap in memory; three hundred million documents grouping into two hundred million keys is not, and that is what spills. Getting this distinction right is most of the senior signal in the question. `explain("executionStats")` tells you which one you have: whether the leading stage was an index scan or a collection scan, and how many documents entered the pipeline. ## Step 1 — cut the input The single highest-leverage change is a selective `$match` as the pipeline's first stage, backed by an index. Only a leading `$match` reaches the query layer and can use an index; a filter written after the `$group` has already paid for everything. Most "grouping the whole collection" pipelines are actually reporting on a window — a date range, a tenant, a status — and simply were not written that way. If the report genuinely needs the full history, the input cannot shrink and you should skip to step 3. ## Step 2 — cut the payload and the accumulator state Documents carry their whole shape through the pipeline unless you narrow them. Projecting to the handful of fields the `$group` actually uses reduces bytes moved per document, which shows up directly in runtime on wide documents. The accumulators matter even more. `$sum`, `$avg`, `$min`, `$max` and `$count` keep a fixed few bytes per group. `$push` and `$addToSet` keep an array per group that grows without bound, and they are the usual reason a `$group` blows its budget at modest cardinality. They can also push a single output document past the 16 MB BSON limit, at which point no amount of disk spilling saves you. If you only need the top few members per group, get them another way rather than pushing everything and slicing afterwards. ## Step 3 — cut, or partition, the key Ask what the grouping key's cardinality actually is. Grouping by `{ userId, minute }` over a year is a vastly bigger key space than `{ userId, day }`, and the report often only needs the coarser one. Rounding a timestamp down before grouping is frequently a hundred-fold memory reduction for identical business meaning. When the key cannot be coarsened, partition the work. Run the pipeline over one time window at a time, so each run's key space is bounded, and combine the partials. ## Step 4 — stop recomputing history An aggregation that reads three hundred million documents on every execution is usually recomputing an answer that barely changed. The durable fix is incremental materialization: run the pipeline over the newest window, end it with `$merge` into a summary collection keyed by the grouping key plus the window, and let reports aggregate the (much smaller) summary. Each run then touches only the new data, and the expensive part of the query disappears from the read path entirely. ## What not to do - **Do not treat `allowDiskUse` as the fix.** It is why the query completes at all. From MongoDB 6.0 it is on by default, which means a spilling pipeline no longer announces itself with an error — it just gets slow and competes for I/O with live traffic. - **Do not add `$limit` after the `$group` and expect relief.** By then all the grouping work is done; the limit only trims the output. - **Do not index the grouping key and assume the group becomes cheap.** An index helps the leading filter and can supply ordering; it does not eliminate the accumulator table. - **Do not tune before measuring.** Confirm from the plan whether you are scan-bound or memory-bound, and change one thing at a time. ## How to present the answer An interviewer is listening for a diagnostic order rather than a list of tricks: identify whether the pain is scan volume or key cardinality, reduce the input with an index-backed leading filter, reduce per-document and per-group state, coarsen or partition the key, and — if the same expensive computation recurs on a schedule — change the shape of the work with incremental materialization instead of optimizing the same full scan forever. Mentioning that heavy analytical passes should be isolated from the nodes serving latency-sensitive traffic rounds it out.
- How do you tell whether the pipeline is scan-bound or memory-bound?Run `explain("executionStats")`. If the leading stage is a collection scan feeding hundreds of millions of documents into the pipeline, you are scan-bound and the fix is an index-backed leading `$match`. If the scan is already targeted but the `$group` still spills, the grouping key's cardinality or the accumulator state is the problem, and coarsening the key or dropping array accumulators is what helps.
- Why does replacing $push with $sum sometimes fix a failing pipeline outright?`$push` retains every contributing value in an array per group, so its memory grows with input size and each output document can exceed the 16 MB BSON limit no matter how much disk is available. `$sum` and `$avg` keep a fixed-size running value per group. Swapping to a scalar accumulator changes the memory profile from unbounded to bounded.
- What does incremental materialization with $merge look like for a daily report?Run the pipeline over the last day only, group to the report's key plus the day, and finish with `$merge` into a summary collection using those fields as the `on` key so a re-run of the same window overwrites rather than duplicates. Reports then aggregate the summary collection, which holds one document per key per day instead of the raw events.
Counting a city's votes: the pile you carry is not every ballot, it is one tally sheet per candidate — so the pain comes from how many candidates there are, not how many ballots.
saying these in an interview costs you the question
- Says enable allowDiskUse and considers the problem solved
- Thinks $group memory scales with input document count
- Adds $limit after $group expecting the group to get cheaper
- Indexes the grouping key and expects the accumulation to be free
- Never separates scan cost from accumulator memory in the diagnosis