Evaluate MongoDB aggregation efficiency with explain, examined counts, early selective stages, memory/disk-spill evidence, and explicit PyMongo cursor cleanup rather than tuning folklore.
Explain Aggregation Pipelines and Push Selective Work Toward Indexed Stages
Batch heterogeneous writes safely, interpret partial success, compare ordered and unordered execution, and use modern cross-namespace bulk APIs without assuming all-or-nothing behavior.
Learning outcomes
AtlasMart has a pipeline that returns the correct dashboard
number, but during peak load it consumes far more resources than
expected. This lesson builds a disciplined workflow for deciding
whether an aggregation is efficient: establish a logical result
oracle, use explain to measure the actual plan,
compare examined work and cardinality, understand blocking-stage
memory behavior, and change one pipeline/index assumption at a
time.
Use queryPlanner versus executionStats explain modes deliberately and understand that explain output format is not guaranteed stable.
Compare a logically expensive late-filter pipeline with an early-selective equivalent when semantics allow the rewrite.
Connect index use to the early $match/$sort prefix rather than claiming indexes accelerate every aggregation stage.
Interpret docs examined, keys examined, returned rows, sort/group memory, and disk-spill evidence as workload-specific signals.
Manage aggregation cursors correctly from PyMongo and keep application-side resource cleanup observable.
This lesson pins MongoDB Community Server
8.3.8 using
mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim, a disposable standalone mongod published only
on loopback 127.0.0.1:27056, and mongosh
2.10.0, and PyMongo 4.17.0. The
database is atlasmart; the collection is
orders_ch08_l5. Authentication and TLS are
disabled only for this isolated learning container. The lab
reads Feature Compatibility Version (FCV) and
allowDiskUseByDefault but does not change them.
Default read/write concern and primary read preference are
used on the standalone. Atlas, Search, KMS, Enterprise
Advanced, and paid services are not required. Commands were
reviewed against current official documentation; runtime
output was not generated here because
Docker/mongod/mongosh/PyMongo are unavailable in this
environment.
Expected result fragments are documentation-derived shapes and invariants, not copied from a generation-time MongoDB process. Exact explain trees, optimizer rewrites, field ordering, cursor IDs, execution counters, memory/disk-spill detail, and error text can vary with server patch, FCV, index state, dataset, and topology. Verify the logical result and the named execution evidence rather than comparing output byte-for-byte.
1. “Correct result” and “efficient execution” are separate acceptance criteria
A pipeline can be logically correct yet wasteful because it expands arrays before a selective filter, groups far more documents than necessary, performs an in-memory sort that an index could have supported earlier, or returns excessive data to the client. Conversely, a low examined count does not prove the result is semantically correct. Maintain two tests: a result oracle for correctness and execution evidence for cost.
| Evidence | Question it answers | What it does not prove |
|---|---|---|
| Winning plan / stage names | Which physical access path/stages MongoDB chose for this execution. | Future plans under different statistics/data are not guaranteed identical. |
totalDocsExamined |
How many collection documents the measured query portion examined. | It does not include every possible CPU/memory cost of later aggregation state. |
totalKeysExamined |
How many index keys were inspected. | Low keys alone do not prove low latency or correct shard targeting. |
| Returned/result cardinality | How much output survived. | It does not reveal hidden intermediate explosion unless you trace/explain it. |
usedDisk diagnostics |
Whether an eligible stage wrote temporary files due to memory pressure. | Absence in a tiny test does not prove production will never spill. |
| Client cursor lifecycle | Whether the app consumes/closes server-side cursor resources responsibly. | It does not tune the server pipeline itself. |
2. Seed a larger but still disposable synthetic workload
docker rm -f atlasmart-mongo-ch08-l5 2>/dev/null || truedocker volume rm atlasmart-mongo-ch08-l5-data 2>/dev/null || truedocker run -d --name atlasmart-mongo-ch08-l5 \ -p 127.0.0.1:27056:27017 \ -v atlasmart-mongo-ch08-l5-data:/data/db \ mongodb/mongodb-community-server:8.3.8-ubuntu2204-slimmongosh "mongodb://127.0.0.1:27056/atlasmart?directConnection=true" --quiet --eval \'printjson({server:db.version(), hello:db.hello().isWritablePrimary}); printjson(db.getSiblingDB("admin").runCommand({getParameter:1,featureCompatibilityVersion:1,allowDiskUseByDefault:1}))'
const c=db.getCollection("orders_ch08_l5");c.drop();const bulk=[];for (let i=0;i<2000;i++) { bulk.push({ _id:`o-${String(i).padStart(5,"0")}`, tenantId:i<1600?"tenant-a":"tenant-b", status:i%7===0?"cancelled":"paid", region:["west","east","north","south"][i%4], createdAt:new Date(Date.UTC(2026,7,1,0,0,i%60)), amountCents:1000+(i%400)*17, lines:[{category:i%2===0?"books":"electronics",qty:1+(i%3)}] });}c.insertMany(bulk,{ordered:true});c.createIndex({tenantId:1,status:1,region:1,createdAt:-1,_id:1});print("seeded",c.countDocuments({}));
Two thousand documents are enough to make examined-count differences visible without pretending to be a production benchmark. The dataset is intentionally synthetic and uniform; production distributions, cache state, storage latency, concurrency, shards, and document sizes will differ.
3. Controlled anti-pattern: expand and group before a selective tenant/region filter
The first pipeline unwinds every order, computes a boolean
wanted value from tenant/status/region, and only
then filters on that computed field. It returns the same
requested groups, but the selective source predicate is no
longer expressed as an index-eligible early $match.
This is especially problematic if arrays are large or the
pre-filter stream is high-cardinality.
const bad=[ { $unwind:"$lines" }, { $set:{ wanted:{ $and:[ {$eq:["$tenantId","tenant-a"]}, {$eq:["$status","paid"]}, {$eq:["$region","west"]} ] } }}, { $match:{wanted:true} }, { $group:{_id:"$lines.category",revenueCents:{$sum:"$amountCents"},rows:{$sum:1}} }, { $sort:{revenueCents:-1,_id:1} }];printjson(c.aggregate(bad).toArray());
4. Repair: push independent selective predicates toward the indexed prefix
The requested business predicate—tenant-a, paid, west—is independent of the line-level category aggregation, so it can safely be expressed directly against stored fields before unwind/group. That preserves the result while giving MongoDB an index-eligible selective prefix and reducing expanded rows and group work.
const good=[ { $match:{tenantId:"tenant-a",status:"paid",region:"west"} }, { $unwind:"$lines" }, { $group:{_id:"$lines.category",revenueCents:{$sum:"$amountCents"},rows:{$sum:1}} }, { $sort:{revenueCents:-1,_id:1} }];printjson(c.aggregate(good).toArray());
This is not a universal “move every match to the front” rule. A
predicate that depends on a field created by $set,
a group result, or another later stage cannot simply be moved
before that value exists. MongoDB’s optimizer can also perform
semantics-preserving rewrites automatically. Start from
correctness, then inspect the optimized execution.
5. Compare both plans with executionStats
function summarize(label,pipeline) { const e=c.explain("executionStats").aggregate(pipeline); const q=e.stages?.find(s => s.$cursor)?.$cursor ?? e; print("\n==",label,"=="); printjson({ explainVersion:e.explainVersion, winningPlan:q.queryPlanner?.winningPlan, totalKeysExamined:q.executionStats?.totalKeysExamined, totalDocsExamined:q.executionStats?.totalDocsExamined, nReturnedFromQueryPrefix:q.executionStats?.nReturned }); // Keep the complete explain object when diagnosing a real system: // printjson(e);}const lateResult=c.aggregate(bad).toArray();const earlyResult=c.aggregate(good).toArray();print("same logical result", EJSON.stringify(lateResult) === EJSON.stringify(earlyResult));summarize("late computed filter",bad);summarize("early selective filter",good);
same logical result trueresult: { _id: "books", revenueCents: 1493172, rows: 342 }The exact winning plans and examined counters are environment-dependent; compare the actual executionStats from your run.
MongoDB warns that explain output format is not guaranteed even under Stable API, so production tooling should avoid brittle assumptions about every nested field. The high-value comparison is conceptual: does the repaired pipeline use the tenant/status/region index prefix, and does it examine materially fewer documents before the expensive cardinality-changing stages?
explain bypasses existing plan-cache entries and
does not populate a new cache entry. It is a diagnostic
execution, not a passive recording of the exact production
request path. Use profiler/telemetry and real request metrics
alongside explain when investigating incidents.
6. Sort/group memory and allowDiskUse are operational controls, not tuning folklore
MongoDB documents that stages such as $group and a
$sort not supported by an index can require
substantial memory. Starting in MongoDB 6.0,
allowDiskUseByDefault controls whether stages that
exceed 100 MB can write temporary files by default; a command
can override this with allowDiskUse. Disk spill
prevents one class of memory error but can increase latency and
temporary I/O. “Just set allowDiskUse:true” is therefore not a
performance strategy.
printjson(db.adminCommand({getParameter:1,allowDiskUseByDefault:1}));print("indexes:"); printjson(c.getIndexes());print("result count",c.aggregate(good).toArray().length);// For a real spill investigation also inspect profiler/diagnostic logs for usedDisk.// Do not infer a spill from a successful aggregation alone.
The lesson deliberately does not allocate enough memory to force
a spill. In production, confirm actual
usedDisk evidence in profiler/diagnostic logs,
monitor temporary storage and latency, and determine why the
stage needs that state—high cardinality, late filtering, missing
sort index, large accumulator arrays, or genuinely unavoidable
workload.
7. Client cursor lifecycle is part of the pipeline
aggregate() returns a cursor. Drivers fetch results
in batches, so application behavior affects connection usage,
network transfer, and server cursor lifetime. The PyMongo
example uses a context manager to make cleanup explicit even
though this pipeline returns only one group document.
from pymongo import MongoClientclient = MongoClient("mongodb://127.0.0.1:27056/atlasmart?directConnection=true", serverSelectionTimeoutMS=3000)coll = client.atlasmart.orders_ch08_l5pipeline = [ {"$match": {"tenantId": "tenant-a", "status": "paid", "region": "west"}}, {"$group": {"_id": None, "count": {"$sum": 1}, "revenueCents": {"$sum": "$amountCents"}}},]with coll.aggregate(pipeline, batchSize=50) as cursor: for doc in cursor: print(doc)client.close()
A small batchSize is not automatically faster; it
changes network round trips and buffering. Tune only after
measuring result sizes, consumer speed, and connection/resource
behavior. Never pull a large collection to Python merely to
reproduce filtering/grouping that the server can execute more
selectively.
8. A practical aggregation performance review
| Review step | Concrete question |
|---|---|
| Result oracle | Do representative fixtures prove the pipeline answers the correct business question? |
| Input selectivity | Can tenant/status/date/security predicates run before unwind/group/sort? |
| Cardinality | How many documents exist after each expansion/reduction stage at p50/p95/p99 workloads? |
| Index evidence | Does explain show an index supporting the early predicate/order, and how many keys/docs are examined? |
| Blocking state | How many groups/sort entries/accumulator values must be retained? |
| Spill/temporary I/O | Did profiler/logs show usedDisk? What is the I/O and latency consequence? |
| Output | How many documents/bytes are returned, and are individual output documents below 16 MiB? |
| Topology | On a sharded cluster, which shards execute which stages and where is merge work performed? |
| Security | Are tenant predicates and authorization enforced server-side before data leaves its boundary? |
| Regression | Are explain/latency/cardinality expectations tested after schema/index/data-distribution changes? |
9. Verification, cleanup, and production judgment
Verification checklist
- The synthetic collection contains exactly 2,000 documents and the compound index exists.
- The late-filter and early-filter pipelines are checked against the intended tenant/status/region semantics.
- Explain is captured for both and the actual winning access path/counters are compared.
- The repaired pipeline filters before unwind/group only because those predicates are semantically independent of later fields.
-
allowDiskUseByDefaultis observed rather than changed, and no false disk-spill claim is made from the small lab. - The PyMongo cursor is consumed in a context manager and the client is closed.
Aggregation optimization is workload engineering, not a one-time query rewrite. Production performance depends on document/array/group cardinality distributions, indexes, memory, cache warmth, storage, concurrent workload, read concern/preference, replica/shard topology, and server version. Track p50/p95/p99 latency, examined/returned ratios, plan changes, spill indicators, temporary-disk pressure, shard fan-out, and client cursor duration. Establish alerting based on your own service-level objectives rather than copying a universal threshold.
Security boundaries remain mandatory: a fast pipeline that scans
multiple tenants and filters in application code is unacceptable
even if the final response is correct. Read-only pipelines do
not need data rollback, but query/index changes need deployment
rollback and regression tests. Chapter 09 builds on these
fundamentals with $lookup, facets, recursive
traversal, windows, and write-back stages—features where
intermediate cardinality and memory discipline become even more
important.
docker rm -f atlasmart-mongo-ch08-l5docker volume rm atlasmart-mongo-ch08-l5-data
Check your understanding
- Why can a pipeline be correct but still operationally expensive?
- When is it safe to move a $match earlier?
- Why should explain JSON be parsed cautiously?
- What does allowDiskUse solve, and what does it not solve?
- Why is explicit client cursor cleanup part of production aggregation practice?
Review the answers
It may process/expand/group/sort far more data than necessary before producing the correct final result.
When the predicate depends only on fields already available and moving it preserves the exact query semantics; optimizer rewrites must also preserve semantics.
MongoDB does not guarantee a fixed explain output format, and plans can vary by version, FCV, indexes, and data. Focus on documented concepts and robust diagnostics.
It permits eligible memory-heavy stages to spill temporary state to disk instead of necessarily failing at the memory threshold. It does not make an intrinsically wasteful pipeline fast or cheap.
Aggregation returns a cursor whose batches and server/client resources can outlive one loop iteration; deterministic cleanup avoids leaking resources when processing stops early or errors.
Authoritative references
- MongoDB 8.3 release notes — Current stable/minor series and patch status; re-check before reproducing the lab.
- MongoDB versioning — Release-series and compatibility context for the pinned server.
- mongosh release notes — Current mongosh release used by the lesson commands.
- PyMongo release notes — Current official Python driver line used where driver cursor behavior is demonstrated.
- Aggregation pipeline — Core ordered-stage document-flow model and expression concepts.
- Aggregation pipeline limits — Stage count, 16 MiB output-document, memory, allowDiskUse, and disk-spill rules.
- Aggregation optimization — Stage reordering/coalescence and index-use opportunities.
- Explain command — Explain verbosity, plan-cache behavior, execution statistics, and output-format warning.
- Use indexes to sort — Index-supported sort and in-memory sort behavior.
- $match stage — Early selection and index eligibility.
- $group stage — Blocking group state, memory threshold, and spill behavior.
- PyMongo aggregation — Driver aggregation API and cursor-oriented execution.
- PyMongo cursors — Cursor iteration, batching, and cleanup context.