Aggregation pipeline — это то место, где MongoDB перестаёт быть «просто хранилищем документов» и превращается в полноценный движок обработки данных внутри БД: конвейер стадий, каждая из которых трансформирует поток документов дальше. Здесь — $match, $group, $unwind, $facet, и отдельно $lookup — псевдо-join между коллекциями и его реальная цена, оптимизатор конвейера, лимиты памяти и allowDiskUse. Клиентский угол показан на двух языках: Go и Java, с контрастом к SQL GROUP BY и JOIN для тех, кто приходит из реляционного мира.
Предыдущая статья серии — Индексы вглубь — закончилась на том, как планировщик выбирает индекс под конкретный запрос. Aggregation опирается ровно на этот механизм: первая стадия $match может пройти по индексу так же, как обычный find, и весь разговор про IXSCAN против COLLSCAN переносится сюда без изменений. Всё измеренное ниже снято на живом replica set (mongo:8.2.11, 3 узла) на сквозном датасете серии (seed=42: 50 000 пользователей, 5000 товаров, 200 000 заказов), стенд mongodb/aggregation и его Java-зеркало mongodb/java/aggregation. Один и тот же пайплайн прогнан двумя драйверами — Go (mongo-go-driver/v2) и Java (org.mongodb:mongodb-driver-sync 5.5.1) — на одном и том же кластере в рамках одного прогона.
В статье
- Стадии конвейера: поток документов
$lookupкак «join» между коллекциями — механика и реальная цена- Оптимизатор конвейера: порядок стадий и индекс на входе
- Лимиты памяти и
allowDiskUse - Expression language внутри стадий
- Go и Java: aggregation pipeline с двух клиентов
- Контраст: pipeline против SQL
GROUP BYиJOIN
Стадии конвейера: поток документов
Pipeline — это массив стадий: db.orders.aggregate([ {$match: …}, {$group: …}, {$sort: …} ]). Каждая стадия принимает на вход поток документов и выдаёт на выход поток документов; следующая стадия видит только то, что вышло из предыдущей. Это ключевое отличие от SQL, где вы описываете что хотите получить одним декларативным выражением, а планировщик сам решает порядок операций. В MongoDB порядок стадий вы задаёте руками — и он имеет значение (хотя оптимизатор кое-что переставляет за вас, см. ниже).
Базовый словарь стадий, вокруг которого построена почти любая агрегация:
$match— фильтр. АналогWHERE. Если стоит первой и поле проиндексировано — проходит по индексу, как обычный запрос.$group— группировка с аккумуляторами ($sum,$avg,$push,$addToSetи т.д.). АналогGROUP BY. Это blocking-стадия: чтобы выдать первую группу, ей нужно увидеть все входные документы.$unwind— разворачивает массив в отдельные документы: один документ с массивом из N элементов превращается в N документов. Прямого аналога в плоском SQL нет — это следствие того, что в документе живут вложенные массивы.$sort— сортировка. Тоже blocking. Может использовать индекс, если стоит рано и поле проиндексировано; иначе сортирует в памяти (с лимитом, о котором отдельный разговор).$facet— несколько под-пайплайнов над одним и тем же входом за один проход: например, сразу и постраничная выборка, и общий счётчик, и гистограмма по цене. Каждая грань (facet) — независимый суб-пайплайн.$lookup— соединение с другой коллекцией. Псевдо-join, самая дорогая стадия из ходовых.
Стадии делятся на streaming ($match, $project, $limit, $skip — пропускают документы по одному, не накапливая) и blocking ($group, $sort — обязаны собрать весь вход). Blocking-стадии — источник почти всех проблем с памятью: именно они могут упереться в лимит и потребовать спилла на диск.
200 000 документов"] --> M["$match: status = paid
IXSCAN(idx_status_agg)
без COLLSCAN"] M --> G["$group: по user_id
слит с $match в один $cursor (SBE):
IXSCAN → FETCH → GROUP"] G --> S["$sort
отдельная поздняя стадия
24 200 групп · usedDisk=false · spills=0"] S --> OUT["результат
keysExamined = docsExamined = 33 251
wall ~128 мс"] style M fill:#c9e4c5,stroke:#5b8a5e style G fill:#f4d9c6,stroke:#c67a4a
flowchart LR
IN["orders
200 000 документов"] --> M["$match: status = paid
IXSCAN(idx_status_agg)
без COLLSCAN"]
M --> G["$group: по user_id
слит с $match в один $cursor (SBE):
IXSCAN → FETCH → GROUP"]
G --> S["$sort
отдельная поздняя стадия
24 200 групп · usedDisk=false · spills=0"]
S --> OUT["результат
keysExamined = docsExamined = 33 251
wall ~128 мс"]
style M fill:#c9e4c5,stroke:#5b8a5e
style G fill:#f4d9c6,stroke:#c67a4a
$lookup как «join» между коллекциями — механика и реальная цена
$lookup соединяет текущий поток документов с другой коллекцией: для каждого входного документа он ищет совпадения в целевой коллекции по локальному и внешнему полю и подкладывает найденное в новое массивное поле. Синтаксически это выглядит как join, и разработчики из SQL-мира рефлекторно тянутся к нему при первой же нормализации данных. Но семантически это не merge-join и не hash-join, а nested-loop join: на каждый входной документ — отдельный поиск в целевой коллекции.
Чтобы измерить цену честно, стенд сравнивает два пайплайна одинакового объёма работы. Первый разворачивает позиции заказа и джойнит каталог товаров:
db.orders.aggregate([
{ $unwind: "$items" },
{ $lookup: {
from: "products",
localField: "items.product_id",
foreignField: "_id",
as: "product_doc"
}},
{ $unwind: "$product_doc" },
{ $group: {
_id: "$product_doc.category",
revenue: { $sum: { $multiply: ["$items.qty", "$items.price"] } }
}}
])Второй считает выручку по имени товара, уже лежащему в позиции заказа, — та же группировка на тех же ~600 тысячах развёрнутых строк, но без обращения к другой коллекции:
db.orders.aggregate([
{ $unwind: "$items" },
{ $group: {
_id: "$items.product_name",
revenue: { $sum: { $multiply: ["$items.qty", "$items.price"] } }
}}
])Замеры на живом кластере (обе цифры — Go и Java, один прогон):
| Go | Java | |
|---|---|---|
$lookup totalDocsExamined |
600462 | 600462 |
$lookup collectionScans |
0 (индекс _id_) |
0 (индекс _id_) |
с $lookup, wall-latency |
12.81 s | 13.15 s |
без $lookup (эквивалент), wall-latency |
412 ms | 433 ms |
| во сколько раз дороже | 31.1× | 30.4× |
Читать эту таблицу нужно внимательно. $lookup здесь не оформлен плохо: collectionScans=0, каждый поиск в products попадает в индекс _id_ (первичный ключ), лишнего скана коллекции нет. И всё равно джойн примерно в 31 раз дороже эквивалентной агрегации без него — потому что nested-loop делает по одному индексированному поиску на каждую из ~600 тысяч развёрнутых строк, и эти ~600 тысяч поисков складываются в 12–13 секунд против 400 миллисекунд. Индекс спасает от катастрофы (без него был бы полный скан products на каждую строку), но не отменяет саму природу цикла.
Честная оговорка про explain: $unwind, стоящий сразу после $lookup, слился внутрь самой $lookup-стадии — в explain это видно по полю unwinding внутри записи $lookup, отдельной записи $unwind после неё в списке стадий нет. Оптимизатор объединяет $lookup + следующий $unwind в одну физическую операцию, чтобы не материализовать промежуточный массив совпадений. Это разумно, но на суть — nested-loop по документу на строку — не влияет.
Практический вывод не «никогда не используйте $lookup», а «знайте его цену». Для сотни-другой входных документов джойн незаметен. Для потока в сотни тысяч строк это доминирующая статья расходов, и часто дешевле продублировать нужное поле в сам документ (как product_name в позиции заказа выше) — классический для документных БД размен нормализации на скорость чтения.
Оптимизатор конвейера: порядок стадий и индекс на входе
MongoDB не выполняет пайплайн буквально стадия-за-стадией — перед выполнением через него проходит оптимизатор конвейера, который переставляет и сливает стадии. Два самых важных преобразования: проброс $match к началу (чтобы отфильтровать раньше и по индексу) и слияние $match + $group в единый план исполнения.
Разберём на пайплайне «выручка и число заказов по пользователям, у которых есть оплаченные заказы, отсортировано». Есть индекс idx_status_agg на orders.status:
db.orders.aggregate([
{ $match: { status: "paid" } },
{ $group: {
_id: "$user_id",
orderCount: { $sum: 1 },
revenue: { $sum: "$total" }
}},
{ $sort: { revenue: -1 } }
])explain с executionStats показывает, что произошло на самом деле:
{
"stages": [
{ "$cursor": {
"queryPlanner": {
"winningPlan": {
"stage": "GROUP",
"inputStage": {
"stage": "FETCH",
"inputStage": {
"stage": "IXSCAN",
"indexName": "idx_status_agg",
"direction": "forward"
}
}
}
},
"executionStats": {
"nReturned": 24200,
"totalKeysExamined": 33251,
"totalDocsExamined": 33251
}
}},
{ "$sort": { "sortKey": { "revenue": -1 } } }
]
}Что здесь важно:
$matchи$groupСЛИЛИСЬ в единый SBE-план внутри одной записи$cursor:IXSCAN(idx_status_agg) → FETCH → GROUP. Это не два прохода, а один: сервер идёт по индексу, забирает документы и тут же группирует.- COLLSCAN отсутствует — вход конвейера прошёл по индексу
idx_status_agg, как обычныйfind({status:"paid"}). keysExamined == docsExamined == 33251— равенство-фильтр по единственному индексному полю: сколько ключей просмотрели, столько же документов достали, ни одного лишнего.- 24200 групп (пользователи, у которых есть оплаченные заказы) — результат группировки.
$sortостался отдельной поздней стадией: 24200 групп — маленький объём, сортировка прошла в памяти (usedDisk=false,spills=0).
Замеры (Go и Java, один прогон):
| Go | Java | |
|---|---|---|
| групп (users с paid-заказами) | 24200 | 24200 |
| keysExamined = docsExamined | 33251 | 33251 |
| index | idx_status_agg (IXSCAN, без COLLSCAN) |
тот же |
| wall-latency | 128 ms | 197 ms |
Отдельная перекрёстная проверка, которую делает стенд: totalDocsExamined из explain (33251) совпадает с суммой orderCount по фактическому результату — то есть одни и те же paid-заказы посчитаны двумя независимыми путями (планировщиком и агрегацией), и цифры сошлись. Это тот же принцип верификации, что и в статье про индексы: доверяй не оценке планировщика, а executionStats реального прогона.
Практическое следствие оптимизатора: ставьте $match как можно раньше. Формально оптимизатор и сам протолкнёт простой $match к началу, но он не всесилен — $match после $group или после $project, переименовавшего поле, протолкнуть уже нельзя, и вы получите полный скан на входе. Индекс работает на первой стадии; всё, что фильтрует поток до первой blocking-стадии, экономит и просмотренные документы, и память.
Теперь $unwind + $group без всяких джойнов — «сколько выручки принёс каждый товар», через разворачивание массива позиций:
db.orders.aggregate([
{ $unwind: "$items" },
{ $group: {
_id: "$items.product_id",
revenue: { $sum: { $multiply: ["$items.qty", "$items.price"] } },
qty: { $sum: "$items.qty" }
}}
])| Go | Java | |
|---|---|---|
| уникальных product_id | 5000 (= весь каталог) | 5000 |
unwound-строк (items[] по всем orders) |
600461 | 600461 |
| revenue из группировки | 376779509.73 | 376779509.73 |
revenue независимой проверки ($sum(orders.total), без unwind) |
376779509.73 | 376779509.73 |
| wall-latency | 462 ms | 492 ms |
Ассерт стенда: сумма revenue, посчитанная через $unwind+$group по позициям, совпадает (diff = 0.0000, допуск 0.01) с независимо посчитанным $sum(orders.total) по неразвёрнутой коллекции. Один и тот же итог — 376779509.73 — получен двумя разными путями агрегации и сошёлся байт-в-байт на обоих клиентах. Это не косметика: совпадение доказывает, что $unwind не потерял и не задвоил ни одной позиции, а $multiply в аккумуляторе даёт ровно то, что уже лежит в orders.total.
Лимиты памяти и allowDiskUse
Blocking-стадии ($group, $sort) обязаны накопить данные в памяти. Исторически на каждую такую стадию действовал жёсткий лимит 100 МиБ (104 857 600 байт); при превышении агрегация падала с ошибкой, и единственным выходом было передать опцию allowDiskUse: true, разрешающую спилл промежуточных данных на диск.
Стенд демонстрирует лимит на $unwind + $sort по полям, не коррелирующим с индексом (~600 тысяч строк, которые нужно отсортировать целиком в памяти):
| Go | Java | |
|---|---|---|
С явным allowDiskUse:false |
ошибка, код 292 (QueryExceededMemoryLimitNoDiskUseAllowed), 167 ms |
ошибка, код 292, 184 ms |
С allowDiskUse:true |
успех, 600461 строк, 3.31 s | успех, 600461 строк, 4.16 s |
Текст ошибки с кодом 292 — дословно с сервера:
{
"code": 292,
"codeName": "QueryExceededMemoryLimitNoDiskUseAllowed",
"errmsg": "Sort exceeded memory limit of 104857600 bytes, but did not opt in to external sorting."
}Обе цифры (Go и Java) сходятся: с allowDiskUse:false — реальная ошибка сервера 292, с allowDiskUse:true — реальный успех, ровно 600461 строк (число совпадает с независимой проверкой $unwind+$count). Оба направления, оба клиента.
Ключевая находка: allowDiskUseByDefault = true с MongoDB 6.0
А вот здесь — самое частоупускаемое место, и подать его нужно точно. Формулировка «без allowDiskUse» в таблице выше означает поле allowDiskUse, явно равное false в команде, а не отсутствие поля.
Причина — легко упускаемый серверный параметр allowDiskUseByDefault, который по умолчанию true, начиная с MongoDB 6.0. Если per-query опция allowDiskUse вообще не передана в команде, сервер подставляет значение этого параметра — то есть спилл на диск разрешён по умолчанию. Старый жёсткий лимит 100 МиБ форсирует только явно переданный на уровне команды allowDiskUse: false.
Это проверено на стенде отдельным сравнением сырых команд через *event.CommandMonitor (перехват того, что реально уходит на сервер) на том же датасете и той же команде:
{"aggregate":"orders", ..., "allowDiskUse": false, ...} // → код 292, ~150-220 ms
{"aggregate":"orders", ...(поля allowDiskUse нет вовсе)...} // → успех, ~3.5 s, 600461 строк
{"aggregate":"orders", ..., "allowDiskUse": true, ...} // → успех, ~3.4 s, 600461 строкПодчеркну, чтобы не было соблазна прочитать это как «сервер противоречит документации»: это документированное, намеренное поведение сервера. allowDiskUseByDefault=true — стандартная настройка с версии 6.0, и она ровно про то, чтобы агрегации по умолчанию не падали на лимите 100 МиБ, а спиллили на диск. «Сюрприз» здесь только в том, что параметр легко не заметить: пропущенная опция ведёт себя как true, и стенд специально передаёт allowDiskUse:false явно — иначе сценарий «без флага» даёт ложный успех, потому что дефолт уже разрешает диск.
Оговорка: allowDiskUse:true спасает не всё (код 146)
И вторая честная оговорка, тоже проверенная вживую: не любой blocking-аккумулятор можно вытолкнуть на диск. Тот же объём данных, но собранный в один аккумулятор на всю коллекцию — $group{ _id: null, docs: { $push: "$$ROOT" } } (одна группа, весь поток в один массив) — падает даже при allowDiskUse:true:
{
"code": 146,
"codeName": "ExceededMemoryLimit",
"errmsg": "$push used too much memory and cannot spill to disk. Memory limit: 104857600 bytes"
}Обратите внимание: код здесь 146 (ExceededMemoryLimit), а не 292. Это другой лимит и другой механизм. allowDiskUse умеет вытолкнуть на диск внешнюю сортировку ($sort) и группировку по многим группам, но не единственный распухший $push-аккумулятор одной группы — такой аккумулятор спиллить на диск сервер не умеет. Именно поэтому демонстрация лимита памяти в этой серии построена на $sort, а не на $group+$push: сортировка честно показывает пару «292 / успех со спиллом», а $push завёл бы в тупик 146, который allowDiskUse не лечит.
Резюме по памяти: включённый по умолчанию спилл (allowDiskUseByDefault=true) — это удобно, но не бесплатно (диск медленнее памяти: 3.3 с против ошибки за 167 мс — это ещё и сигнал, что стадия не влезает в память), и не универсально (код 146 остаётся). Явный allowDiskUse:false полезен как предохранитель: если вы хотите, чтобы тяжёлая незапланированная агрегация упала быстро, а не молча съела диск и время.
Expression language внутри стадий
Внутри стадий MongoDB живёт отдельный выразительный язык — expression operators. Это то, что стоит внутри $group, $project, $addFields: $sum, $avg, $multiply, $cond, $switch, $dateToString, $map, $filter, $reduce и десятки других. Строка "$field" — это ссылка на поле (path), "$$ROOT" — весь текущий документ, "$$NOW" — серверное время.
В пайплайнах выше уже встречались $sum: 1 (счётчик), $sum: "$total" (сумма поля) и $multiply: ["$items.qty", "$items.price"] (произведение двух полей позиции внутри аккумулятора). Разница между $sum: 1 и $sum: "$total" ровно та же, что между COUNT(*) и SUM(total) в SQL, только записана как выражение-аргумент аккумулятора.
Ключевое отличие от SQL: expression language работает над вложенной структурой документа. $items.qty — это обращение к полю внутри элемента массива после $unwind; $map/$filter/$reduce умеют пройти по массиву без разворачивания его в отдельные документы. В плоском реляционном мире аналога нет — там сначала пришлось бы нормализовать позиции в отдельную таблицу. В документной модели массив — часть документа, и язык выражений умеет с ним работать на месте. Это и делает $unwind + $group таким естественным для позиций заказа, и одновременно объясняет, почему цена $lookup так велика: джойн — это выход за пределы текущего документа, всё остальное expression language делает внутри него.
Go и Java: aggregation pipeline с двух клиентов
Пайплайн — это данные (BSON-массив стадий), а не код. Поэтому любой драйвер отправляет на сервер одну и ту же команду aggregate, и вся работа происходит на сервере одинаково. Отличается только то, как клиент собирает этот BSON и как передаёт опцию allowDiskUse. Именно поэтому все цифры выше приведены парами Go/Java и совпадают по сути: разница в wall-latency (128 против 197 мс, 3.31 против 4.16 с) — это накладные расходы драйвера и JVM-прогрев, а не разная работа сервера. Число групп, просмотренные документы, revenue, коды ошибок — идентичны.
Go, mongo-go-driver/v2 — тот же $match → $group → $sort:
pipeline := mongo.Pipeline{
{{"$match", bson.D{{"status", "paid"}}}},
{{"$group", bson.D{
{"_id", "$user_id"},
{"orderCount", bson.D{{"$sum", 1}}},
{"revenue", bson.D{{"$sum", "$total"}}},
}}},
{{"$sort", bson.D{{"revenue", -1}}}},
}
cur, err := db.Collection("orders").Aggregate(ctx, pipeline)Java, mongodb-driver-sync 5.5.1 — тот же пайплайн через хелперы Aggregates:
List<Bson> pipeline = List.of(
Aggregates.match(Filters.eq("status", "paid")),
Aggregates.group("$user_id",
Accumulators.sum("orderCount", 1),
Accumulators.sum("revenue", "$total")),
Aggregates.sort(Sorts.descending("revenue"))
);
AggregateIterable<Document> result = db.getCollection("orders").aggregate(pipeline);И решающая деталь про память — обе опции передаются per-query, симметрично. В Go:
opts := options.Aggregate().SetAllowDiskUse(false) // явный false → код 292
cur, err := coll.Aggregate(ctx, pipeline, opts)
// SetAllowDiskUse(true) → успех, 600461 строкВ Java:
coll.aggregate(pipeline).allowDiskUse(false); // явный false → код 292
// .allowDiskUse(true) → успех, 600461 строкВажный клиентский вывод из находки про allowDiskUseByDefault: если вы не вызываете SetAllowDiskUse / .allowDiskUse(…), драйвер вообще не кладёт поле в команду — и сервер применяет свой дефолт true. То есть «я не трогал allowDiskUse» на MongoDB 6.0+ означает «диск разрешён», а не «действует лимит 100 МиБ». Чтобы получить жёсткий лимит как предохранитель, false нужно передать явно, на обоих клиентах одинаково.
Контраст: pipeline против SQL GROUP BY и JOIN
Для тех, кто приходит из реляционного мира, полезно разложить соответствие по полочкам — и заодно увидеть, где аналогия ломается.
| SQL | Aggregation pipeline | Комментарий |
|---|---|---|
WHERE status='paid' |
{ $match: { status: "paid" } } |
Идёт по индексу, если стоит первой стадией |
GROUP BY user_id |
{ $group: { _id: "$user_id", … } } |
_id группы = ключ группировки |
SUM(total), COUNT(*) |
{ $sum: "$total" }, { $sum: 1 } |
Аккумуляторы внутри $group |
ORDER BY revenue DESC |
{ $sort: { revenue: -1 } } |
Blocking, если не по индексу |
JOIN products ON … |
{ $lookup: { from: "products", … } } |
Nested-loop, не hash/merge-join |
| — (нет прямого аналога) | { $unwind: "$items" } |
Следствие вложенных массивов |
HAVING … |
$match после $group |
Фильтр по результату группировки |
Три различия, которые стоит держать в голове:
-
Порядок задаёте вы. SQL декларативен — планировщик сам выбирает порядок соединений и фильтров. В pipeline порядок стадий — часть запроса; оптимизатор переставит простой
$matchк началу и сольёт$match+$group, но не спасёт от$match, поставленного после$group. Раннее фильтрование — ваша ответственность. -
$lookup— не полноценный join. Реляционный оптимизатор дляJOINвыберет hash-join или merge-join в зависимости от объёмов и индексов.$lookupв этом типичномlocalField/foreignField-сценарии ведёт себя как indexed nested-loop: документ-на-строку. Даже с индексом (collectionScans=0) это обошлось в 31× относительно эквивалента без join. В реляционной БД тот же join на тех же объёмах, скорее всего, выбрал бы hash-join и был бы куда ближе к безджойновому варианту. Отсюда документный рефлекс: денормализуй то, что читаешь часто. -
$unwindне имеет аналога в плоском SQL. Массивы внутри документа — норма для MongoDB;$unwindразворачивает их в строки, чтобы группировать. В реляционной модели этих данных просто не было бы в одной таблице — они жили бы в отдельнойorder_items, и$unwindбыл бы не нужен, зато был бы нужен ещё одинJOIN. Это тот же самый размен: нормализация против вложенности.
Общий вывод: aggregation pipeline даёт мощь SQL-агрегаций плюс работу с вложенной структурой документа, но переносит на разработчика два решения, которые в SQL берёт на себя планировщик, — порядок операций и стратегию соединения. Знание реальной цены стадий (особенно $lookup) и серверных умолчаний (особенно allowDiskUseByDefault) — ровно то, что отличает работающий на игрушечных данных пайплайн от пайплайна, который держит продакшн.
Дальше в серии — Replica sets, oplog и write/read concerns: как эти же агрегации ведут себя на кластере из нескольких узлов, откуда читать тяжёлую аналитику, чтобы не мешать записи, и какие гарантии даёт oplog. Живые стенды к этой статье — Go: mongodb/aggregation, Java: mongodb/java/aggregation.
Комментарии