聚合管道基础
find 只能做「过滤 + 投影 + 排序」,做不了分组统计、字段计算、跨集合关联。聚合管道(aggregation pipeline)补上了这块,它是 MongoDB 里对应 SQL 的 GROUP BY、JOIN、窗口函数的那套东西,而且能力更强。
1. 管道模型
一条聚合是一个阶段数组,文档像流水一样依次穿过每个阶段:
集合中的文档
│
▼
┌─────────┐
│ $match │ 过滤,相当于 WHERE
└────┬────┘
▼
┌─────────┐
│ $group │ 分组统计,相当于 GROUP BY
└────┬────┘
▼
┌─────────┐
│ $sort │ 排序,相当于 ORDER BY
└────┬────┘
▼
┌─────────┐
│ $limit │ 截断,相当于 LIMIT
└────┬────┘
▼
结果文档关键在于「流」这个字:每个阶段的输出是下一个阶段的输入,文档在管道里会被反复改造,最终形态可能和原始文档完全不同。
db.posts.aggregate([
{ $match: { status: "published" } },
{ $group: { _id: "$authorId", total: { $sum: 1 }, views: { $sum: "$stats.views" } } },
{ $sort: { views: -1 } },
{ $limit: 5 }
])[
{ "_id": "alice", "total": 42, "views": 128500 },
{ "_id": "bob", "total": 17, "views": 93200 },
{ "_id": "carol", "total": 33, "views": 71800 },
{ "_id": "dave", "total": 8, "views": 42100 },
{ "_id": "eve", "total": 21, "views": 38900 }
]对应 SQL:
SELECT author_id, COUNT(*) AS total, SUM(views) AS views
FROM posts WHERE status = 'published'
GROUP BY author_id ORDER BY views DESC LIMIT 5;2. 表达式语法
进管道之前必须先掌握表达式语法,它和查询语法不一样。
2.1 字段引用要加美元前缀
"$authorId" // 引用字段 authorId 的值
"$stats.views" // 引用内嵌字段
"$$ROOT" // 引用整个当前文档
"$$NOW" // 当前时间(4.2+)
"authorId" // 没有美元符号 = 字符串字面量 "authorId"这是新手最常见的困惑来源:忘了写美元符号,结果分组键变成了一个固定字符串,所有文档被分到同一组。
2.2 常用表达式操作符
// 算术
{ $add: ["$a", "$b"] }
{ $subtract: ["$end", "$start"] }
{ $multiply: ["$price", "$qty"] }
{ $divide: ["$total", "$count"] }
{ $round: ["$score", 2] }
// 字符串
{ $concat: ["$firstName", " ", "$lastName"] }
{ $toUpper: "$code" }
{ $substrCP: ["$title", 0, 20] }
{ $split: ["$path", "/"] }
{ $strLenCP: "$body" }
// 条件
{ $cond: { if: { $gte: ["$score", 60] }, then: "pass", else: "fail" } }
{ $ifNull: ["$nickname", "$username"] }
{ $switch: {
branches: [
{ case: { $gte: ["$views", 10000] }, then: "hot" },
{ case: { $gte: ["$views", 1000] }, then: "warm" }
],
default: "cold"
} }
// 数组
{ $size: "$tags" }
{ $arrayElemAt: ["$tags", 0] }
{ $filter: { input: "$comments", as: "c", cond: { $gt: ["$$c.likes", 3] } } }
{ $map: { input: "$scores", as: "s", in: { $multiply: ["$$s", 10] } } }
{ $in: ["mongodb", "$tags"] }
// 日期
{ $year: "$createdAt" }
{ $dateToString: { date: "$createdAt", format: "%Y-%m-%d", timezone: "Asia/Shanghai" } }
{ $dateTrunc: { date: "$createdAt", unit: "day", timezone: "Asia/Shanghai" } }
{ $dateDiff: { startDate: "$createdAt", endDate: "$$NOW", unit: "day" } }
// 类型转换
{ $toInt: "$viewsStr" }
{ $toDate: "$tsMillis" }
{ $toObjectId: "$idStr" }
{ $convert: { input: "$v", to: "int", onError: 0, onNull: 0 } }$$c、$$s、$$ROOT 里的双美元符号表示「系统或用户定义的变量」,用于 $filter、$map、$let 这类需要引入临时变量的表达式里。单美元符号引用文档字段,双美元符号引用变量,记住这个区分。
3. 核心阶段详解
3.1 $match
{ $match: { status: "published", "stats.views": { $gte: 100 } } }$match 用的是查询语法(和 find 完全一致),不是表达式语法。它是唯一能利用索引的过滤阶段,所以位置至关重要(见第 5 节)。
要在 $match 里用表达式,需要套 $expr:
{ $match: { $expr: { $gt: ["$stats.likes", "$stats.views"] } } }3.2 $project
db.posts.aggregate([
{ $project: {
_id: 0,
title: 1,
author: "$authorId", // 重命名
views: "$stats.views", // 提升内嵌字段
titleLen: { $strLenCP: "$title" }, // 计算新字段
isHot: { $gte: ["$stats.views", 1000] }, // 布尔计算
day: { $dateToString: { date: "$createdAt", format: "%Y-%m-%d" } }
} }
])[
{ "title": "MongoDB 入门", "author": "alice", "views": 1200,
"titleLen": 11, "isHot": true, "day": "2024-06-01" }
]$project 是白名单模式:没列出来的字段会被丢弃。这经常不是你想要的,所以更常用的是 $set / $addFields:
// $set 和 $addFields 完全等价,保留所有原字段,只新增/覆盖指定字段
{ $set: { isHot: { $gte: ["$stats.views", 1000] } } }
// $unset 删字段
{ $unset: ["internalNote", "draftBody"] }$project 需要你把想保留的字段全部列出来,一旦文档结构变化就容易漏。$set 只声明「增加什么」,语义更清晰、更不易出错。只在管道末尾整理输出格式时才用 $project。
3.3 $group
{ $group: {
_id: "$authorId", // 分组键,null 表示全部聚成一组
count: { $sum: 1 },
totalViews: { $sum: "$stats.views" },
avgViews: { $avg: "$stats.views" },
maxViews: { $max: "$stats.views" },
firstPost: { $first: "$title" }, // 需要先 $sort 才有意义
allTags: { $push: "$tags" },
uniqueTags: { $addToSet: "$tags" }
} }_id 可以是复合的:
// 按「作者 + 月份」分组
{ $group: {
_id: {
author: "$authorId",
month: { $dateToString: { date: "$createdAt", format: "%Y-%m" } }
},
count: { $sum: 1 }
} }[
{ "_id": { "author": "alice", "month": "2024-05" }, "count": 8 },
{ "_id": { "author": "alice", "month": "2024-06" }, "count": 12 },
{ "_id": { "author": "bob", "month": "2024-06" }, "count": 5 }
]累加器一览:
| 累加器 | 说明 | SQL 对应 |
|---|---|---|
$sum | 求和;$sum: 1 就是计数 | SUM / COUNT |
$avg | 平均值,自动忽略非数字 | AVG |
$min / $max | 最小/最大 | MIN / MAX |
$first / $last | 组内第一/最后一个 | 需要窗口函数 |
$push | 收集成数组 | GROUP_CONCAT |
$addToSet | 收集去重数组 | DISTINCT 聚合 |
$count | 计数(独立阶段) | COUNT(*) |
$stdDevPop / $stdDevSamp | 标准差 | STDDEV |
$top / $bottom | 组内 Top N(5.2+) | 窗口函数 |
$mergeObjects | 合并文档 | — |
$group 的输出顺序是不确定的,不要假设它按分组键有序。需要有序结果必须在 $group 之后显式 $sort。同理,$first / $last 只有在 $group 之前做过 $sort 时才有确定含义。
3.4 $sort / $limit / $skip / $count
{ $sort: { totalViews: -1, _id: 1 } }
{ $limit: 10 }
{ $skip: 20 }
{ $count: "total" } // 输出形如一个只有 total 字段的文档$count 是 $group 加 $project 的语法糖:
// 这两个等价
{ $count: "total" }
// 等于
{ $group: { _id: null, total: { $sum: 1 } } }
{ $project: { _id: 0 } }4. 六个实战示例
基于 community 库。
4.1 每个城市的用户数与平均年龄
db.users.aggregate([
{ $group: {
_id: "$profile.city",
userCount: { $sum: 1 },
avgAge: { $avg: "$age" }
} },
{ $set: { avgAge: { $round: ["$avgAge", 1] } } },
{ $sort: { userCount: -1 } }
])[
{ "_id": "Beijing", "userCount": 1240, "avgAge": 29.4 },
{ "_id": "Shanghai", "userCount": 980, "avgAge": 31.2 },
{ "_id": "Shenzhen", "userCount": 651, "avgAge": 27.8 }
]4.2 标签热度排行(数组展开统计)
db.posts.aggregate([
{ $match: { status: "published" } },
{ $unwind: "$tags" },
{ $group: { _id: "$tags", posts: { $sum: 1 }, views: { $sum: "$stats.views" } } },
{ $sort: { posts: -1 } },
{ $limit: 10 }
])[
{ "_id": "mongodb", "posts": 320, "views": 892000 },
{ "_id": "index", "posts": 145, "views": 331000 },
{ "_id": "redis", "posts": 98, "views": 210000 }
]$unwind 把一个含 N 个元素的数组文档拆成 N 个文档,每个只带一个元素。下一章详讲。
4.3 每日发帖趋势
db.posts.aggregate([
{ $match: { createdAt: { $gte: new Date("2024-06-01") } } },
{ $group: {
_id: { $dateToString: { date: "$createdAt", format: "%Y-%m-%d", timezone: "Asia/Shanghai" } },
count: { $sum: 1 }
} },
{ $sort: { _id: 1 } }
])[
{ "_id": "2024-06-01", "count": 42 },
{ "_id": "2024-06-02", "count": 38 },
{ "_id": "2024-06-03", "count": 55 }
]4.4 分档统计(用 $switch 做区间归类)
db.posts.aggregate([
{ $set: {
level: { $switch: {
branches: [
{ case: { $gte: ["$stats.views", 10000] }, then: "S" },
{ case: { $gte: ["$stats.views", 1000] }, then: "A" },
{ case: { $gte: ["$stats.views", 100] }, then: "B" }
],
default: "C"
} }
} },
{ $group: { _id: "$level", count: { $sum: 1 } } },
{ $sort: { _id: 1 } }
])[
{ "_id": "A", "count": 4200 },
{ "_id": "B", "count": 31500 },
{ "_id": "C", "count": 60100 },
{ "_id": "S", "count": 380 }
]4.5 每个作者的最热门帖子
db.posts.aggregate([
{ $match: { status: "published" } },
{ $sort: { authorId: 1, "stats.views": -1 } },
{ $group: {
_id: "$authorId",
topTitle: { $first: "$title" },
topViews: { $first: "$stats.views" },
totalPosts: { $sum: 1 }
} },
{ $sort: { topViews: -1 } },
{ $limit: 5 }
])先按「作者 + 浏览量倒序」排序,再 $group 取 $first,这是「分组取 Top 1」的经典写法。5.2+ 也可以直接用 $top:
{ $group: {
_id: "$authorId",
best: { $top: { output: ["$title", "$stats.views"], sortBy: { "stats.views": -1 } } }
} }4.6 去重计数
// 有多少个不同的作者发过帖
db.posts.aggregate([
{ $group: { _id: "$authorId" } },
{ $count: "distinctAuthors" }
])[ { "distinctAuthors": 12483 } ]5. 管道优化(重点)
5.1 $match 越早越好
// 慢:先分组 100 万条,再过滤
db.posts.aggregate([
{ $group: { _id: "$authorId", views: { $sum: "$stats.views" } } },
{ $match: { _id: "alice" } }
])
// 快:先过滤到 42 条,再分组
db.posts.aggregate([
{ $match: { authorId: "alice" } },
{ $group: { _id: "$authorId", views: { $sum: "$stats.views" } } }
])只有位于管道最前面的 $match 能利用索引。一旦经过 $group / $unwind / $project 等改变文档结构的阶段,后续的 $match 只能在内存里逐条过滤。
5.2 优化器会做的自动改写
MongoDB 的聚合优化器会尝试一些重排,可以用 explain 看到最终执行的管道:
db.posts.explain().aggregate([ ... ])常见的自动优化:
| 优化 | 说明 |
|---|---|
$match 前移 | 把 $match 挪到 $project / $unwind / $sort 之前(前提是引用的字段没被改动) |
$sort + $limit 合并 | 变成 Top-K 排序,只保留 K 个文档在内存里 |
$limit 前移与合并 | 连续的 $limit 取最小值 |
$skip + $limit 交换 | 合并成一次 |
$project 下推 | 提前裁剪字段,减少后续阶段的数据量 |
优化器的重排有严格前提:如果 $project 重命名了 $match 要用的字段,它就不敢前移。养成手动把 $match 写在第一位的习惯,比祈祷优化器聪明可靠得多。
5.3 100MB 内存限制
每个阻塞阶段($group、$sort、$bucket 等需要看到全部输入才能输出的阶段)有 100MB 内存上限:
MongoServerError: PlanExecutor error during aggregation ::
Exceeded memory limit for $group, but didn't allow external sort处理方式:
// 允许落盘(慢,但不会失败)
db.posts.aggregate([ ... ], { allowDiskUse: true })从 6.0 起 allowDiskUse 默认为 true,但仍然要尽量避免落盘——磁盘排序比内存慢一到两个数量级。真正的解法是把 $match 前移、用索引支持 $sort、或者减少 $group 的分组数。
按 userId 分组 1000 万用户,会在内存里维护 1000 万个累加器,必然超限。这类需求应该改成:一、按时间分批跑;二、用 $merge 增量写入结果集合;三、干脆用预聚合模式在写入时就维护好统计(第 12 章)。
5.4 让 $sort 走索引
db.posts.createIndex({ status: 1, createdAt: -1 })
db.posts.aggregate([
{ $match: { status: "published" } },
{ $sort: { createdAt: -1 } }, // 走索引,零内存开销
{ $limit: 20 }
])规则和 find 完全一样:$sort 只有紧跟在能用索引的 $match 之后(中间没有改变文档的阶段),才可能利用索引。
6. aggregate 的其他选项
db.posts.aggregate(pipeline, {
allowDiskUse: true,
maxTimeMS: 30000, // 超时保护,强烈建议对线上聚合设置
hint: { status: 1, createdAt: -1 },
comment: "daily-report", // 会出现在慢日志和 currentOp 里,便于定位
collation: { locale: "zh", numericOrdering: true }
})一、写一个聚合,统计每个作者的帖子数、总浏览数、平均点赞数,按总浏览数倒序取前 10;二、统计 posts 中每个标签的使用次数(提示:需要 $unwind);三、统计每个月每个城市的新增用户数(需要 $lookup 吗?想想为什么不需要);四、故意把 $match 写在 $group 之后,用 explain 对比两种写法的执行时间;五、写一个聚合把 posts 按浏览数分成 S/A/B/C 四档并统计各档数量,然后加上每档的平均点赞数;六、用 $group 加 $count 统计有多少个不同的标签。
小结
- 聚合是阶段数组,文档像流水一样依次穿过,形态可被彻底改造
- 表达式语法里字段要加单美元前缀,变量用双美元前缀,别和查询语法混淆
$match用查询语法且是唯一能走索引的过滤阶段,必须放在最前面$set/$addFields比$project更安全,后者是白名单会丢字段$group输出无序,$first/$last依赖前置$sort- 阻塞阶段有 100MB 内存限制,分组基数过大是最常见的超限原因
- 用
maxTimeMS给线上聚合加超时保护 - 下一章讲
$lookup、$facet、窗口函数等进阶阶段 →