MongoDBNotes

第 08 章:聚合管道

zjc 于 2026-01-08 发布

这是《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 }
])

性能建议:

  1. $match 放前;
  2. $limit 尽量提前;
  3. 排序字段进索引;
  4. 避免大 $skip
  5. 用游标条件分页;
  6. 只投影必要字段。

8.10 内存限制

某些阶段有内存限制。排序可以显式允许落盘:

db.orders.aggregate(
  [
    { $sort: { created_at: -1 } }
  ],
  { allowDiskUse: true }
)

注意:

  1. 落盘会更慢;
  2. 占用临时空间;
  3. 可能影响其他请求;
  4. 应先优化索引和过滤;
  5. 大分析任务考虑 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 的数据处理流水线。写管道时坚持先过滤、再投影、再分组或关联,排序和分页尽量依靠索引。小型实时统计适合管道,大规模分析应考虑专用分析引擎。

思考题

  1. 为什么 $match 应放在管道前面?
  2. $unwind 为什么会放大数据量?
  3. $lookup 的外键为什么需要索引?
  4. $out$merge 有什么区别?
  5. allowDiskUse 能解决所有性能问题吗?