Aggregation pipeline вглубь

Стадии aggregation pipeline, цена $lookup как join, оптимизатор конвейера, allowDiskUse и лимиты памяти — с Go и Java клиентами на живых данных, в контрасте с SQL GROUP BY и JOIN

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) — на одном и том же кластере в рамках одного прогона.

Агрегационный конвейер данных: маскот-лист запускает рычаг, стрелка IXSCAN на входе; стадии 1)$match 2)$group 3)$unwind 4)$facet 5)$lookup (join со слоном-Postgres), сброс на диск allowDiskUse при переполнении памяти; роботы Go и Java с одинаковым результатом; панель управления пропускной способностью

В статье

Стадии конвейера: поток документов

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-стадии — источник почти всех проблем с памятью: именно они могут упереться в лимит и потребовать спилла на диск.

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

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
Пайплайн $match → $group → $sort на orders (status = paid): где включается индекс и что сливает оптимизатор. Числа — с живого стенда mongodb/aggregation, mongo:8.2.11, seed=42

$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 Фильтр по результату группировки

Три различия, которые стоит держать в голове:

  1. Порядок задаёте вы. SQL декларативен — планировщик сам выбирает порядок соединений и фильтров. В pipeline порядок стадий — часть запроса; оптимизатор переставит простой $match к началу и сольёт $match+$group, но не спасёт от $match, поставленного после $group. Раннее фильтрование — ваша ответственность.

  2. $lookup — не полноценный join. Реляционный оптимизатор для JOIN выберет hash-join или merge-join в зависимости от объёмов и индексов. $lookup в этом типичном localField/foreignField-сценарии ведёт себя как indexed nested-loop: документ-на-строку. Даже с индексом (collectionScans=0) это обошлось в 31× относительно эквивалента без join. В реляционной БД тот же join на тех же объёмах, скорее всего, выбрал бы hash-join и был бы куда ближе к безджойновому варианту. Отсюда документный рефлекс: денормализуй то, что читаешь часто.

  3. $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.

Обсуждение в Telegram

Присоединиться →

Комментарии