Learn
MongoDB/11-aggregation-advanced

聚合管道进阶

上一章的阶段解决了「单集合统计」。这一章解决三类更难的问题:跨集合关联、一次查询产出多组结果、以及行间计算(窗口函数)。这些能力让 MongoDB 能承担相当一部分原本要靠 SQL 或数据仓库做的工作。

1. $unwind:展开数组

1.1 基本行为

db.posts.aggregate([
  { $match: { _id: 11 } },
  { $unwind: "$tags" }
])

原始文档:

{ "_id": 11, "title": "A", "tags": ["mongodb", "index", "perf"] }

输出:

[
  { "_id": 11, "title": "A", "tags": "mongodb" },
  { "_id": 11, "title": "A", "tags": "index" },
  { "_id": 11, "title": "A", "tags": "perf" }
]

一个文档变三个,tags 从数组变成单值。这相当于关系模型里把「一对多」的宽表展平。

1.2 完整语法

{ $unwind: {
    path: "$tags",
    includeArrayIndex: "tagIdx",          // 输出元素下标
    preserveNullAndEmptyArrays: true      // 数组为空/字段缺失时保留文档
} }
⚠️默认会丢掉空数组的文档

简写形式 { $unwind: "$tags" } 在遇到 tags 为 []、null 或字段不存在时,会直接丢弃整个文档。统计「有多少帖子」时如果先 unwind 了 tags,没有标签的帖子就凭空消失了。需要保留必须显式写 preserveNullAndEmptyArrays: true。

1.3 文档膨胀风险

100 万篇帖子,每篇平均 5 个标签
      ↓ $unwind
500 万个中间文档
 
100 万篇帖子,每篇平均 200 条内嵌评论
      ↓ $unwind
2 亿个中间文档   ← 后续任何 $group 都会爆内存
💡unwind 之前先 match,之后立刻 group

把 $match 放在 $unwind 前面尽量减少输入文档数,$unwind 之后立刻用 $group 收敛。永远不要让展开后的文档流经多个阶段。

2. $lookup:关联查询

2.1 基础写法(等值关联)

db.posts.aggregate([
  { $match: { status: "published" } },
  { $lookup: {
      from: "users",              // 关联的集合
      localField: "authorId",     // 本集合字段
      foreignField: "_id",        // 目标集合字段
      as: "author"                // 结果放进哪个字段(永远是数组)
  } },
  { $unwind: "$author" },         // 一对一时展开成对象
  { $project: { title: 1, "author.username": 1, "author.profile.city": 1 } }
])
[
  {
    "_id": 11,
    "title": "MongoDB 入门",
    "author": { "username": "alice", "profile": { "city": "Beijing" } }
  }
]

对应 SQL:

SELECT p.title, u.username, u.city
FROM posts p LEFT JOIN users u ON p.author_id = u.id
WHERE p.status = 'published';

注意 $lookup 的语义是 LEFT OUTER JOIN:没匹配到时 as 字段是空数组,而不是丢弃文档。

2.2 子管道写法(更强大)

只想关联一部分数据、或者关联条件不是简单等值时,用子管道形式:

db.users.aggregate([
  { $match: { username: "alice" } },
  { $lookup: {
      from: "posts",
      let: { uid: "$_id" },                       // 把外层字段传进子管道
      pipeline: [
        { $match: {
            $expr: { $eq: ["$authorId", "$$uid"] },   // 关联条件
            status: "published",                       // 额外过滤
            "stats.views": { $gte: 100 }
        } },
        { $sort: { createdAt: -1 } },
        { $limit: 5 },
        { $project: { title: 1, "stats.views": 1, _id: 0 } }
      ],
      as: "topPosts"
  } }
])
[
  {
    "_id": "alice",
    "username": "alice",
    "topPosts": [
      { "title": "索引原理", "stats": { "views": 3400 } },
      { "title": "聚合管道", "stats": { "views": 1200 } }
    ]
  }
]

子管道的三个优势:

  1. 可以在关联时就过滤和排序,只拉回需要的数据
  2. 支持非等值关联条件(范围、多字段组合)
  3. 可以嵌套 $lookup,做多级关联

let 定义的变量在子管道里用双美元符号引用($$uid),子管道内的字段用单美元符号($authorId)。

2.3 多字段关联

{ $lookup: {
    from: "inventory",
    let: { sku: "$sku", wh: "$warehouse" },
    pipeline: [
      { $match: { $expr: { $and: [
        { $eq: ["$sku", "$$sku"] },
        { $eq: ["$warehouse", "$$wh"] }
      ] } } }
    ],
    as: "stock"
} }

2.4 性能:$lookup 是嵌套循环

这是最重要的一点。$lookup 的实现是对左侧每一个文档,去右侧集合查一次:

左侧 1000 篇帖子
   │
   ├─ 帖子 1 → 查 users 集合一次
   ├─ 帖子 2 → 查 users 集合一次
   ├─ ...
   └─ 帖子 1000 → 查 users 集合一次
 
  总共 1000 次查询。
  如果 users._id 有索引 → 1000 次索引查找,可接受
  如果没索引 → 1000 次全表扫描,灾难

三条铁律:

  1. foreignField 必须有索引(_id 天然有)
  2. 先 $match 缩小左侧文档数,再 $lookup
  3. 不要在几十万文档的结果上做 $lookup
⚠️$lookup 不是 SQL 的 JOIN

关系数据库的优化器会根据数据量在 hash join、merge join、nested loop join 之间选择,还能做谓词下推和 join 重排序。MongoDB 的 $lookup 只有嵌套循环一种实现,也不会重排 join 顺序。如果你的查询需要三个以上集合关联,说明建模需要调整——第 12 章会讲怎么用内嵌和扩展引用模式消灭这些关联。

2.5 $lookup 后常用的整形技巧

// 一对一:展开成对象
{ $unwind: { path: "$author", preserveNullAndEmptyArrays: true } }
 
// 只取第一个
{ $set: { author: { $arrayElemAt: ["$author", 0] } } }
 
// 只要计数,不要内容
{ $set: { commentCount: { $size: "$comments" } } }
{ $unset: "comments" }

3. $facet:一次扫描,多组结果

分页接口通常需要「当前页数据 + 总条数」,两次查询意味着两次扫描。$facet 让一份输入同时喂给多条子管道:

db.posts.aggregate([
  { $match: { status: "published", tags: "mongodb" } },
  { $facet: {
      data: [
        { $sort: { createdAt: -1 } },
        { $skip: 0 },
        { $limit: 20 },
        { $project: { title: 1, authorId: 1, createdAt: 1 } }
      ],
      total: [ { $count: "value" } ],
      byLevel: [
        { $group: { _id: "$level", n: { $sum: 1 } } }
      ]
  } }
])
[
  {
    "data": [
      { "_id": 11, "title": "MongoDB 入门", "authorId": "alice", "createdAt": ISODate("2024-06-01T00:00:00Z") }
    ],
    "total": [ { "value": 320 } ],
    "byLevel": [ { "_id": "A", "n": 42 }, { "_id": "B", "n": 278 } ]
  }
]

典型用途是电商的「搜索结果 + 各维度筛选项计数」:一次聚合同时算出商品列表、品牌分布、价格区间分布、总数。

⚠️$facet 内部无法使用索引

$facet 之前的阶段可以用索引,但每个子管道内部的 $sort / $match 都是内存操作。所以务必在 $facet 之前用 $match 把数据量压到几千条以内。用 $facet 处理百万级数据是常见的性能陷阱。

4. 分桶:$bucket 与 $bucketAuto

4.1 $bucket:手动指定边界

db.posts.aggregate([
  { $bucket: {
      groupBy: "$stats.views",
      boundaries: [0, 100, 1000, 10000, 100000],
      default: "other",                              // 落在边界外的
      output: {
        count: { $sum: 1 },
        avgLikes: { $avg: "$stats.likes" },
        samples: { $push: "$title" }
      }
  } }
])
[
  { "_id": 0,     "count": 60100, "avgLikes": 1.2,   "samples": ["..."] },
  { "_id": 100,   "count": 31500, "avgLikes": 12.4,  "samples": ["..."] },
  { "_id": 1000,  "count": 4200,  "avgLikes": 88.1,  "samples": ["..."] },
  { "_id": 10000, "count": 380,   "avgLikes": 512.6, "samples": ["..."] }
]

_id 是桶的下界,区间是左闭右开。

4.2 $bucketAuto:自动等频分桶

db.posts.aggregate([
  { $bucketAuto: {
      groupBy: "$stats.views",
      buckets: 5,
      granularity: "R20"          // 可选,让边界取整成好看的数字
  } }
])
[
  { "_id": { "min": 0,    "max": 45   }, "count": 20112 },
  { "_id": { "min": 45,   "max": 180  }, "count": 20098 },
  { "_id": { "min": 180,  "max": 620  }, "count": 20105 },
  { "_id": { "min": 620,  "max": 2400 }, "count": 20090 },
  { "_id": { "min": 2400, "max": 98000}, "count": 20075 }
]

$bucket 适合「已知业务分档」(比如价格区间),$bucketAuto 适合「先看看数据分布」(做直方图)。

5. $graphLookup:递归关联

处理树形和图结构。内容社区里的典型场景是多级评论和关注链。

5.1 多级评论树

// comments: { _id, postId, body, parentId }
db.comments.aggregate([
  { $match: { _id: ObjectId("665f1a2b3c4d5e6f70819301") } },
  { $graphLookup: {
      from: "comments",
      startWith: "$_id",              // 从哪个值开始
      connectFromField: "_id",        // 用当前文档的哪个字段继续找
      connectToField: "parentId",     // 去目标集合匹配哪个字段
      as: "replies",
      maxDepth: 5,                    // 最多递归 5 层
      depthField: "level"             // 输出每个节点的深度
  } }
])
[
  {
    "_id": ObjectId("...301"),
    "body": "根评论",
    "replies": [
      { "_id": ObjectId("...302"), "body": "一级回复", "level": 0 },
      { "_id": ObjectId("...303"), "body": "二级回复", "level": 1 },
      { "_id": ObjectId("...304"), "body": "三级回复", "level": 2 }
    ]
  }
]

5.2 二度人脉

// follows: { followerId, followeeId }
db.users.aggregate([
  { $match: { username: "alice" } },
  { $graphLookup: {
      from: "follows",
      startWith: "$_id",
      connectFromField: "followeeId",
      connectToField: "followerId",
      as: "network",
      maxDepth: 1,                    // 0 = 直接关注,1 = 加上二度
      depthField: "degree"
  } },
  { $project: {
      username: 1,
      secondDegree: { $filter: {
        input: "$network", as: "n", cond: { $eq: ["$$n.degree", 1] }
      } }
  } }
])
⚠️graphLookup 必须限制 maxDepth

不设 maxDepth 时它会一直递归到没有新节点为止。在关注关系这种连通性很强的图上,从任意一个用户出发几步就能覆盖全网,结果文档会瞬间超过 16MB 限制并报错。真正的社交图计算应该用专门的图数据库或离线计算,$graphLookup 适合的是「层级不深、分支不多」的树形结构。

6. $setWindowFields:窗口函数(5.0+)

SQL 的 OVER (PARTITION BY ... ORDER BY ...) 终于有了对应实现。它和 $group 的区别是:不减少文档数量,只是给每个文档附加一个基于「窗口内其他行」计算出的值。

6.1 排名

db.posts.aggregate([
  { $match: { status: "published" } },
  { $setWindowFields: {
      partitionBy: "$authorId",                     // 分区,相当于 PARTITION BY
      sortBy: { "stats.views": -1 },                // 窗口内排序
      output: {
        rankInAuthor: { $rank: {} },
        denseRank: { $denseRank: {} },
        rowNum: { $documentNumber: {} }
      }
  } },
  { $match: { rankInAuthor: { $lte: 3 } } }         // 每个作者的 Top 3
])
[
  { "title": "索引原理",   "authorId": "alice", "stats": { "views": 3400 }, "rankInAuthor": 1 },
  { "title": "聚合管道",   "authorId": "alice", "stats": { "views": 1200 }, "rankInAuthor": 2 },
  { "title": "MongoDB 入门","authorId": "alice", "stats": { "views": 900 }, "rankInAuthor": 3 }
]

「分组取 Top N」在没有窗口函数时非常难写,现在只要两个阶段。

6.2 移动平均与累计

db.dailyStats.aggregate([
  { $setWindowFields: {
      sortBy: { date: 1 },
      output: {
        ma7: {
          $avg: "$pv",
          window: { documents: [-6, 0] }          // 当前行往前 6 行,共 7 天
        },
        cumulative: {
          $sum: "$pv",
          window: { documents: ["unbounded", "current"] }
        },
        prevDay: { $shift: { output: "$pv", by: -1, default: 0 } }
      }
  } },
  { $set: { growth: { $subtract: ["$pv", "$prevDay"] } } }
])
[
  { "date": "2024-06-01", "pv": 1200, "ma7": 1200,   "cumulative": 1200, "growth": 1200 },
  { "date": "2024-06-02", "pv": 1350, "ma7": 1275,   "cumulative": 2550, "growth": 150  },
  { "date": "2024-06-03", "pv": 1180, "ma7": 1243.3, "cumulative": 3730, "growth": -170 }
]

窗口的两种定义方式:

类型写法含义
文档窗口documents: [-6, 0]按行数,前 6 行到当前行
范围窗口range: [-7, 0], unit: "day"按 sortBy 字段的值域,近 7 天

常用窗口操作符:$sum、$avg、$min、$max、$count、$rank、$denseRank、$documentNumber、$shift、$first、$last、$derivative、$integral、$expMovingAvg。

💡窗口函数省掉了自关联

在窗口函数出现之前,「和上一天比增长了多少」只能靠把集合和自己 $lookup 一次,或者拉到应用层算。现在 $shift 一行搞定,而且是流式处理,内存友好。

7. 输出阶段:$out 与 $merge

7.1 $out:整体覆盖

db.posts.aggregate([
  { $group: { _id: "$authorId", posts: { $sum: 1 }, views: { $sum: "$stats.views" } } },
  { $out: "authorStats" }
])

管道结果完全替换 authorStats 集合(先写临时集合,成功后原子改名)。适合每天全量重算的报表。

7.2 $merge:增量合并(4.2+)

db.posts.aggregate([
  { $match: { createdAt: { $gte: new Date("2024-06-04") } } },
  { $group: {
      _id: { author: "$authorId", day: { $dateToString: { date: "$createdAt", format: "%Y-%m-%d" } } },
      posts: { $sum: 1 },
      views: { $sum: "$stats.views" }
  } },
  { $merge: {
      into: "authorDailyStats",
      on: "_id",                        // 匹配键,需要唯一索引
      whenMatched: "merge",             // replace | keepExisting | merge | fail | 自定义管道
      whenNotMatched: "insert"          // insert | discard | fail
  } }
])

whenMatched 还能写成一个管道,实现「累加」而不是「覆盖」:

{ $merge: {
    into: "authorDailyStats",
    on: "_id",
    whenMatched: [
      { $set: {
          posts: { $add: ["$posts", "$$new.posts"] },
          views: { $add: ["$views", "$$new.views"] }
      } }
    ],
    whenNotMatched: "insert"
} }

$$new 引用的是管道产出的新文档,$posts 引用的是目标集合里已有的值。

对比$out$merge
目标集合整体替换增量合并
能否写入其他库4.4+ 可以可以
能否写入分片集合不行可以
是否必须是最后阶段是是
典型场景全量重算增量更新物化视图
💡用 $merge 做物化视图

每 5 分钟跑一次增量 $merge,把统计结果写进一个小集合,前端直接查这个小集合。这是把「慢聚合」变成「快查询」的标准做法,也是 MongoDB 版本的物化视图。

8. 其他实用阶段

{ $sample: { size: 100 } }                      // 随机抽样
{ $sortByCount: "$tags" }                       // group + count + sort 的语法糖
{ $replaceRoot: { newRoot: "$author" } }        // 把某个内嵌文档提升为顶层
{ $replaceWith: "$author" }                     // 同上,简写
{ $densify: { field: "date", range: { step: 1, unit: "day", bounds: "full" } } }  // 补齐缺失日期
{ $fill: { sortBy: { date: 1 }, output: { pv: { method: "linear" } } } }          // 填充空值
{ $unionWith: { coll: "archivedPosts", pipeline: [ { $match: { year: 2023 } } ] } } // 合并另一个集合
🎯练习

一、用 $lookup 子管道,查询 alice 的资料并附带她最近 5 篇已发布帖子的标题;二、用 $facet 实现一个分页接口,一次返回当前页 20 条数据、总条数、以及按标签分组的计数;三、用 $bucketAuto 画出 posts.stats.views 的 10 桶直方图;四、用 $graphLookup 查出某条根评论下的完整回复树,限制 3 层;五、用 $setWindowFields 找出每个作者浏览量最高的前 2 篇帖子;六、写一个 $merge 管道,把每日作者统计增量写入 authorDailyStats,要求重复执行不会导致重复累加(提示:whenMatched 用 replace 而不是自定义累加管道)。

小结

  • $unwind 展开数组,默认丢弃空数组文档,注意文档膨胀
  • $lookup 是 LEFT JOIN 且实现为嵌套循环,foreignField 必须有索引
  • 子管道形式的 $lookup 能在关联时过滤排序,是生产环境的首选写法
  • 需要三个以上集合关联时,应该回头调整建模而不是硬写聚合
  • $facet 一次扫描出多组结果,但子管道内部不走索引,前面要先压数据量
  • $graphLookup 处理树形结构,必须设 maxDepth
  • $setWindowFields 提供排名、移动平均、行间差值,替代了自关联
  • $merge 做增量物化视图,是把慢聚合变快查询的标准手法
  • 下一章讲数据建模,前面所有性能问题的根源都在这里 →