Skip to content

MongoDB 聚合管道

MongoDB 聚合管道用于对文档进行过滤、分组、字段转换、排序、关联和统计。它像一条流水线:每个阶段接收上一阶段输出的数据,再交给下一阶段。

如果普通 find 解决的是“查哪些文档”,聚合管道解决的就是“查出来以后如何加工和统计”。

学习目标

学完本页,你要能做到:

目标能力
看懂管道模型知道文档如何一阶段一阶段流动
会写常见统计能写 $match$project$group$sort$unwind$lookup
会优化性能知道为什么要先过滤、先裁剪字段,为什么大排序和大分组危险
会处理商业报表能做资产统计、采集日志统计、部门排行、异常设备 TopN
知道边界能判断实时聚合、预聚合、数仓、ES 分别适合什么
会面试回答能讲清 $lookup 风险、内存限制、allowDiskUse、深分页和预聚合

基本流程

mermaid
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展开数组行展开

示例数据

javascript
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"] }
])

统计每个用户的已支付金额

javascript
db.orders.aggregate([
  { $match: { status: "PAID" } },
  {
    $group: {
      _id: "$userId",
      totalAmount: { $sum: "$amount" },
      orderCount: { $sum: 1 }
    }
  },
  { $sort: { totalAmount: -1 } }
])

逐步解释:

  1. $match 只保留已支付订单。
  2. $groupuserId 分组。
  3. $sum: "$amount" 统计金额。
  4. $sum: 1 统计订单数量。
  5. $sort 按总金额倒序排列。

聚合管道怎么执行

聚合不是把所有阶段一次性完成,而是上一阶段输出文档流,下一阶段继续处理。

mermaid
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 字段转换

javascript
db.orders.aggregate([
  {
    $project: {
      userId: 1,
      amount: 1,
      discountedAmount: { $multiply: ["$amount", 0.9] },
      _id: 0
    }
  }
])

$project 可以选择字段,也可以生成新字段。

商业场景:医疗资产按科室统计

需求:统计某医院已启用资产在各科室的数量,并按数量倒序显示 Top 20。

javascript
db.assets.aggregate([
  {
    $match: {
      hospitalId: "H001",
      status: "USED",
      deleted: false
    }
  },
  {
    $group: {
      _id: "$departmentId",
      total: { $sum: 1 }
    }
  },
  { $sort: { total: -1 } },
  { $limit: 20 }
])

推荐索引:

javascript
db.assets.createIndex({
  hospitalId: 1,
  status: 1,
  deleted: 1,
  departmentId: 1
})

为什么这样写:

  1. $match 放最前面,让索引先过滤某医院、状态、未删除数据。
  2. $group 只处理过滤后的文档,不扫全库。
  3. $sort 排的是分组后的科室结果,而不是全量资产文档。
  4. $limit 控制返回 TopN,避免接口返回过大。

如果这个统计每个页面都高频访问,并且数据量很大,更推荐做预聚合,把“每次实时扫全量统计”变成“写入或定时任务维护统计表”。

unwind 展开数组

javascript
db.orders.aggregate([
  { $unwind: "$items" },
  {
    $group: {
      _id: "$items",
      count: { $sum: 1 }
    }
  }
])

原来一个订单有多个商品,$unwind 会把数组拆成多条记录,方便统计每个商品出现次数。

unwind 的风险

$unwind 会把一条包含数组的文档展开成多条文档。如果数组平均长度是 20,10 万条文档展开后就是 200 万条中间文档。

mermaid
flowchart TD
    A["10万条订单"] --> B["每单20个items"]
    B --> C["$unwind 后 200万条中间文档"]
    C --> D["再 $group / $sort"]
    D --> E["CPU和内存压力放大"]

所以 $unwind 前要尽量 $match 缩小范围,并用 $project 只保留必要字段。

错误写法:

javascript
db.orders.aggregate([
  { $unwind: "$items" },
  { $match: { createdAt: { $gte: ISODate("2026-07-01") } } },
  { $group: { _id: "$items.skuId", total: { $sum: 1 } } }
])

更合理:

javascript
db.orders.aggregate([
  { $match: { createdAt: { $gte: ISODate("2026-07-01") } } },
  { $project: { items: 1 } },
  { $unwind: "$items" },
  { $group: { _id: "$items.skuId", total: { $sum: 1 } } }
])

lookup 关联集合

假设有用户集合:

javascript
db.users.insertMany([
  { _id: "u1", name: "Tom" },
  { _id: "u2", name: "Alice" }
])

关联订单和用户:

javascript
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 往往说明建模可能不合适。

mermaid
flowchart TD
    A["主集合过滤后文档"] --> B["对每条文档查关联集合"]
    B --> C{"关联字段是否有索引"}
    C -- "有" --> D["按索引查关联文档"]
    C -- "无" --> E["关联集合扫描,成本放大"]
    D --> F["合并结果"]
    E --> F

商业例子:资产列表要展示科室名称。如果列表 QPS 很高,科室名称变化又很少,可以考虑把 departmentName 冗余到资产文档。冗余不是没有代价,科室改名时要通过事件、后台任务或读时修正同步字段。

$lookup 优化建议:

建议原因
$lookup 前先 $match减少主集合参与关联的文档数
foreignField 建索引避免关联集合扫描
控制返回字段不要把关联文档整块带出来
高频简单展示考虑冗余减少实时关联成本
复杂报表走预聚合/数仓在线库不适合承担所有分析

性能原则

mermaid
flowchart TD
    A["聚合需求"] --> B["尽早 $match"]
    B --> C["只保留需要的字段"]
    C --> D["减少 $group 和 $sort 的数据量"]
    D --> E["必要时建立索引"]
    E --> F["使用 explain 分析"]

建议:

  1. $match 尽量放前面,减少后续阶段的数据量。
  2. 能命中索引的过滤条件要放在管道前部。
  3. 大集合排序要有索引或限制数据范围。
  4. $group$sort 可能消耗较多内存。
  5. 报表类复杂统计可以考虑离线预聚合。

常见问题

现象可能原因建议
聚合很慢前面没有过滤,大量数据进入管道先加 $match
内存不足大量 $group$sort缩小范围或做预聚合
关联很慢$lookup 数据量太大建索引或减少关联
结果字段太多没有 $project只返回必要字段
分页越往后越慢大量 $skip改游标分页

allowDiskUse 是不是万能解

聚合中 $group$sort 可能需要较多内存。allowDiskUse: true 允许部分阶段把临时数据写到磁盘。

javascript
db.asset_collect_logs.aggregate([
  { $match: { hospitalId: "H001" } },
  { $group: { _id: "$assetNo", total: { $sum: 1 } } },
  { $sort: { total: -1 } }
], { allowDiskUse: true })

它能避免部分内存限制错误,但不是性能优化本身:

理解说明
能让大聚合继续跑但可能变慢,因为写磁盘比内存慢
不能替代索引$match$sort 仍要尽量利用索引
不能替代预聚合高频报表每次跑全量仍会压垮在线库
需要监控磁盘大量临时文件会影响整体性能

如果一个接口必须长期依赖 allowDiskUse 才能返回,通常要重新评估:是否应该缩小范围、异步导出、预聚合或进数仓。

聚合 explain 怎么看

聚合也可以 explain:

javascript
db.assets.explain("executionStats").aggregate([
  { $match: { hospitalId: "H001", status: "USED" } },
  { $group: { _id: "$departmentId", total: { $sum: 1 } } }
])

重点:

指标说明
是否有 COLLSCAN前置 $match 是否没用上索引
totalDocsExamined扫描文档数是否远大于业务预期
nReturned返回结果数
是否有大范围 SORT排序是否无法利用索引或数据量过大
每个 stage 的输入输出哪一段把数据量放大或消耗最多

排查聚合慢,先问三个问题:

  1. $match 是否足够靠前并命中索引?
  2. $group/$sort/$lookup 前的数据量有多大?
  3. 这个统计是否应该实时算,还是应该预聚合?

练习

  1. 统计每个订单状态的订单数量。
  2. 统计每个用户的订单总金额和订单数。
  3. 使用 $unwind 统计商品出现次数。
  4. 使用 $project 生成折扣后金额。
  5. 思考一个复杂报表是否应该实时聚合,还是提前生成统计表。

小结

聚合管道是 MongoDB 的数据加工流水线。学习时要掌握 $match$project$group$sort$lookup$unwind。写聚合时要先减少数据量,再做分组和排序,并用索引和 explain 验证性能。