MongoDB聚合管道(Aggregation Pipeline)是处理复杂数据查询和分析的核心工具。相比简单的find查询,聚合管道通过多个阶段(stage)串联处理数据流,支持分组统计、数据转换、关联查询和窗口函数等操作。数据库运维中,聚合管道是生成报表和数据分析的主要手段。
聚合管道基础语法与执行流程
聚合管道将文档依次通过多个处理阶段,每个阶段接收上一阶段的输出作为输入,最终产生结果集。常用阶段包括$match过滤、$group分组、$project投影、$sort排序和$limit限制。
// 聚合管道基础示例
// 需求:统计每个分类下文章数量和平均阅读量,按数量降序排列
db.articles.aggregate([
// 第1阶段: 过滤已发布的文章
{ $match: {
status: "published",
publishedAt: { $gte: ISODate("2026-01-01") }
}},
// 第2阶段: 按分类分组统计
{ $group: {
_id: "$categoryId",
count: { $sum: 1 },
avgViews: { $avg: "$views" },
maxViews: { $max: "$views" },
minViews: { $min: "$views" },
totalViews: { $sum: "$views" }
}},
// 第3阶段: 重命名字段
{ $project: {
_id: 0,
categoryId: "$_id",
articleCount: "$count",
avgViews: { $round: ["$avgViews", 2] },
maxViews: 1,
minViews: 1,
totalViews: 1
}},
// 第4阶段: 按文章数量降序
{ $sort: { articleCount: -1 }},
// 第5阶段: 限制返回10条
{ $limit: 10 }
])
// 等价SQL:
// SELECT categoryId,
// COUNT(*) as articleCount,
// ROUND(AVG(views), 2) as avgViews,
// MAX(views) as maxViews,
// MIN(views) as minViews,
// SUM(views) as totalViews
// FROM articles
// WHERE status = 'published' AND publishedAt >= '2026-01-01'
// GROUP BY categoryId
// ORDER BY articleCount DESC
// LIMIT 10
$group分组操作符与统计函数
$group阶段是聚合管道中最常用的分组阶段,支持多种累加操作符:
// 分组统计操作符大全
db.orders.aggregate([
{ $group: {
// 分组键(null表示全表统计)
_id: "$customerId",
// 计数操作符
totalCount: { $sum: 1 }, // 总数
distinctCount: { $addToSet: "$productId" }, // 去重集合
// 求和与平均
totalAmount: { $sum: "$amount" }, // 求和
avgAmount: { $avg: "$amount" }, // 平均值
// 极值
maxAmount: { $max: "$amount" }, // 最大值
minAmount: { $min: "$amount" }, // 最小值
// 数组操作符
allProducts: { $push: "$productId" }, // 收集到数组(保留重复)
uniqueProducts: { $addToSet: "$productId" }, // 收集到数组(去重)
firstOrder: { $first: "$orderDate" }, // 分组内第一条
lastOrder: { $last: "$orderDate" }, // 分组内最后一条
// 合并文档
orderList: { $mergeObjects: "$$ROOT" } // 合并文档
}}
])
// 多字段分组
db.orders.aggregate([
{ $group: {
_id: {
customerId: "$customerId",
year: { $year: "$orderDate" },
month: { $month: "$orderDate" }
},
monthlyTotal: { $sum: "$amount" },
orderCount: { $sum: 1 }
}}
])
$lookup关联查询与数据合并
$lookup阶段实现类似SQL的LEFT JOIN操作,支持两个集合之间的关联查询:
// 基本关联查询
db.orders.aggregate([
{ $lookup: {
from: "customers", // 关联的集合
localField: "customerId", // 本地字段
foreignField: "_id", // 外键字段
as: "customer" // 输出字段名
}},
// customer字段是数组,取第一个元素
{ $unwind: "$customer" },
{ $project: {
orderId: "$_id",
customerName: "$customer.name",
customerEmail: "$customer.email",
amount: 1
}}
])
// 多条件关联查询(MongoDB 3.6+)
db.orders.aggregate([
{ $lookup: {
from: "products",
let: {
pid: "$productId",
orderQty: "$quantity"
},
pipeline: [
{ $match: {
$expr: {
$and: [
{ $eq: ["$_id", "$$pid"] },
{ $gte: ["$stock", "$$orderQty"] }
]
}
}},
{ $project: { name: 1, price: 1, stock: 1 }}
],
as: "productInfo"
}}
])
$unwind数组展开与$facet并行管道
$unwind将数组字段展开为多条文档,每个数组元素对应一条文档。$facet允许在同一阶段并行执行多个聚合管道:
// $unwind数组展开
// 原始文档: { _id: 1, tags: ["A", "B", "C"] }
// 展开后: 3条文档,分别对应tag A, B, C
db.articles.aggregate([
{ $match: { status: "published" }},
{ $unwind: "$tags" },
{ $group: {
_id: "$tags",
count: { $sum: 1 }
}},
{ $sort: { count: -1 }}
])
// 统计每个标签被多少篇文章使用
// $facet并行管道(一次查询返回多种统计结果)
db.orders.aggregate([
{ $facet: {
// 管道1: 按状态统计
"byStatus": [
{ $group: {
_id: "$status",
count: { $sum: 1 },
totalAmount: { $sum: "$amount" }
}}
],
// 管道2: 按月份统计趋势
"monthlyTrend": [
{ $group: {
_id: {
year: { $year: "$orderDate" },
month: { $month: "$orderDate" }
},
count: { $sum: 1 },
total: { $sum: "$amount" }
}},
{ $sort: { "_id.year": 1, "_id.month": 1 }}
],
// 管道3: TOP 5 大额订单
"topOrders": [
{ $sort: { amount: -1 }},
{ $limit: 5 },
{ $project: {
orderId: "$_id",
amount: 1,
customerId: 1
}}
],
// 管道4: 总计
"summary": [
{ $group: {
_id: null,
totalOrders: { $sum: 1 },
totalAmount: { $sum: "$amount" },
avgAmount: { $avg: "$amount" }
}}
]
}}
])
// 一次查询返回4种统计结果,减少数据库往返
窗口函数与$bucket分桶统计
MongoDB 5.0引入窗口函数,支持在分组内进行排名、累计和移动平均等计算。$bucket阶段可自动将数据分桶:
// 窗口函数示例:计算每篇文章在同分类中的阅读量排名
db.articles.aggregate([
{ $match: { status: "published" }},
{ $setWindowFields: {
// 分区:按分类分组
partitionBy: "$categoryId",
// 排序:按阅读量降序
sortBy: { views: -1 },
// 窗口输出
output: {
rankInCategory: { $rank: {} },
// 累计阅读量
cumulativeViews: {
$sum: "$views",
window: { documents: ["unbounded", "current"] }
},
// 移动平均(前3条到当前)
movingAvg: {
$avg: "$views",
window: { documents: [-3, 0] }
},
// 同分类总文章数
categoryTotal: { $count: {} }
}
}}
])
// $bucket自动分桶统计
db.products.aggregate([
{ $bucket: {
groupBy: "$price",
boundaries: [0, 50, 100, 500, 1000, Infinity],
default: "Other",
output: {
count: { $sum: 1 },
avgPrice: { $avg: "$price" },
products: { $push: "$name" }
}
}}
])
// 输出:
// { _id: 0, count: 15, avgPrice: 25.3 } // 0-50元
// { _id: 50, count: 22, avgPrice: 78.1 } // 50-100元
// { _id: 100, count: 8, avgPrice: 310.5 } // 100-500元
// ...
聚合管道性能优化与索引使用
聚合管道的性能优化关键在于尽早过滤数据和使用合适的索引:
// 优化原则1: $match尽量放在管道最前面
// $match在管道开头时可以使用索引,减少后续阶段处理的数据量
db.orders.aggregate([
{ $match: { status: "completed" }}, // 使用status索引过滤
{ $group: { _id: "$customerId", total: { $sum: "$amount" }}}
])
// 确保 status 字段有索引
// db.orders.createIndex({ status: 1 })
// 优化原则2: $project尽早减少字段
db.orders.aggregate([
{ $match: { status: "completed" }},
{ $project: { customerId: 1, amount: 1 }}, // 只保留需要的字段
{ $group: { _id: "$customerId", total: { $sum: "$amount" }}}
])
// 优化原则3: 复合索引支持$match + $sort
// 如果$match和$sort相邻,可使用复合索引
db.orders.createIndex({ status: 1, orderDate: -1 })
db.orders.aggregate([
{ $match: { status: "completed" }},
{ $sort: { orderDate: -1 }}, // 使用(status, orderDate)复合索引
{ $limit: 100 }
])
// 优化原则4: allowDiskUse处理大数据集
// 聚合管道默认内存限制100MB,超过需要启用磁盘缓存
db.largeCollection.aggregate([
{ $group: { _id: "$field", count: { $sum: 1 }}}
], { allowDiskUse: true })
// 优化原则5: $indexHints强制使用索引
db.orders.aggregate([
{ $indexHint: { status: 1 }},
{ $match: { status: "completed" }},
{ $group: { _id: "$customerId", total: { $sum: "$amount" }}}
])
MongoDB聚合管道的功能覆盖了SQL GROUP BY、JOIN、窗口函数和CASE WHEN等操作。对于NoSQL选型应用场景,聚合管道提供了灵活的数据分析能力。实际使用中需注意管道阶段的执行顺序对性能的影响:$match和$project尽早减少数据量,$group和$lookup等内存密集型操作放在靠后位置。SQL查询优化经验在聚合管道中同样适用——索引覆盖、字段投影和分批处理是三大优化方向。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/mongodb-ju-he-guan-dao-shi-zhan-fen-zu-tong-ji-yu-shu-ju/