这是《MongoDB 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 聚合管道把文档流经过一组阶段处理,每个阶段接收上一步输出并生成下一步输入。它适合统计、分组、变形、关联和多集合组装,但大数据量时必须优先匹配和索引,必要时使用允许落盘的排序。
8.1 管道思想
collection
-> $match
-> $project
-> $group
-> $sort
-> $limit
每个阶段都在内存中构造新文档流,因此阶段顺序和过滤位置直接决定性能。
8.2 准备数据
db.orders.insertMany([
{ order_no: "O1", user_id: "u1", status: "PAID", amount: NumberDecimal("200.00"), items: [{ sku: "A", qty: 2 }, { sku: "B", qty: 1 }] },
{ order_no: "O2", user_id: "u1", status: "PAID", amount: NumberDecimal("300.00"), items: [{ sku: "A", qty: 1 }] },
{ order_no: "O3", user_id: "u2", status: "CREATED", amount: NumberDecimal("100.00"), items: [{ sku: "C", qty: 5 }] }
])
8.3 $match 与 $project
先过滤:
db.orders.aggregate([
{ $match: {
status: "PAID",
created_at: { $gte: new Date("2026-08-01T00:00:00Z") }
}}
])
再投影:
db.orders.aggregate([
{ $match: { status: "PAID" } },
{ $project: {
_id: 0,
order_no: 1,
user_id: 1,
amount: 1
}}
])
$match 应尽量放在管道前面,才能使用索引并减少后续数据量。
8.4 $group
按用户统计:
db.orders.aggregate([
{ $match: { status: "PAID" } },
{ $group: {
_id: "$user_id",
order_count: { $sum: 1 },
total_amount: { $sum: "$amount" },
max_amount: { $max: "$amount" },
last_order_no: { $last: "$order_no" }
}},
{ $sort: { total_amount: -1 } },
{ $limit: 10 }
])
常见累加器:
| 累加器 | 说明 |
|---|---|
$sum |
求和 |
$avg |
平均 |
$min / $max |
极值 |
$push |
收集值到数组 |
$addToSet |
收集去重值 |
$first / $last |
首尾值 |
$push 可能形成大数组,要限制分组规模。
8.5 $unwind
展开订单明细:
db.orders.aggregate([
{ $match: { status: "PAID" } },
{ $unwind: "$items" },
{ $group: {
_id: "$items.sku",
quantity: { $sum: "$items.qty" },
order_count: { $sum: 1 }
}},
{ $sort: { quantity: -1 } }
])
空数组处理:
{ $unwind: { path: "$items", preserveNullAndEmptyArrays: true } }
$unwind 会放大文档数量,放在 $match 之后能减少处理量。
8.6 $lookup
关联用户:
db.orders.aggregate([
{ $match: { status: "PAID" } },
{ $lookup: {
from: "users",
localField: "user_id",
foreignField: "_id",
as: "user"
}},
{ $unwind: "$user" },
{ $project: {
order_no: 1,
amount: 1,
user_name: "$user.name",
user_level: "$user.level"
}}
])
关联集合应确保外键有索引。高频 $lookup 通常说明模型冗余或拆分需要调整。
8.7 $facet
一次返回列表和统计:
db.products.aggregate([
{ $match: { status: "ON_SALE" } },
{ $facet: {
page: [
{ $sort: { created_at: -1 } },
{ $skip: 0 },
{ $limit: 20 }
],
summary: [
{ $group: {
_id: null,
total: { $sum: 1 },
avg_price: { $avg: "$price" }
}}
]
}}
])
$facet 会共享前置阶段输出,但仍要控制数据规模。
8.8 窗口函数
按用户计算订单时间序:
db.orders.aggregate([
{ $match: { status: "PAID" } },
{ $sort: { user_id: 1, created_at: 1 } },
{ $setWindowFields: {
partitionBy: "$user_id",
sortBy: { created_at: 1 },
output: {
user_order_seq: { $documentNumber: {} },
rolling_amount: {
$sum: "$amount",
window: { documents: [-2, 0] }
}
}
}}
])
窗口算子适合报表和排名,复杂窗口仍应评估是否由分析引擎承担。
8.9 分页与排序
基础分页:
db.orders.aggregate([
{ $match: { user_id: "u1" } },
{ $sort: { created_at: -1 } },
{ $skip: 0 },
{ $limit: 20 }
])
性能建议:
$match放前;$limit尽量提前;- 排序字段进索引;
- 避免大
$skip; - 用游标条件分页;
- 只投影必要字段。
8.10 内存限制
某些阶段有内存限制。排序可以显式允许落盘:
db.orders.aggregate(
[
{ $sort: { created_at: -1 } }
],
{ allowDiskUse: true }
)
注意:
- 落盘会更慢;
- 占用临时空间;
- 可能影响其他请求;
- 应先优化索引和过滤;
- 大分析任务考虑 ClickHouse 等分析系统。
8.11 输出到集合
保存结果:
db.orders.aggregate([
{ $match: { status: "PAID" } },
{ $group: {
_id: "$user_id",
total_amount: { $sum: "$amount" }
}},
{ $merge: {
into: "user_order_stats",
on: "_id",
whenMatched: "replace",
whenNotMatched: "insert"
}}
])
$out 覆盖输出集合,$merge 支持合并策略。定时任务需要幂等和调度审计。
8.12 常见问题
| 问题 | 排查 |
|---|---|
| 聚合很慢 | $match 太晚、无索引、数据放大 |
| 内存超限 | 分组或排序过大 |
$lookup 慢 |
外键无索引 |
$unwind 后数量暴涨 |
数组过大 |
| 结果金额不准 | 使用 double 而非 Decimal128 |
| 随机不一致 | $first / $last 前未排序 |
本章小结
聚合管道是 MongoDB 的数据处理流水线。写管道时坚持先过滤、再投影、再分组或关联,排序和分页尽量依靠索引。小型实时统计适合管道,大规模分析应考虑专用分析引擎。
思考题
- 为什么
$match应放在管道前面? $unwind为什么会放大数据量?$lookup的外键为什么需要索引?$out与$merge有什么区别?allowDiskUse能解决所有性能问题吗?