Lesson 10 of 18

Aggregation Pipeline

A Pipeline Is a Conveyor Belt

find() can filter and sort, but it cannot answer questions like "what were total sales per category last month?" or "which five customers spent the most?". Those need grouping and arithmetic across many documents, and that is what the aggregation pipeline is for. It is MongoDB's equivalent of SQL's GROUP BY, HAVING and JOIN — but expressed as a list of steps rather than one statement.

The mental model is a conveyor belt. Documents enter the first stage, that stage transforms or filters them, and whatever comes out is fed into the next stage. Each stage is a document with a single key naming the operation. Stages can repeat: it is perfectly normal to have two $match stages in one pipeline, at different points.

The reason to learn this rather than pulling documents into Node.js and looping over them is simple. Aggregation runs inside the database, next to the data, using the indexes. Doing the same work in your application means moving every document across the network first, which is slower, uses far more memory, and gets worse as your data grows.

Example
// Total sales per category, best first, top 5
db.orders.aggregate([
  { $match: { status: "completed" } },
  { $group: {
      _id: "$category",
      totalSales: { $sum: "$amount" },
      avgSale:    { $avg: "$amount" },
      orders:     { $sum: 1 }
  }},
  { $sort: { totalSales: -1 } },
  { $limit: 5 }
])
// [
//   { _id: 'Electronics', totalSales: 482300, avgSale: 6031.25, orders: 80 },
//   { _id: 'Furniture',   totalSales: 195000, avgSale: 4875,    orders: 40 },
//   ...
// ]
Notes
  • Build a pipeline one stage at a time. Run it with only $match, look at the output, then add $group, look again. Debugging a six-stage pipeline you wrote all at once is much harder than growing it a stage at a time.

$group and Its Accumulators

$group is the stage that does the real work, and its _id field is the part people misread. Inside $group, _id does not mean the document's identifier — it means the key you are grouping by. Setting _id: "$city" produces one output document per distinct city. Setting _id: null puts everything into a single group, which is how you get a grand total across the whole collection.

The dollar sign in "$city" matters too. A bare string like "city" is the literal text; "$city" means "the value of the field named city in the current document". Forgetting the dollar sign gives you one group whose key is the word "city", which is a confusing result until you know why.

Every other field in $group is built with an accumulator. $sum adds values — and $sum: 1 adds one per document, which is how you count. $avg, $min and $max do what their names say. $push collects values into an array, and $addToSet does the same while removing duplicates. You can also group by several fields at once by making _id an object.

Example
// Count users per city
db.users.aggregate([
  { $group: { _id: "$city", users: { $sum: 1 } } },
  { $sort: { users: -1 } }
])

// One grand total for the whole collection
db.orders.aggregate([
  { $group: { _id: null, revenue: { $sum: "$amount" }, count: { $sum: 1 } } }
])

// Group by two keys, and collect a list
db.orders.aggregate([
  { $group: {
      _id: { city: "$city", status: "$status" },
      revenue:  { $sum: "$amount" },
      orderNos: { $push: "$orderNo" }
  }}
])

// Sales per month, using a real date field
db.orders.aggregate([
  { $group: {
      _id: { $dateToString: { format: "%Y-%m", date: "$placedAt" } },
      revenue: { $sum: "$amount" }
  }},
  { $sort: { _id: 1 } }
])
Notes
  • The monthly example only works because placedAt was stored as a real date. If it had been saved as the string "15/01/2026", none of the date operators would apply. This is where the advice about types from the inserting lesson pays off.

Stage Order Decides Your Performance

The same stages in a different order can produce the same answer at wildly different cost. The rule that matters most: filter early. A $match at the very start of a pipeline can use an index and can throw away most of the collection before any expensive work begins. The identical $match placed after a $group runs over results that are already computed, so you paid to process every document and then discarded most of the output.

The second rule is to shrink early. If your pipeline only needs three fields, a $project near the front means every later stage moves smaller documents. On documents carrying large arrays or long text, this makes a visible difference.

There is a nuance worth knowing: MongoDB's optimiser will move some stages around on its own, and it can push a $match earlier when it is safe to do so. But it cannot do this once your fields have been renamed or computed, because it can no longer prove the filter means the same thing. Write the pipeline in the efficient order yourself and you never depend on the optimiser's judgement.

Example
// Slow: groups the whole collection, then throws most of it away
db.orders.aggregate([
  { $group: { _id: "$customerId", spend: { $sum: "$amount" } } },
  { $match: { _id: { $in: [101, 102, 103] } } }
])

// Fast: filter first (uses an index on customerId), then group
db.orders.aggregate([
  { $match: { customerId: { $in: [101, 102, 103] } } },
  { $group: { _id: "$customerId", spend: { $sum: "$amount" } } }
])

// Filter, shrink, then work
db.orders.aggregate([
  { $match: { placedAt: { $gte: ISODate("2026-01-01") } } },
  { $project: { customerId: 1, amount: 1, _id: 0 } },
  { $group: { _id: "$customerId", spend: { $sum: "$amount" } } },
  { $sort: { spend: -1 } },
  { $limit: 10 }
])
Notes
  • You can call .explain() on an aggregation exactly as you can on a find(). It shows whether the leading $match used an index, which is usually the first thing to check when a pipeline is slow.

$unwind: One Document per Array Element

Orders usually hold an array of items. To answer "how many units of each product did we sell?", you first need each item to be its own document, and that is what $unwind does: an order with three items becomes three documents, identical except that items now holds a single item instead of the array.

Once flattened, the rest is an ordinary $group over the unwound field. This pairing — $unwind followed by $group — is one of the most common shapes in real pipelines.

The gotcha is what happens to documents whose array is empty or missing. By default $unwind drops them entirely, because zero elements means zero output documents. If an order with no items should still appear in your report, pass the object form with preserveNullAndEmptyArrays: true. Silent disappearance of rows is a hard bug to spot in a total, so decide deliberately which behaviour you want.

Example
// Units sold per product
db.orders.aggregate([
  { $match: { status: "completed" } },
  { $unwind: "$items" },
  { $group: {
      _id: "$items.sku",
      unitsSold: { $sum: "$items.qty" },
      revenue:   { $sum: { $multiply: ["$items.qty", "$items.price"] } }
  }},
  { $sort: { unitsSold: -1 } }
])

// Keep orders that have no items at all
db.orders.aggregate([
  { $unwind: { path: "$items", preserveNullAndEmptyArrays: true } }
])
Notes
  • After $unwind, a count of documents is a count of array elements, not of orders. If you need both numbers, count the orders before unwinding, or use $addToSet on the order id and take the size of that set.

$lookup: Joining Two Collections

$lookup is MongoDB's join. You name the other collection in from, the field on the current documents in localField, the field on the other collection in foreignField, and the name of the new field in as. For each incoming document, MongoDB finds all the matching documents in the other collection.

The result is always an array, even when exactly one document matched and even when none did. That trips people up: after a lookup joining an order to its single customer, customer is a one-element array, not an object. Follow it with $unwind to flatten it, and use preserveNullAndEmptyArrays if orders whose customer was deleted should still show up.

One performance rule outweighs everything else here: the foreignField in the other collection must be indexed. Without an index, MongoDB looks up each incoming document against the other collection separately, and a pipeline that ran in milliseconds on your test data takes minutes in production. If foreignField is _id, you are already covered.

Example
// Attach each order's customer document
db.orders.aggregate([
  { $match: { status: "completed" } },
  { $lookup: {
      from: "users",
      localField: "customerId",
      foreignField: "_id",
      as: "customer"
  }},
  { $unwind: { path: "$customer", preserveNullAndEmptyArrays: true } },
  { $project: { orderNo: 1, amount: 1, "customer.name": 1, "customer.city": 1 } }
])

// The other direction: every user with their orders attached
db.users.aggregate([
  { $lookup: { from: "orders", localField: "_id", foreignField: "customerId", as: "orders" } },
  { $addFields: { orderCount: { $size: "$orders" } } },
  { $match: { orderCount: { $gt: 0 } } }
])

// Make the join fast
db.orders.createIndex({ customerId: 1 })
Notes
  • If you find yourself writing $lookup in almost every query, that is a signal about your schema rather than about aggregation. The next two lessons cover when related data should have been embedded in the first place.

Other Stages and Practical Details

A handful of remaining stages round out most day-to-day work. You will not need all of them at once, but knowing they exist saves you from writing awkward workarounds.

  • $project — choose, rename and compute fields, exactly like a projection in find() but with arithmetic available
  • $addFields (also spelled $set) — add computed fields while keeping everything already there
  • $sort, $skip, $limit — the same operations as on a cursor, usable at any point in the pipeline
  • $count: "total" — replace the stream with a single document holding the count
  • $sample: { size: 5 } — pick documents at random, useful for a "featured products" strip
  • $merge and $out — write the pipeline's results into a collection; both must be the final stage, and $out replaces the target collection entirely
Notes
  • Aggregation stages have a memory limit per stage. A large $group or $sort can exceed it; passing { allowDiskUse: true } as the second argument to aggregate() lets those stages spill to temporary files instead of failing. Reach for it when a report over historical data errors out, but treat it as a hint that an earlier $match or a supporting index would serve you better.
Ask AI