MongoDB 聚合管道
MongoDB 聚合管道用于对文档进行过滤、分组、字段转换、排序、关联和统计。它像一条流水线:每个阶段接收上一阶段输出的数据,再交给下一阶段。
如果普通 find 解决的是“查哪些文档”,聚合管道解决的就是“查出来以后如何加工和统计”。
学习目标
学完本页,你要能做到:
| 目标 | 能力 |
|---|---|
| 看懂管道模型 | 知道文档如何一阶段一阶段流动 |
| 会写常见统计 | 能写 $match、$project、$group、$sort、$unwind、$lookup |
| 会优化性能 | 知道为什么要先过滤、先裁剪字段,为什么大排序和大分组危险 |
| 会处理商业报表 | 能做资产统计、采集日志统计、部门排行、异常设备 TopN |
| 知道边界 | 能判断实时聚合、预聚合、数仓、ES 分别适合什么 |
| 会面试回答 | 能讲清 $lookup 风险、内存限制、allowDiskUse、深分页和预聚合 |
基本流程
flowchart TD
A["原始集合"] --> B["$match 过滤"]
B --> C["$project 选择或计算字段"]
C --> D["$group 分组统计"]
D --> E["$sort 排序"]
E --> F["$limit 限制数量"]
F --> G["最终结果"]常见阶段
| 阶段 | 作用 | 类比 SQL |
|---|---|---|
$match | 过滤文档 | where |
$project | 选择字段、计算字段 | select |
$group | 分组统计 | group by |
$sort | 排序 | order by |
$limit | 限制数量 | limit |
$skip | 跳过数量 | offset |
$lookup | 关联其他集合 | join |
$unwind | 展开数组 | 行展开 |
示例数据
db.orders.insertMany([
{ userId: "u1", status: "PAID", amount: 100, items: ["book", "pen"] },
{ userId: "u1", status: "PAID", amount: 50, items: ["book"] },
{ userId: "u2", status: "CANCEL", amount: 80, items: ["bag"] }
])统计每个用户的已支付金额
db.orders.aggregate([
{ $match: { status: "PAID" } },
{
$group: {
_id: "$userId",
totalAmount: { $sum: "$amount" },
orderCount: { $sum: 1 }
}
},
{ $sort: { totalAmount: -1 } }
])逐步解释:
$match只保留已支付订单。$group按userId分组。$sum: "$amount"统计金额。$sum: 1统计订单数量。$sort按总金额倒序排列。
聚合管道怎么执行
聚合不是把所有阶段一次性完成,而是上一阶段输出文档流,下一阶段继续处理。
flowchart TD
A["Collection 原始文档"] --> B["Stage 1: $match"]
B --> C["输出较少文档"]
C --> D["Stage 2: $project"]
D --> E["输出较窄文档"]
E --> F["Stage 3: $group"]
F --> G["输出分组结果"]
G --> H["Stage 4: $sort"]
H --> I["最终结果"]这就解释了为什么阶段顺序很重要:
| 顺序 | 结果 |
|---|---|
先 $match | 后面处理的文档少,CPU、内存、IO 都小 |
先 $project | 后面携带字段少,网络和内存压力小 |
先 $group/$sort 再过滤 | 可能先处理全量数据,最容易慢 |
大量 $lookup 放前面 | 关联规模放大,内存和查询压力高 |
一句话:聚合优化的第一原则是让重操作尽量晚发生,让进入重操作的数据尽量少。
project 字段转换
db.orders.aggregate([
{
$project: {
userId: 1,
amount: 1,
discountedAmount: { $multiply: ["$amount", 0.9] },
_id: 0
}
}
])$project 可以选择字段,也可以生成新字段。
商业场景:医疗资产按科室统计
需求:统计某医院已启用资产在各科室的数量,并按数量倒序显示 Top 20。
db.assets.aggregate([
{
$match: {
hospitalId: "H001",
status: "USED",
deleted: false
}
},
{
$group: {
_id: "$departmentId",
total: { $sum: 1 }
}
},
{ $sort: { total: -1 } },
{ $limit: 20 }
])推荐索引:
db.assets.createIndex({
hospitalId: 1,
status: 1,
deleted: 1,
departmentId: 1
})为什么这样写:
$match放最前面,让索引先过滤某医院、状态、未删除数据。$group只处理过滤后的文档,不扫全库。$sort排的是分组后的科室结果,而不是全量资产文档。$limit控制返回 TopN,避免接口返回过大。
如果这个统计每个页面都高频访问,并且数据量很大,更推荐做预聚合,把“每次实时扫全量统计”变成“写入或定时任务维护统计表”。
unwind 展开数组
db.orders.aggregate([
{ $unwind: "$items" },
{
$group: {
_id: "$items",
count: { $sum: 1 }
}
}
])原来一个订单有多个商品,$unwind 会把数组拆成多条记录,方便统计每个商品出现次数。
unwind 的风险
$unwind 会把一条包含数组的文档展开成多条文档。如果数组平均长度是 20,10 万条文档展开后就是 200 万条中间文档。
flowchart TD
A["10万条订单"] --> B["每单20个items"]
B --> C["$unwind 后 200万条中间文档"]
C --> D["再 $group / $sort"]
D --> E["CPU和内存压力放大"]所以 $unwind 前要尽量 $match 缩小范围,并用 $project 只保留必要字段。
错误写法:
db.orders.aggregate([
{ $unwind: "$items" },
{ $match: { createdAt: { $gte: ISODate("2026-07-01") } } },
{ $group: { _id: "$items.skuId", total: { $sum: 1 } } }
])更合理:
db.orders.aggregate([
{ $match: { createdAt: { $gte: ISODate("2026-07-01") } } },
{ $project: { items: 1 } },
{ $unwind: "$items" },
{ $group: { _id: "$items.skuId", total: { $sum: 1 } } }
])lookup 关联集合
假设有用户集合:
db.users.insertMany([
{ _id: "u1", name: "Tom" },
{ _id: "u2", name: "Alice" }
])关联订单和用户:
db.orders.aggregate([
{
$lookup: {
from: "users",
localField: "userId",
foreignField: "_id",
as: "user"
}
},
{ $unwind: "$user" },
{
$project: {
userName: "$user.name",
amount: 1,
status: 1
}
}
])$lookup 很方便,但不要滥用。大量复杂关联可能更适合在业务层、搜索引擎或数仓中处理。
lookup 为什么要谨慎
$lookup 类似 Join,但 MongoDB 的文档模型本来就鼓励把经常一起读、数量有限的数据嵌入或冗余。高频复杂 $lookup 往往说明建模可能不合适。
flowchart TD
A["主集合过滤后文档"] --> B["对每条文档查关联集合"]
B --> C{"关联字段是否有索引"}
C -- "有" --> D["按索引查关联文档"]
C -- "无" --> E["关联集合扫描,成本放大"]
D --> F["合并结果"]
E --> F商业例子:资产列表要展示科室名称。如果列表 QPS 很高,科室名称变化又很少,可以考虑把 departmentName 冗余到资产文档。冗余不是没有代价,科室改名时要通过事件、后台任务或读时修正同步字段。
$lookup 优化建议:
| 建议 | 原因 |
|---|---|
$lookup 前先 $match | 减少主集合参与关联的文档数 |
| foreignField 建索引 | 避免关联集合扫描 |
| 控制返回字段 | 不要把关联文档整块带出来 |
| 高频简单展示考虑冗余 | 减少实时关联成本 |
| 复杂报表走预聚合/数仓 | 在线库不适合承担所有分析 |
性能原则
flowchart TD
A["聚合需求"] --> B["尽早 $match"]
B --> C["只保留需要的字段"]
C --> D["减少 $group 和 $sort 的数据量"]
D --> E["必要时建立索引"]
E --> F["使用 explain 分析"]建议:
$match尽量放前面,减少后续阶段的数据量。- 能命中索引的过滤条件要放在管道前部。
- 大集合排序要有索引或限制数据范围。
$group和$sort可能消耗较多内存。- 报表类复杂统计可以考虑离线预聚合。
常见问题
| 现象 | 可能原因 | 建议 |
|---|---|---|
| 聚合很慢 | 前面没有过滤,大量数据进入管道 | 先加 $match |
| 内存不足 | 大量 $group 或 $sort | 缩小范围或做预聚合 |
| 关联很慢 | $lookup 数据量太大 | 建索引或减少关联 |
| 结果字段太多 | 没有 $project | 只返回必要字段 |
| 分页越往后越慢 | 大量 $skip | 改游标分页 |
allowDiskUse 是不是万能解
聚合中 $group、$sort 可能需要较多内存。allowDiskUse: true 允许部分阶段把临时数据写到磁盘。
db.asset_collect_logs.aggregate([
{ $match: { hospitalId: "H001" } },
{ $group: { _id: "$assetNo", total: { $sum: 1 } } },
{ $sort: { total: -1 } }
], { allowDiskUse: true })它能避免部分内存限制错误,但不是性能优化本身:
| 理解 | 说明 |
|---|---|
| 能让大聚合继续跑 | 但可能变慢,因为写磁盘比内存慢 |
| 不能替代索引 | $match 和 $sort 仍要尽量利用索引 |
| 不能替代预聚合 | 高频报表每次跑全量仍会压垮在线库 |
| 需要监控磁盘 | 大量临时文件会影响整体性能 |
如果一个接口必须长期依赖 allowDiskUse 才能返回,通常要重新评估:是否应该缩小范围、异步导出、预聚合或进数仓。
聚合 explain 怎么看
聚合也可以 explain:
db.assets.explain("executionStats").aggregate([
{ $match: { hospitalId: "H001", status: "USED" } },
{ $group: { _id: "$departmentId", total: { $sum: 1 } } }
])重点:
| 指标 | 说明 |
|---|---|
是否有 COLLSCAN | 前置 $match 是否没用上索引 |
totalDocsExamined | 扫描文档数是否远大于业务预期 |
nReturned | 返回结果数 |
是否有大范围 SORT | 排序是否无法利用索引或数据量过大 |
| 每个 stage 的输入输出 | 哪一段把数据量放大或消耗最多 |
排查聚合慢,先问三个问题:
$match是否足够靠前并命中索引?$group/$sort/$lookup前的数据量有多大?- 这个统计是否应该实时算,还是应该预聚合?
练习
- 统计每个订单状态的订单数量。
- 统计每个用户的订单总金额和订单数。
- 使用
$unwind统计商品出现次数。 - 使用
$project生成折扣后金额。 - 思考一个复杂报表是否应该实时聚合,还是提前生成统计表。
小结
聚合管道是 MongoDB 的数据加工流水线。学习时要掌握 $match、$project、$group、$sort、$lookup、$unwind。写聚合时要先减少数据量,再做分组和排序,并用索引和 explain 验证性能。
