MongoDB 系列 · 第四篇

聚合管道:一条 pipeline 是怎么被拆开执行的

写一条 aggregate() pipeline 的时候,很容易默认它就是"stage 1 跑完再跑 stage 2,顺序执行"。真实情况不是这样——查询规划器会重排、合并、甚至把某些 stage 直接下推成 index 查询的一部分,跟你写的顺序未必一样。这一篇用真实 explain() 输出和真实撞上的一次内存报错,把这些"看不见的改写"一条条挖出来,再对照 MongoDB Server 真实源码确认每一条都不是巧合。

本机 blogdemo 数据库上跑了 6 组真实 pipeline,又特意造了一批大文档去真实触发 $push 的内存上限——错误信息里的每一个数字都能在源码常量里对上。

1 次真实崩溃
刻意造出 150MB 数据触发 $push 的 100MB 单数组上限,拿到真实报错信息
4 处真实源码
$match 与前置 stage 的交换规则、$sort 吸收 $limit 的机制、$push 内存上限常量

真实实验:写了 3 个 stage,顶层只剩 2 个

在建过索引的 inventory 集合上跑一条三段 pipeline:按分类过滤、按状态分组求和、按总量排序。

真实实测
$ mongosh --quiet blogdemo --eval '
  db.inventory.aggregate([
    { $match: { category: "electronics" } },
    { $group: { _id: "$status", totalQty: { $sum: "$qty" }, count: { $sum: 1 } } },
    { $sort: { totalQty: -1 } }
  ], { explain: true })
'
顶层 stages: [ "$cursor", "$sort" ] $cursor.queryPlanner.winningPlan.queryPlan: { stage: 'GROUP', inputStage: { stage: 'FETCH', inputStage: { stage: 'IXSCAN', keyPattern: { category: 1, qty: 1 }, indexBounds: { category: [ '["electronics", "electronics"]' ], qty: [ '[MinKey, MaxKey]' ] } } } }

写的是 3 个 stage,顶层 stages 数组却只有 $cursor$sort 两项——$match 变成了 IXSCAN 的索引边界(category: ["electronics","electronics"]),$group 直接接在 FETCH 之上,跟 $match 一起组成了同一棵查询执行计划,根本不是管道里单独的一步。只有 $sort 因为要等分组结果出来才能排,留在了管道顶层。

为什么会被重排——不是玄学,是一条可读的规则

下面两节各用一对真实的"能推 / 不能推"例子验证同一条规则:能不能交换,只看字段有没有冲突。

规则一

$match 能不能往前挪

只要 $match 用到的字段,跟它前面那个 stage 会修改的字段完全不沾边,规划器就能把 $match 交换到前面去,一路推到查询层。

规则二

$sort 能不能吃掉后面的 $limit

$sort 会看一眼自己后面紧跟的是不是 $limit,是的话直接把这个数字吸收进排序执行器,变成只维护 top-k 的堆排序。

规则三

$push 的内存上限不认 allowDiskUse

单个 $push 累加数组有独立的 100MB 硬上限,源码里的报错原文直接写着"不能溢出到磁盘"——这条跟 allowDiskUse 无关。

真实源码:$match 能不能跟前一个 stage 交换

splitSourceBy 用一个字段集合的"是否相交"判断来决定能不能交换。

src/mongo/db/pipeline/document_source_match.cpp · mongodb/mongo @ r8.3.7, L446 pair<intrusive_ptr<DocumentSourceMatch>, intrusive_ptr<DocumentSourceMatch>> DocumentSourceMatch::splitSourceBy(const OrderedPathSet& fields, const StringMap<std::string>& renames) && { ... if (!newExpr.second && renames.empty()) { // This $match is entirely independent of 'fields' and there were no renames to apply. In // this case, the current stage can swap with its predecessor without modification. We // simply return this as the first stage in the pair. _matchProcessor->setExpression(std::move(newExpr.first)); return {this, nullptr}; }

fields 是前一个 stage(比如 $addFields)会修改或产生的字段集合。如果 $match 的谓词跟这个集合完全不相交,就直接整体挪到前面(return {this, nullptr});挪不动的部分留在原地,挪得动的部分继续网前找,一直到不能再挪为止——这正是上面 $match → IXSCAN 边界 这一路下推的起点。

真实实测
$ mongosh --quiet blogdemo --eval '
  db.inventory.aggregate([
    { $addFields: { qtyDoubled: { $multiply: ["$qty", 2] } } },
    { $match: { category: "electronics" } }   // 只依赖 category,跟 addFields 无关
  ], { explain: true }).stages[0].$cursor.queryPlanner.parsedQuery
'
{ category: { '$eq': 'electronics' } } // 非空 -- $match 被交换到 $addFields 前面,推进了 $cursor # 换成 $match 依赖 addFields 刚算出来的字段: $ mongosh --quiet blogdemo --eval ' db.inventory.aggregate([ { $addFields: { qtyDoubled: { $multiply: ["$qty", 2] } } }, { $match: { qtyDoubled: { $gt: 100 } } } // 依赖 qtyDoubled,跟 addFields 冲突 ], { explain: true }).stages[0].$cursor.queryPlanner.parsedQuery ' {} // 空 -- 交换不成立,$match 老实留在 $addFields 后面

两条 pipeline 只有 $match 依赖的字段不一样,结果天差地别:第一条 parsedQuery 直接就是完整的过滤条件,第二条整个是空的——跟 splitSourceBy 判断"字段是否相交"的逻辑完全对上。

真实源码:$sort 怎么把 $limit 吃进去

按未建索引的字段排序再截取前几条,是个很常见的写法——规划器不会真的把全部数据排完。

src/mongo/db/pipeline/document_source_sort.cpp · mongodb/mongo @ r8.3.7, L292 DocumentSourceContainer::iterator DocumentSourceSort::optimizeAt( DocumentSourceContainer::iterator itr, DocumentSourceContainer* container) { ... auto stageItr = std::next(itr); auto limit = extractLimitForPushdown(stageItr, container); if (limit) _sortExecutor->setLimit(*limit);

$sortoptimizeAt 会去看紧跟在自己后面的下一个 stage,如果能提取出一个 $limit 值,就直接设置到自己的排序执行器上。排序算法一旦知道"我只需要前 N 个",就能用一个大小为 N 的堆维护候选,不需要在内存里对全部输入做一次完整排序。

真实实测
$ mongosh --quiet blogdemo --eval '
  db.inventory.aggregate([
    { $sort: { price: 1 } },
    { $limit: 5 }
  ], { explain: true }).queryPlanner.winningPlan
'
{ stage: "SORT", sortPattern: { price: 1 }, limitAmount: 5, // <- $limit 的 5 被直接吸收进了这个 SORT 节点 ... } # 整条 pipeline 甚至没有独立的 "stages" 数组了,跟一次普通的 # find().sort({price:1}).limit(5) 变成了同一棵执行计划

limitAmount: 5 直接出现在 SORT 节点里——这条两段的 pipeline 完全没有留下 $sort/$limit 分开的痕迹,规划器把它整个压成了一次 top-5 堆排序。

真实实验:$push 撞上 100MB 上限——allowDiskUse 救不了

故意造 5000 篇每篇 30KB 的文档(总共约 150MB),用 $push 把它们塞进同一个数组。

src/mongo/db/query/query_execution_knobs.idl · mongodb/mongo @ r8.3.7, L214 internalQueryMaxPushBytes: description: >- Limits the vector of values pushed into a single array while grouping with the $push accumulator. default: expr: 100 * 1024 * 1024
src/mongo/db/pipeline/accumulator_push.cpp · mongodb/mongo @ r8.3.7, L56 _memUsageTracker.add(input.getApproximateSize()); uassert(ErrorCodes::ExceededMemoryLimit, str::stream() << "$push used too much memory and cannot spill to disk. Memory limit: " << _memUsageTracker.maxAllowedMemoryUsageBytes() << " bytes", _memUsageTracker.withinMemoryLimit());
真实实测
$ mongosh --quiet blogdemo --eval '
  db.bigblob.aggregate([
    { $group: { _id: "$group", items: { $push: "$pad" } } }
  ]).toArray()
'
MongoServerError: Executor error during aggregate command on namespace: blogdemo.bigblob :: caused by :: Used too much memory for a single array. Memory limit: 104857600 bytes. The array contains 3411 elements and is of size 104833674 bytes. The element being added has size 30734 bytes. # 加上 allowDiskUse:true 原样重跑: $ mongosh --quiet blogdemo --eval ' db.bigblob.aggregate([ { $group: { _id: "$group", items: { $push: "$pad" } } } ], { allowDiskUse: true }).toArray() ' MongoServerError: ... Used too much memory for a single array. Memory limit: 104857600 bytes. The array contains 3411 elements and is of size 104833674 bytes. The element being added has size 30734 bytes. // 一模一样,allowDiskUse 完全没起作用

104857600 正是源码里 internalQueryMaxPushBytes 的默认值(100 * 1024 * 1024);报错原文"cannot spill to disk"直接对应真实报错里 allowDiskUse:true 完全不起作用这件事——这不是 bug,是 $push 这个累加器自己的设计:它的内存占用等于最终要输出的那个 BSON 数组本身,而单个 BSON 值终归要整个进内存,没有"边算边落盘"这回事。真正能被 allowDiskUse 缓解的,是 $group/$sort 需要同时维护大量分组或候选的场景,不是单个累加数组本身超限。

交互演示:6 个真实场景,一次走完三条规则

把上面全部真实实验按发生顺序串成一条演示,每一步标注命中了哪条规则、真实证据是什么。

聚合管道优化实录未开始
点击"下一步"或"播放"开始。

判断逻辑由 Python 脚本按上面三段真实源码逐条转写:字段相交判断直接照抄 splitSourceBy 的不相交检查;$sort 吸收 $limit 的结果和真实 explain()limitAmount 断言相等;$push 内存上限的模拟用真实报错里的元素大小(30734 字节)反推,算出的"3411 个元素、104833674 字节"跟真实报错逐字节一致,还有一个用整除法独立算出同一个数字的交叉校验。

参考与说明

  • 本文源码引用(DocumentSourceMatch::splitSourceByDocumentSourceSort::optimizeAtAccumulatorPush::processInternalinternalQueryMaxPushBytes)均取自 mongodb/mongo 仓库 r8.3.7 标签,与本机安装的 MongoDB Community 8.3.7 版本一致,直接从 GitHub 拉取源文件核对过函数签名、关键调用和行号。
  • 全部 explain() 输出、真实报错信息均为本机真实 mongosh 会话产生,未做删改;bigblob 集合是专门为触发 100MB 上限造的测试数据(5000 篇 × 30KB)。
  • 演示数据的自检:两个"能不能交换"场景的模型判定,和真实 $cursor.queryPlanner.parsedQuery 是否为空逐一断言相等;$push 上限模拟用一个完全独立的整除法重新算过一遍"能装下几个元素",跟累加循环的结果一致,而且两者都跟真实报错的 3411/104833674 这两个数字精确相等。
  • 没有涉及:分片集群上聚合管道的拆分与合并(mongos 与各 shard 各跑一段,留给最后一篇分片策略)、$lookup/$graphLookup 的执行方式、窗口函数($setWindowFields)、聚合管道的 SBE 与经典引擎选择逻辑细节。
☕ 如果这篇文章帮到你,可以请作者喝杯咖啡 · 爱发电