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 会严重影响性能。更根本的优化策略是:
- 在
$match阶段充分利用索引,减少进入后续阶段的文档数。 - 使用
$limit尽早截断数据流。如果只需要前 N 条结果,在$sort之后立即加$limit。 - 减少
$project阶段的字段数量。传输和处理的字段越少,内存占用越低。 - 对于大集合的分组统计,考虑使用预聚合。维护一个统计集合,在数据写入时增量更新聚合结果,查询时直接读取预计算数据。
优化原则四:$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 放在管道前端,尽可能利用索引减少后续阶段的输入文档量。对于跨集合查询,确保 $lookup 的 foreignField 上有索引,避免被连接集合成为性能瓶颈。
阶段选择层面,根据场景选择 $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 数据分析中最有力的工具。从简单的统计报表到复杂的实时分析仪表盘,聚合管道都能胜任。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。