MongoDB 聚合管道深度解析:从分组到多面搜索

深入 MongoDB 聚合管道,覆盖核心阶段、$lookup 联表、$facet 多面搜索、性能优化

MongoDB 的聚合管道(Aggregation Pipeline)是其数据处理能力的精华所在。如果说 CRUD 操作是对单个文档的增删改查,那么聚合管道就是将整个集合作为输入,经过一系列有序的数据转换阶段(Stage),最终输出经过加工、重组和计算的汇总结果。它相当于关系型数据库中 SQL 的强大补充(GROUP BY、JOIN、子查询、窗口函数的综合替代),但拥有更直观的管道式语法和远超传统 SQL 的灵活性。在生产环境中,聚合管道承担了报表统计、实时分析、数据清洗、ETL 转换等核心场景。本文将从执行原理出发,深入解析每个核心阶段的工作机制,通过电商订单分析的完整实战案例,揭示如何写出高效、可维护的聚合查询。

聚合管道核心概念与执行流程

管道的基本结构

聚合管道的核心思想是管道(Pipeline)——数据像水流一样,依次通过一个个处理阶段(Stage),每个阶段的输出作为下一个阶段的输入。整个管道以数组形式定义,数组中的每个元素都是一个阶段对象,对象键为阶段操作符,值为该阶段的配置参数。

// 管道基本骨架:一个由多个阶段组成的数组
db.collection.aggregate([
  { $stage1: { /* 阶段1的配置 */ } },
  { $stage2: { /* 阶段2的配置 */ } },
  { $stage3: { /* 阶段3的配置 */ } }
])

以统计各品类商品平均价格为例,一个典型的三阶段管道如下:

// 统计每个商品品类的平均价格、商品数量和最高价格
db.products.aggregate([
  // 阶段1:过滤——只分析上架状态的商品
  { $match: { status: "active" } },
  // 阶段2:分组——按品类汇总统计
  { $group: {
      _id: "$category",
      avgPrice: { $avg: "$price" },
      count: { $sum: 1 },
      maxPrice: { $max: "$price" }
  }},
  // 阶段3:排序——按商品数量降序排列
  { $sort: { count: -1 } }
])

管道的执行顺序严格遵循数组定义的顺序。数据不会跳过任何阶段,也不会回头执行。这种线性的处理模型使得开发者可以像搭积木一样,逐步构建复杂的数据处理逻辑,每个阶段只做一件事,组合起来就能完成复杂的分析任务。

执行引擎与优化器简介

MongoDB 的聚合管道由专门的聚合框架(Aggregation Framework)执行,它内置了一个查询优化器(Query Optimizer),会根据数据分布、索引情况和阶段特性自动调整执行计划。优化器的主要能力包括:

  • 阶段重组(Stage Reordering):将 $match 阶段尽可能提前到管道前端,利用索引过滤数据,减少后续阶段需要处理的文档数量。
  • 阶段合并(Stage Coalescence):将相邻的 $sort + $limit 合并为 Top-K 操作,避免全量排序;将连续的 $project 合并为单次投影。
  • 索引利用(Index Intersection):当 $match 包含复合索引的前缀条件时,优先使用索引扫描(IXSCAN)而非全集合扫描(COLLSCAN)。
// 优化器会自动将早期管道中的 $match 下推
// 原始管道:$group 在前,$match 在后
db.orders.aggregate([
  { $group: { _id: "$status", total: { $sum: "$amount" } } },
  { $match: { "_id": "completed" } }
])

// 优化后等价于:先 match 再 group(如果集合有 status 索引的话)
db.orders.aggregate([
  { $match: { status: "completed" } },
  { $group: { _id: "$status", total: { $sum: "$amount" } } }
])

需要注意的是,优化器并非万能。对于 $group 后的 $match(即对分组结果进行过滤,对应 SQL 的 HAVING),优化器无法将其提前到 $group 之前。此外,涉及 $lookup$unwind 等关联操作的阶段,顺序变更会影响语义正确性,优化器不会随意重排。

管道运算符分类体系

聚合管道的运算符分为两大类:阶段操作符(Stage Operators)表达式操作符(Expression Operators)

阶段操作符 出现在管道的顶层,定义数据处理的主要阶段。常见的阶段操作符包括:

阶段操作符作用SQL 类比
$match过滤文档WHERE
$group按指定键分组聚合GROUP BY
$sort对文档排序ORDER BY
$project选择字段、计算新字段SELECT
$lookup左外连接其他集合LEFT JOIN
$unwind展开数组字段为多个文档LATERAL UNNEST
$addFields添加计算字段(保留原有字段)SELECT *, expr AS new_col
$facet在同一输入上并行执行多个子管道多面搜索
$limit限制输出文档数量LIMIT
$skip跳过指定数量文档OFFSET
$count统计文档数并输出单一值COUNT(*)
$replaceRoot将子文档提升为根文档SELECT subdoc.*

表达式操作符 用于在字段级别进行计算和转换,通常嵌套在 $project$group$addFields 等阶段的值中。例如 $avg$sum$concat$dateToString 等。

// 表达式操作符嵌套在阶段中
db.orders.aggregate([
  { $project: {
      year: { $year: "$createdAt" },           // 日期提取表达式
      quarter: { $ceil: { $divide: [{ $month: "$createdAt" }, 3] } },  // 数学表达式
      displayAmount: { $concat: ["¥", { $toString: "$amount" }] }      // 字符串表达式
  }}
])

理解这两类操作符的区别至关重要:阶段操作符改变的是文档的集合形态(过滤、分组、排序、连接),表达式操作符改变的是单个文档内部的字段值(计算、转换、提取)。

核心阶段深度解析

一、$match:精准过滤,管道优化第一站

$match 阶段用于过滤文档,其语法与 find() 的查询条件完全一致,支持所有查询操作符($eq$gt$in$regex$exists 等)。作为管道中最常见的第一阶段,它的核心价值在于利用索引减少后续阶段的输入量。

// 基础过滤:查找 2024 年内的已完成订单
db.orders.aggregate([
  { $match: {
      status: "completed",
      createdAt: { $gte: new Date("2024-01-01"), $lt: new Date("2025-01-01") }
  }}
])

$match 放在管道前端的优势是显著的。假设 orders 集合有 1000 万条文档,其中 2024 年 completed 状态的订单只有 50 万条。如果先 $match,后续所有阶段只需处理 50 万文档;如果先其他阶段再 $match,优化器不保证能将其提前,可能导致全量处理。

// 反例:$group 放在前面,$match 放在后面过滤分组结果
db.orders.aggregate([
  // 先对所有数据分组(全量扫描 + 全量分组)
  { $group: { _id: "$status", totalAmount: { $sum: "$amount" } } },
  // 再过滤分组结果(此时 $match 不能提前,因为需要先 group 才有 _id)
  { $match: { _id: "completed" } }
])

// 正例:先过滤再分组(索引友好)
db.orders.aggregate([
  { $match: { status: "completed" } },
  { $group: { _id: "$status", totalAmount: { $sum: "$amount" } } }
])

$match 中支持使用查询计划相关的操作符:

// 结合 $expr 使用聚合表达式进行过滤
db.orders.aggregate([
  { $match: {
      $expr: { $gt: ["$amount", "$discountThreshold"] }
  }}
])

// $expr 支持在 $match 中引用聚合表达式变量
// 注意:使用 $expr 时,普通索引无法被利用,可能退化为全集合扫描
db.orders.aggregate([
  { $match: {
      $and: [
        { status: "completed" },           // 这部分可走索引
        { $expr: { $gt: ["$amount", 1000] } }  // 这部分不能走索引
      ]
  }}
])

在生产环境中,如果 $match 需要使用 $expr 且查询频繁,可考虑预先计算相关字段并建立索引,或使用 $addFields 计算后再 $match(但这会导致全量扫描)。

二、$group:数据分组的灵魂

$group 是聚合管道中最核心的阶段之一,它按指定的 _id 表达式将文档分组,并对每组应用聚合函数。$group 的工作方式与 SQL 的 GROUP BY 类似,但功能更灵活。

// 基础分组:按支付方式统计订单数量和总金额
db.orders.aggregate([
  { $match: { status: "completed" } },
  { $group: {
      _id: "$paymentMethod",
      orderCount: { $sum: 1 },                    // 每组文档数
      totalAmount: { $sum: "$amount" },          // 金额汇总
      avgAmount: { $avg: "$amount" },            // 平均金额
      minAmount: { $min: "$amount" },            // 最小金额
      maxAmount: { $max: "$amount" }             // 最大金额
  }}
])

$group 阶段支持的全部累加器(Accumulator)函数:

累加器功能适用类型
$sum求和(参数为 1 时计数)数值
$avg求平均值数值
$min / $max最小值 / 最大值可比较类型
$first / $last组内第一条/最后一条文档的字段值任意
$push将字段值收集到数组任意
$addToSet将字段值收集到数组(去重)任意
$stdDevPop / $stdDevSamp总体/样本标准差数值
$mergeObjects合并组内文档为单个对象对象
$top / $bottom按排序取 top/bottom 文档的字段值任意(5.2+)
$topN / $bottomN按排序取 top/bottom N 个值任意(5.2+)

多字段分组(复合键分组)能提供更细粒度的分析维度:

// 按年和月份分组统计(复合键分组)
db.orders.aggregate([
  { $match: { status: "completed" } },
  { $group: {
      _id: {
        year: { $year: "$createdAt" },
        month: { $month: "$createdAt" }
      },
      totalAmount: { $sum: "$amount" },
      orderCount: { $sum: 1 }
  }},
  { $sort: { "_id.year": -1, "_id.month": -1 } }
])

条件聚合——$sum 配合 $cond:在 $group 中,累加器的参数可以是任意表达式,这使得条件统计成为可能。

// 统计各品类订单:总数、成功数、退款数、成功金额、退款金额
db.orders.aggregate([
  { $group: {
      _id: "$category",
      totalOrders: { $sum: 1 },
      successOrders: { $sum: { $cond: [{ $eq: ["$status", "completed"] }, 1, 0] } },
      refundOrders: { $sum: { $cond: [{ $eq: ["$status", "refunded"] }, 1, 0] } },
      successAmount: { $sum: { $cond: [{ $eq: ["$status", "completed"] }, "$amount", 0] } },
      refundAmount: { $sum: { $cond: [{ $eq: ["$status", "refunded"] }, "$amount", 0] } }
  }}
])

$cond 是三元条件表达式,格式为 { $cond: { if: <boolean-expression>, then: <true-value>, else: <false-value> } },也可简写为数组形式 { $cond: [ <boolean>, <true>, <false> ] }。与 $cond 类似的条件表达式还有 $switch(多分支选择),适合复杂条件场景:

// 使用 $switch 进行多分支评级
db.orders.aggregate([
  { $project: {
      amount: 1,
      grade: { $switch: {
        branches: [
          { case: { $gte: ["$amount", 10000] }, then: "VIP" },
          { case: { $gte: ["$amount", 5000] }, then: "高级" },
          { case: { $gte: ["$amount", 1000] }, then: "普通" }
        ],
        default: "小额"
      }}
  }}
])

三、$sort / $project / $limit / $skip:管道的基础变换

这四个阶段是管道中最基础但最常用的数据变换工具。

$sort 对输入文档按指定字段排序,可以一次指定多个字段。排序方向 1 为升序,-1 为降序。当排序字段上有匹配的索引时,MongoDB 可以避免内存排序,直接从索引获取有序数据。

// 多字段排序:先按状态升序,再按创建时间降序
db.orders.aggregate([
  { $sort: { status: 1, createdAt: -1 } }
])

// 排序配合分页:获取第 3 页,每页 20 条
db.orders.aggregate([
  { $sort: { createdAt: -1 } },
  { $skip: 40 },
  { $limit: 20 }
])

排序的内存消耗是一个需要注意的问题。MongoDB 默认对排序操作施加 100 MB 的内存限制(由 internalQueryExecMaxBlockingSortBytes 控制)。当排序数据量超过此限制时,必须使用 allowDiskUse: true 选项,允许 MongoDB 将排序溢出数据写入磁盘。磁盘排序的性能远低于内存排序,应尽量避免。

// 大数据量排序时启用磁盘溢出
db.orders.aggregate(
  [
    { $sort: { createdAt: -1 } }
  ],
  { allowDiskUse: true }
)

$project 用于控制输出文档的字段结构,可以包含、排除、重命名、计算字段。投影阶段的灵活性远超 SQL 的 SELECT

// 包含特定字段,排除 _id
db.orders.aggregate([
  { $project: {
      _id: 0,                   // 0 表示排除
      orderNo: 1,               // 1 表示包含(原样输出)
      amount: 1,
      year: { $year: "$createdAt" },    // 计算新字段
      month: { $month: "$createdAt" },
      amountWithTax: { $multiply: ["$amount", 1.06] }
  }}
])

$project 中有一个重要的限制:不能在同一层级混合使用包含(1)和排除(0)模式,唯一的例外是 _id 字段可以同时排除。如果要排除少数字段而保留大多数,使用 $addFields 配合 $project 的排除模式,或者直接使用 $unset 阶段(MongoDB 3.6+):

// $unset 专门用于排除字段,比 $project 更直观
db.orders.aggregate([
  { $unset: ["internalNotes", "rawPayload", "ipAddress"] }
])

// 等价写法(使用 $project 的排除模式)
db.orders.aggregate([
  { $project: { internalNotes: 0, rawPayload: 0, ipAddress: 0 } }
])

$limit$skip 是分页的基础工具。但要注意,$skip 在大偏移量时性能急剧下降,因为它需要扫描并丢弃大量文档。对于深度分页,推荐使用游标或基于排序键的范围查询替代 $skip

// 优化深度分页:使用游标记(cursor-based pagination)
// 第一页
db.orders.aggregate([
  { $match: { status: "completed" } },
  { $sort: { createdAt: -1 } },
  { $limit: 20 }
])
// 拿到最后一条的 createdAt,作为下一页的 "cursor"

// 第二页(用 cursor 替代 skip)
db.orders.aggregate([
  { $match: {
      status: "completed",
      createdAt: { $lt: lastPageCursor }   // 只查找比 cursor 更早的
  }},
  { $sort: { createdAt: -1 } },
  { $limit: 20 }
])

四、$lookup:联表查询的桥梁

$lookup 实现了 MongoDB 跨集合的连接查询,相当于 SQL 的 LEFT OUTER JOIN。在文档模型的设计中,虽然推荐在适用场景下使用嵌入文档减少关联需求,但在多对多关系、频繁变更的关联数据、需要独立查询的实体等场景中,$lookup 是不可替代的。

MongoDB 3.2 引入了基础版 $lookup,5.1 引入了等值连接改进版。最常用的语法如下:

// $lookup 基础语法(等值连接)
db.orders.aggregate([
  { $lookup: {
      from: "users",              // 被连接的集合名
      localField: "userId",       // 本集合的关联字段
      foreignField: "_id",        // 被连接集合的关联字段
      as: "userInfo"              // 输出数组字段名
  }}
])

$lookup 的关键特性是总是返回数组。即使匹配结果只有一条文档,也会包装成单元素数组 [{...}];如果没有匹配,则返回空数组 []。这一特性决定了 $lookup 后通常需要配合 $unwind 将数组展开为文档。

// $lookup 后展开数组
db.orders.aggregate([
  { $match: { status: "completed" } },
  { $lookup: {
      from: "users",
      localField: "userId",
      foreignField: "_id",
      as: "userInfo"
  }},
  // 将 userInfo 数组展开(preserveNullAndEmptyArrays 保留无匹配的订单)
  { $unwind: { path: "$userInfo", preserveNullAndEmptyArrays: true } },
  // 现在可以像访问普通字段一样访问 userInfo 的子字段
  { $project: {
      orderNo: 1,
      amount: 1,
      userName: "$userInfo.name",
      userEmail: "$userInfo.email",
      userLevel: "$userInfo.vipLevel"
  }}
])

$unwind 的使用需要谨慎:

  • preserveNullAndEmptyArrays: true 保留无匹配或数组为空的文档,相当于 LEFT JOIN
  • 省略此选项(默认 false)时,无匹配或空数组的文档会被过滤掉,相当于 INNER JOIN
  • 如果源字段不是数组而是 null 或缺失字段,$unwind 的行为取决于 preserveNullAndEmptyArrays 设置。

管道式 $lookup(MongoDB 3.6+)允许在 $lookup 内部定义子管道(let 变量绑定 + pipeline 自定义),实现更复杂的连接逻辑:

// 管道式 $lookup:连接时附加筛选和计算
db.orders.aggregate([
  { $match: { status: "completed" } },
  { $lookup: {
      from: "users",
      // let 定义本阶段传入子管道的变量
      let: { orderUserId: "$userId", orderTime: "$createdAt" },
      pipeline: [
        // 在子管道中使用 $$ 引用 let 变量
        { $match: {
            $expr: { $eq: ["$_id", "$$orderUserId"] }
        }},
        // 在子管道中执行额外的聚合操作
        { $project: {
            name: 1,
            email: 1,
            vipLevel: 1,
            // 判断该用户是否在订单创建时已经是 VIP
            wasVipAtOrderTime: {
              $cond: {
                if: { $gte: ["$vipSince", "$$orderTime"] },
                then: true,
                else: false
              }
            }
        }}
      ],
      as: "userInfo"
  }},
  { $unwind: { path: "$userInfo", preserveNullAndEmptyArrays: true } }
])

管道式 $lookup 的强大之处在于子管道中可以包含任意聚合阶段:额外的 $match 过滤、 $project 投影、 $group 分组等。这使得 $lookup 不再局限于简单的等值连接,可以实现条件连接、聚合连接甚至半连接(Semi-Join)。

// 半连接:只返回有至少一条关联评论的文章
db.articles.aggregate([
  { $lookup: {
      from: "comments",
      let: { articleId: "$_id" },
      pipeline: [
        { $match: { $expr: { $eq: ["$articleId", "$$articleId"] } } },
        { $limit: 1 }  // 只需要知道存在即可
      ],
      as: "hasComments"
  }},
  { $match: { "hasComments.0": { $exists: true } } },
  { $unset: "hasComments" }
])

$lookup 的性能需要特别注意:被连接的集合(from)会在 foreignField 上执行等值查询。如果该字段没有索引,每次 $lookup 都会触发被连接集合的全集合扫描,性能随数据量线性下降。在生产环境中,务必确保 foreignField 上有索引

// 为 users._id 建立索引(通常 _id 默认已有索引,这里仅作演示)
db.users.createIndex({ _id: 1 })

// 为 orders.userId 建立索引(用于从 orders 方向查询时)
db.orders.createIndex({ userId: 1 })

五、$addFields:增强文档而非替换

$addFields(MongoDB 3.4+)与 $project 类似,都是添加计算字段,但 $addFields保留原文档的所有字段,仅添加新字段或覆盖已有字段。这比 $project 更适合"在现有文档基础上做增强"的场景。

// $addFields 保留所有原字段,只添加/覆盖指定字段
db.orders.aggregate([
  { $match: { status: "completed" } },
  { $addFields: {
      // 添加新字段
      discountRate: { $cond: {
        if: { $gte: ["$amount", 5000] },
        then: 0.15,
        else: { $cond: { if: { $gte: ["$amount", 1000] }, then: 0.10, else: 0.05 } }
      }},
      finalAmount: {
        $multiply: [
          "$amount",
          { $subtract: [1, { $cond: {
            if: { $gte: ["$amount", 5000] },
            then: 0.15,
            else: { $cond: { if: { $gte: ["$amount", 1000] }, then: 0.10, else: 0.05 } }
          }}]
        }]
      },
      // 覆盖已有字段(转换格式)
      createdAt: { $dateToString: { format: "%Y-%m-%d %H:%M:%S", date: "$createdAt" } }
  }}
])

相比 $project$addFields 更适合增量式添加字段的场景,尤其是管道中已经有很多字段需要保留时。$addFields 也可以接受文档字面量,将子文档整体添加到根级别。

六、$facet:多面搜索的利器

$facet(MongoDB 3.4+)是聚合管道中最具特色的阶段之一。它允许在同一组输入文档上并行执行多个独立的子管道(Sub-pipeline),每个子管道的输出作为结果对象的一个字段。这在统计仪表板、商品筛选页、搜索结果页等多维度数据展示场景中极为实用。

一个典型的电商商品筛选页面需要同时返回:分类统计、价格区间分布、品牌列表、分页商品结果。传统做法需要发送多次查询,而 $facet 可以在单次聚合中完成所有计算。

// 电商商品筛选页的多面搜索
db.products.aggregate([
  // 先根据用户的筛选条件过滤
  { $match: {
      category: "手机",
      price: { $gte: 2000, $lte: 8000 },
      status: "active"
  }},
  // 在同一批商品上并行计算多个维度
  { $facet: {
      // 子管道1:按品牌统计数量
      brandStats: [
        { $group: { _id: "$brand", count: { $sum: 1 } } },
        { $sort: { count: -1 } }
      ],
      // 子管道2:按价格区间统计
      priceRanges: [
        { $bucket: {
            groupBy: "$price",
            boundaries: [2000, 3000, 4000, 5000, 6000, 7000, 8000],
            default: "8000+",
            output: { count: { $sum: 1 } }
        }}
      ],
      // 子管道3:分页商品列表
      products: [
        { $sort: { salesCount: -1 } },
        { $skip: 0 },
        { $limit: 20 },
        { $project: { name: 1, brand: 1, price: 1, rating: 1, imageUrl: 1 } }
      ],
      // 子管道4:总数量(用于分页)
      totalCount: [
        { $count: "count" }
      ],
      // 子管道5:可用筛选项(按属性统计)
      attributeFilters: [
        { $unwind: "$attributes" },
        { $group: {
            _id: "$attributes.key",
            values: { $addToSet: "$attributes.value" }
        }}
      ]
  }}
])

$facet 的输出是一个单一文档,每个子管道的结果作为文档的一个字段:

// $facet 输出示例
{
  brandStats: [
    { _id: "Apple", count: 15 },
    { _id: "Samsung", count: 12 },
    { _id: "Xiaomi", count: 8 }
  ],
  priceRanges: [
    { _id: 2000, count: 5 },
    { _id: 3000, count: 12 },
    { _id: 4000, count: 8 }
  ],
  products: [
    { _id: ObjectId("..."), name: "iPhone 15", brand: "Apple", price: 5999, ... },
    // ... 共 20 条
  ],
  totalCount: [{ count: 100 }],
  attributeFilters: [
    { _id: "color", values: ["黑色", "白色", "蓝色"] },
    { _id: "storage", values: ["128GB", "256GB", "512GB"] }
  ]
}

$facet 的限制和注意事项:

  • 每个子管道接收相同的输入(即 $facet 前一阶段的输出)。
  • 子管道内部不能使用 $facet(不能嵌套),也不能使用 $out$merge
  • 子管道之间没有依赖关系,是真正并行执行的。
  • $facet 的结果总是单一文档,如果原始输入为空,它会返回包含空数组字段的文档,而非空结果集。

$bucket 和 $bucketAuto$facet 的常用搭档,用于数据分桶统计:

// $bucket 手动指定边界
db.orders.aggregate([
  { $bucket: {
      groupBy: "$amount",
      boundaries: [0, 100, 500, 1000, 5000, 10000],
      default: "10000+",                           // 超出最大边界的归入 default
      output: {
        count: { $sum: 1 },
        totalAmount: { $sum: "$amount" },
        avgAmount: { $avg: "$amount" }
      }
  }}
])

// $bucketAuto 自动等频分桶(分成指定数量的桶,每桶文档数大致相等)
db.orders.aggregate([
  { $bucketAuto: {
      groupBy: "$amount",
      buckets: 5,                                 // 分成 5 个桶
      output: {
        count: { $sum: 1 },
        avgAmount: { $avg: "$amount" }
      }
  }}
])

管道性能优化:从索引到内存

优化原则一:尽早过滤

在管道前端尽可能多用 $match 缩小数据范围,这是性能提升的第一要点。理想情况下,$match 应当利用索引进行过滤,将后续阶段需要处理的文档量降到最低。

// 差:先做全量分组和排序,最后才过滤
db.orders.aggregate([
  { $group: { _id: "$userId", total: { $sum: "$amount" } } },
  { $sort: { total: -1 } },
  { $match: { total: { $gte: 10000 } } }    // 在末尾过滤
])

// 好:如果目的是找大额用户,先过滤订单再分组更高效
db.orders.aggregate([
  { $match: { amount: { $gte: 1000 } } },   // 先过滤(可走索引)
  { $group: { _id: "$userId", total: { $sum: "$amount" } } },
  { $match: { total: { $gte: 10000 } } },   // HAVING 过滤无法提前
  { $sort: { total: -1 } }
])

优化原则二:索引与排序的协同

如果 $sort 字段与 $match 的过滤字段可以共享同一个复合索引,MongoDB 可能直接利用索引获取已排序的数据,避免内存排序。

// 创建复合索引(status 用于过滤,createdAt 用于排序)
db.orders.createIndex({ status: 1, createdAt: -1 })

// 此管道可能直接使用索引顺序输出排序结果
db.orders.aggregate([
  { $match: { status: "completed" } },
  { $sort: { createdAt: -1 } },    // 如果优化器选择此索引,可避免内存排序
  { $limit: 100 }
])

优化原则三:限制内存使用

聚合管道的每个阶段默认有 100 MB 的内存使用限制。当 $group$sort$lookup 等阶段处理的数据量超过此限制时,管道会报错终止,除非显式启用 allowDiskUse

// 启用磁盘溢出(允许排序/group 使用磁盘临时文件)
db.orders.aggregate(
  [
    { $group: { _id: "$category", total: { $sum: "$amount" } } },
    { $sort: { total: -1 } }
  ],
  { allowDiskUse: true }
)

allowDiskUse 虽然解决了内存溢出的问题,但磁盘 I/O 会严重影响性能。更根本的优化策略是:

  1. $match 阶段充分利用索引,减少进入后续阶段的文档数。
  2. 使用 $limit 尽早截断数据流。如果只需要前 N 条结果,在 $sort 之后立即加 $limit
  3. 减少 $project 阶段的字段数量。传输和处理的字段越少,内存占用越低。
  4. 对于大集合的分组统计,考虑使用预聚合。维护一个统计集合,在数据写入时增量更新聚合结果,查询时直接读取预计算数据。

优化原则四:$lookup 的索引策略

$lookup 的性能瓶颈主要在被连接集合上。每次 $lookup 都会向 from 集合发起一次查询,如果 foreignField 没有索引,被连接集合会成为性能瓶颈。

// 确保 users._id 有索引(_id 默认有唯一索引)
// 如果 foreignField 不是 _id,必须手动创建索引
db.products.createIndex({ categoryId: 1 })   // 为 $lookup 的 foreignField 建索引

此外,如果在 $lookup 之前先用 $match 大幅缩小输入集,$lookup 需要执行的次数也会减少,这是优化 $lookup 管道的关键。

优化原则五:避免不必要的 $unwind

$unwind 将一条文档展开为多条文档,会使后续阶段的输入文档数量膨胀。如果仅需要关联数据的某些字段,考虑在 $lookup 的子管道中使用 $project 精简输出,而不是先展开再投影。

// 低效:展开后再投影
db.orders.aggregate([
  { $lookup: { from: "users", localField: "userId", foreignField: "_id", as: "userInfo" } },
  { $unwind: "$userInfo" },
  { $project: { orderNo: 1, userName: "$userInfo.name" } }
])

// 稍好:在 $lookup 子管道中直接投影,减少传输的数据量
db.orders.aggregate([
  { $lookup: {
      from: "users",
      let: { uid: "$userId" },
      pipeline: [
        { $match: { $expr: { $eq: ["$_id", "$$uid"] } } },
        { $project: { _id: 0, name: 1 } }
      ],
      as: "userInfo"
  }},
  { $unwind: { path: "$userInfo", preserveNullAndEmptyArrays: true } },
  { $project: { orderNo: 1, userName: "$userInfo.name" } }
])

实战案例:电商订单分析系统

以一个中等规模的电商订单分析需求为例,综合运用上述所有核心阶段,构建完整的聚合管道。

数据模型

// orders 集合(模拟千万级订单)
{
  _id: ObjectId("..."),
  orderNo: "ORD-2026-0123456",
  userId: ObjectId("..."),
  status: "completed",              // completed / paid / shipped / refunded / cancelled
  items: [
    { sku: "SKU-001", name: "机械键盘", category: "数码外设", price: 498, quantity: 1 },
    { sku: "SKU-002", name: "鼠标垫", category: "数码外设", price: 59, quantity: 2 }
  ],
  amount: 616,
  paymentMethod: "alipay",         // alipay / wechat / card / cod
  region: "north",
  createdAt: ISODate("2026-08-13T08:30:00Z"),
  deliveredAt: ISODate("2026-08-15T14:20:00Z")
}

// users 集合
db.users.insertMany([
  { _id: ObjectId("64a1b2c3d4e5f6a7b8c9d0e1"), name: "张三", email: "zhangsan@example.com", vipLevel: 2, registerDate: ISODate("2024-03-15") },
  { _id: ObjectId("64a1b2c3d4e5f6a7b8c9d0e2"), name: "李四", email: "lisi@example.com", vipLevel: 1, registerDate: ISODate("2025-01-20") }
])

需求一:月度销售报表(按品类 + 支付方式多维度)

db.orders.aggregate([
  // 阶段1:过滤有效订单
  { $match: {
      status: "completed",
      createdAt: { $gte: new Date("2026-01-01"), $lt: new Date("2027-01-01") }
  }},
  // 阶段2:展开商品明细(一条订单可能含多个品类商品)
  { $unwind: "$items" },
  // 阶段3:计算每条商品明细的单品金额
  { $addFields: {
      itemTotal: { $multiply: ["$items.price", "$items.quantity"] }
  }},
  // 阶段4:按年-月-品类分组统计
  { $group: {
      _id: {
        year: { $year: "$createdAt" },
        month: { $month: "$createdAt" },
        category: "$items.category"
      },
      totalRevenue: { $sum: "$itemTotal" },
      totalQuantity: { $sum: "$items.quantity" },
      orderCount: { $addToSet: "$_id" },   // 去重统计涉及订单数
      avgItemPrice: { $avg: "$items.price" }
  }},
  // 阶段5:计算真实订单数($addToSet 的结果是数组,取长度)
  { $addFields: {
      uniqueOrderCount: { $size: "$orderCount" }
  }},
  // 阶段6:格式化输出
  { $project: {
      _id: 0,
      year: "$_id.year",
      month: "$_id.month",
      category: "$_id.category",
      totalRevenue: { $round: ["$totalRevenue", 2] },
      totalQuantity: 1,
      uniqueOrderCount: 1,
      avgItemPrice: { $round: ["$avgItemPrice", 2] }
  }},
  // 阶段7:排序输出
  { $sort: { year: -1, month: -1, totalRevenue: -1 } }
])

需求二:用户消费行为分析(含 $lookup 关联)

db.orders.aggregate([
  // 阶段1:只分析已完成订单,且排除测试数据
  { $match: {
      status: "completed",
      createdAt: { $gte: new Date("2025-01-01") },
      amount: { $gt: 0 }
  }},
  // 阶段2:计算每个用户的消费统计
  { $group: {
      _id: "$userId",
      totalSpent: { $sum: "$amount" },
      orderCount: { $sum: 1 },
      avgOrderValue: { $avg: "$amount" },
      firstOrderDate: { $min: "$createdAt" },
      lastOrderDate: { $max: "$createdAt" },
      favoritePayment: { $first: "$paymentMethod" },    // 简化处理,实际需要更复杂的众数计算
      categories: { $addToSet: "$items.category" }       // 收集所有购买过的品类
  }},
  // 阶段3:关联用户表获取用户信息
  { $lookup: {
      from: "users",
      localField: "_id",
      foreignField: "_id",
      as: "user"
  }},
  { $unwind: { path: "$user", preserveNullAndEmptyArrays: true } },
  // 阶段4:计算用户生命周期价值相关指标
  { $addFields: {
      customerLifetimeDays: {
        $ceil: {
          $divide: [
            { $subtract: ["$lastOrderDate", "$firstOrderDate"] },
            1000 * 60 * 60 * 24
          ]
        }
      },
      customerSegment: {
        $switch: {
          branches: [
            { case: { $gte: ["$totalSpent", 50000] }, then: "钻石" },
            { case: { $gte: ["$totalSpent", 20000] }, then: "白金" },
            { case: { $gte: ["$totalSpent", 5000] }, then: "黄金" }
          ],
          default: "普通"
        }
      }
  }},
  // 阶段5:格式化输出
  { $project: {
      _id: 0,
      userId: "$_id",
      userName: "$user.name",
      userEmail: "$user.email",
      vipLevel: "$user.vipLevel",
      customerSegment: 1,
      totalSpent: { $round: ["$totalSpent", 2] },
      orderCount: 1,
      avgOrderValue: { $round: ["$avgOrderValue", 2] },
      customerLifetimeDays: 1,
      categoryCount: { $size: "$categories" }
  }},
  // 阶段6:按消费额排序
  { $sort: { totalSpent: -1 } },
  { $limit: 100 }
])

需求三:多面搜索仪表盘——$facet 综合应用

db.orders.aggregate([
  // 阶段1:先按时间范围过滤(所有子面共享的过滤条件)
  { $match: {
      status: "completed",
      createdAt: { $gte: new Date("2026-07-01"), $lt: new Date("2026-08-01") }
  }},
  // 阶段2:并行计算所有统计面
  { $facet: {
      // 面1:每日销售趋势
      dailyTrend: [
        { $group: {
            _id: { $dateToString: { format: "%Y-%m-%d", date: "$createdAt" } },
            revenue: { $sum: "$amount" },
            orderCount: { $sum: 1 }
        }},
        { $sort: { _id: 1 } }
      ],
      // 面2:地区销售分布
      regionDistribution: [
        { $group: {
            _id: "$region",
            revenue: { $sum: "$amount" },
            orderCount: { $sum: 1 }
        }},
        { $sort: { revenue: -1 } }
      ],
      // 面3:支付方式占比
      paymentStats: [
        { $group: {
            _id: "$paymentMethod",
            orderCount: { $sum: 1 },
            totalAmount: { $sum: "$amount" }
        }},
        { $sort: { orderCount: -1 } }
      ],
      // 面4:Top 10 热销商品(需要展开 items)
      topProducts: [
        { $unwind: "$items" },
        { $group: {
            _id: "$items.sku",
            name: { $first: "$items.name" },
            category: { $first: "$items.category" },
            totalSold: { $sum: "$items.quantity" },
            totalRevenue: { $sum: { $multiply: ["$items.price", "$items.quantity"] } }
        }},
        { $sort: { totalSold: -1 } },
        { $limit: 10 }
      ],
      // 面5:汇总统计(单值)
      summary: [
        { $group: {
            _id: null,
            totalRevenue: { $sum: "$amount" },
            totalOrders: { $sum: 1 },
            avgOrderValue: { $avg: "$amount" }
        }},
        { $project: {
            _id: 0,
            totalRevenue: { $round: ["$totalRevenue", 2] },
            totalOrders: 1,
            avgOrderValue: { $round: ["$avgOrderValue", 2] }
        }}
      ]
  }}
])

需求四:配送时效分析——日期运算与条件统计

db.orders.aggregate([
  { $match: {
      status: { $in: ["completed", "shipped"] },
      deliveredAt: { $exists: true }
  }},
  // 计算配送天数(ceil 向上取整)
  { $addFields: {
      deliveryDays: {
        $ceil: {
          $divide: [
            { $subtract: ["$deliveredAt", "$createdAt"] },
            1000 * 60 * 60 * 24
          ]
        }
      }
  }},
  // 按地区和配送时效分组统计
  { $group: {
      _id: "$region",
      avgDeliveryDays: { $avg: "$deliveryDays" },
      minDays: { $min: "$deliveryDays" },
      maxDays: { $max: "$deliveryDays" },
      totalOrders: { $sum: 1 },
      // 统计各时效区间的订单数
      within3Days: { $sum: { $cond: [{ $lte: ["$deliveryDays", 3] }, 1, 0] } },
      within7Days: { $sum: { $cond: [{ $lte: ["$deliveryDays", 7] }, 1, 0] } },
      over14Days: { $sum: { $cond: [{ $gt: ["$deliveryDays", 14] }, 1, 0] } }
  }},
  // 计算百分比
  { $addFields: {
      pctWithin3Days: { $round: [{ $multiply: [{ $divide: ["$within3Days", "$totalOrders"] }, 100] }, 2] },
      pctWithin7Days: { $round: [{ $multiply: [{ $divide: ["$within7Days", "$totalOrders"] }, 100] }, 2] }
  }},
  { $sort: { avgDeliveryDays: 1 } },
  { $project: {
      _id: 0,
      region: "$_id",
      avgDeliveryDays: { $round: ["$avgDeliveryDays", 1] },
      totalOrders: 1,
      pctWithin3Days: 1,
      pctWithin7Days: 1,
      over14Days: 1,
      minDays: 1,
      maxDays: 1
  }}
])

SQL 与 MongoDB Aggregation 对照速查

对于从 SQL 背景转向 MongoDB 的开发者,以下对照表能够加速认知转换:

SQL 概念MongoDB Aggregation说明
WHERE$match过滤文档
GROUP BY$group按指定键分组
SELECT col1, col2$project / $addFields选择或计算字段
ORDER BY$sort排序文档
LIMIT n$limit限制结果数量
OFFSET n$skip跳过文档
LEFT JOIN$lookup + $unwind左外连接
INNER JOIN$lookup + $unwind(不带 preserveNullAndEmptyArrays内连接
SUM() / AVG() / MIN() / MAX()$sum / $avg / $min / $max聚合函数
COUNT(*){ $sum: 1 }计数
COUNT(DISTINCT){ $addToSet: ... } + $size$group + $sum: 1去重计数
HAVING$match 放在 $group 之后过滤分组结果
CASE WHEN$cond / $switch条件表达式
EXTRACT(YEAR FROM date){ $year: "$date" }日期提取
CONCAT()$concat字符串拼接
SUBSTRING()$substr / $substrBytes / $substrCP子串提取
UNNEST()$unwind展开数组
子查询(Scalar Subquery)$lookup 子管道嵌套查询
WITH (CTE)定义 view 或拆分多个聚合公用表表达式

典型 SQL 到 Aggregation 的转换示例

-- SQL:统计每个品类 2024 年的销售额、订单数、平均客单价
SELECT
    items.category,
    SUM(items.price * items.quantity) AS total_revenue,
    COUNT(DISTINCT o._id) AS order_count,
    AVG(o.amount) AS avg_order_value
FROM orders o
CROSS JOIN UNNEST(o.items) AS items
WHERE o.status = 'completed'
  AND o.created_at >= '2024-01-01'
  AND o.created_at < '2025-01-01'
GROUP BY items.category
ORDER BY total_revenue DESC;
// MongoDB Aggregation 等价实现
db.orders.aggregate([
  { $match: {
      status: "completed",
      createdAt: { $gte: new Date("2024-01-01"), $lt: new Date("2025-01-01") }
  }},
  { $unwind: "$items" },
  { $group: {
      _id: "$items.category",
      totalRevenue: { $sum: { $multiply: ["$items.price", "$items.quantity"] } },
      orderCount: { $sum: 1 },
      avgOrderValue: { $avg: "$amount" }
  }},
  { $project: {
      _id: 0,
      category: "$_id",
      totalRevenue: 1,
      orderCount: 1,
      avgOrderValue: { $round: ["$avgOrderValue", 2] }
  }},
  { $sort: { totalRevenue: -1 } }
])

总结

MongoDB 聚合管道是一个功能全面且表达力极强的数据处理框架。从简单的过滤、排序、分页,到复杂的分组统计、跨集合连接、多面搜索,管道式的数据流模型让复杂的数据分析任务变得层次分明、易于理解和维护。

在实际生产环境中,掌握以下几点能够帮助写出高性能的聚合管道:

管道设计层面,牢记"尽早过滤"原则。把 $match 放在管道前端,尽可能利用索引减少后续阶段的输入文档量。对于跨集合查询,确保 $lookupforeignField 上有索引,避免被连接集合成为性能瓶颈。

阶段选择层面,根据场景选择 $project 还是 $addFields:需要精简字段结构时用 $project,需要在现有文档上增量添加字段时用 $addFields。使用 $facet 实现多面搜索时,确保所有子面共享的前置 $match 已经充分过滤数据。

性能优化层面,关注 100 MB 内存限制。当数据量可能超出限制时,优先在前端过滤缩小数据量,而非直接启用 allowDiskUse。对于深度分页,使用基于游标的范围查询替代 $skip。对于大集合的复杂分组统计,考虑预聚合(定期用聚合管道计算并写入统计集合)或 Change Streams 实时增量更新。

可读性层面,复杂管道建议拆分为多个视图(db.createView()),每个视图封装一个逻辑层。这不仅提升代码可维护性,还能让 MongoDB 的查询优化器有更多机会对每个子视图独立优化。

// 创建视图:将复杂管道封装为可复用的虚拟集合
db.createView("monthly_sales_report", "orders", [
  { $match: { status: "completed" } },
  { $unwind: "$items" },
  { $group: {
      _id: { year: { $year: "$createdAt" }, month: { $month: "$createdAt" }, category: "$items.category" },
      revenue: { $sum: { $multiply: ["$items.price", "$items.quantity"] } },
      count: { $sum: 1 }
  }},
  { $sort: { "_id.year": -1, "_id.month": -1, revenue: -1 } }
])

// 像查询普通集合一样查询视图
db.monthly_sales_report.find({ "_id.year": 2026 })

聚合管道的学习曲线虽然比基础 CRUD 陡峭,但一旦掌握,它将成为 MongoDB 数据分析中最有力的工具。从简单的统计报表到复杂的实时分析仪表盘,聚合管道都能胜任。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「database」更多文章

  1. 缓存架构演进之路:从单机 Redis 到亿级分布式多级缓存体系
  2. Redis 7.x 重大新特性与架构升级深度解析
  3. Redis 消息队列深度对比:Pub/Sub、Streams 与 Kafka/RabbitMQ 选型指南