聚合管道进阶
上一章的阶段解决了「单集合统计」。这一章解决三类更难的问题:跨集合关联、一次查询产出多组结果、以及行间计算(窗口函数)。这些能力让 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 都会爆内存把 $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 } }
]
}
]子管道的三个优势:
- 可以在关联时就过滤和排序,只拉回需要的数据
- 支持非等值关联条件(范围、多字段组合)
- 可以嵌套
$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 次全表扫描,灾难三条铁律:
foreignField必须有索引(_id天然有)- 先
$match缩小左侧文档数,再$lookup - 不要在几十万文档的结果上做
$lookup
关系数据库的优化器会根据数据量在 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 之前的阶段可以用索引,但每个子管道内部的 $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] }
} }
} }
])不设 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+ 可以 | 可以 |
| 能否写入分片集合 | 不行 | 可以 |
| 是否必须是最后阶段 | 是 | 是 |
| 典型场景 | 全量重算 | 增量更新物化视图 |
每 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做增量物化视图,是把慢聚合变快查询的标准手法- 下一章讲数据建模,前面所有性能问题的根源都在这里 →