35. MongoDB 聚合管道高级算子与窗口函数

聚合管道高级算子:$setWindowFields 窗口函数与排名位移算子、$lookup 的 pipeline 形式与 $graphLookup、$facet 与 $bucket 分面统计、$merge 物化视图,以及 $unwind 性能陷阱。

当 $match、$group、$sort 已经不够用时,聚合管道还有一批高级算子:$setWindowFields 提供类 SQL 的窗口计算,$graphLookup 支持递归图遍历,$facet 一次扫描产出多份统计,$merge 能把结果物化成视图。这些算子的共同特点是能力强但代价不透明,用错地方会让单次聚合吃掉整个副本集的资源。本文逐个拆解它们的语义与性能边界。

1. $setWindowFields 基础

$setWindowFields 是 MongoDB 5.0 引入的窗口计算阶段,等价于 SQL 的 OVER (PARTITION BY ... ORDER BY ...)。它按分区和排序把文档组织成"窗口",在每个窗口内计算聚合值并追加到文档上,而不改变文档数量。

1.1 基本结构

db.sales.aggregate([
  { $setWindowFields: {
      partitionBy: "$region",          // 分区键:按地区分组
      sortBy: { orderDate: 1 },        // 排序键:窗口内按日期升序
      output: {
        runningTotal: {
          $sum: "$amount",
          window: { documents: ["unbounded", "current"] }
        },
        rankInRegion: { $rank: {} },
        prevAmount: { $shift: { output: "$amount", by: -1, default: null } }
      }
  } }
])

输出文档会多出 runningTotal、rankInRegion、prevAmount 三个字段,文档本身与排序都被保留。

1.2 窗口范围

window 定义了参与计算的文档范围,有两种表达:

// 按文档序号:从分区开头到当前文档(累计)
{ window: { documents: ["unbounded", "current"] } }

// 按数值范围:以排序字段的值为基准,前后各一天(含单位)
{ window: { range: [-1, 1], unit: "day" } }

// 按文档偏移:当前文档前 2 条到后 1 条
{ window: { documents: [-2, 1] } }
窗口写法语义边界关键字
documents 区间按文档序号计数unbounded、current
range 区间按排序字段数值数字偏移,可带 unit
省略 window整个分区默认行为

注意:range 窗口要求 sortBy 只有一个字段,且该字段必须是数值或日期类型,否则会报错。多字段排序时只能用 documents 窗口。

1.3 分区与排序的选择

partitionBy 决定窗口的切分粒度。省略 partitionBy 时整个集合是一个分区,适合算全局排名;sortBy 决定窗口内的顺序,是 $rank、$shift、$derivative 等算子的前提。

// 全局排名,不分区
db.scores.aggregate([
  { $setWindowFields: {
      sortBy: { score: -1 },
      output: { globalRank: { $rank: {} } }
  } }
])

2. 窗口算子精讲

$setWindowFields 的 output 支持两类算子:排名类与累积计算类。

2.1 排名算子

db.scores.aggregate([
  { $setWindowFields: {
      partitionBy: "$class",
      sortBy: { score: -1 },
      output: {
        rank:         { $rank: {} },          // 并列同名次,跳号 1,2,2,4
        denseRank:    { $denseRank: {} },     // 并列同名次,不跳号 1,2,2,3
        docNumber:    { $documentNumber: {} } // 分区内行号 1,2,3,4
      }
  } }
])
算子并列处理示例结果
$rank同名次,后续跳号1, 2, 2, 4
$denseRank同名次,不跳号1, 2, 2, 3
$documentNumber不区分并列1, 2, 3, 4

2.2 位移与差分算子

$shift 取窗口内的偏移文档,$derivative 与 $integral 做微积分近似,$expMovingAvg 做指数移动平均。

db.metrics.aggregate([
  { $setWindowFields: {
      partitionBy: "$deviceId",
      sortBy: { ts: 1 },
      output: {
        prevTemp: { $shift: { output: "$temp", by: -1, default: null } },
        nextTemp: { $shift: { output: "$temp", by: 1 } },
        tempDelta: { $derivative: { input: "$temp", unit: "hour" },
                     window: { range: [-1, 0], unit: "hour" } },
        smooth: { $expMovingAvg: { input: "$temp", N: 5 } }
      }
  } }
])

2.3 累积与统计算子

常见的累积算子包括 $sum、$avg、$min、$max、$count、$stdDevPop、$covariancePop,都可配合窗口做累计或滑动:

db.sales.aggregate([
  { $setWindowFields: {
      partitionBy: "$region",
      sortBy: { orderDate: 1 },
      output: {
        cumSum:   { $sum: "$amount", window: { documents: ["unbounded", "current"] } },
        movAvg7:  { $avg: "$amount", window: { documents: [-6, 0] } },
        cumCount: { $count: {},     window: { documents: ["unbounded", "current"] } }
      }
  } }
])

2.4 $accumulator 自定义窗口聚合

当内置算子不够用时,用 $accumulator 定义自定义累积逻辑:

db.orders.aggregate([
  { $setWindowFields: {
      partitionBy: "$customerId",
      sortBy: { createdAt: 1 },
      output: {
        weightedAvg: {
          $accumulator: {
            init: function() { return { sum: 0, weight: 0 } },
            accumulate: function(state, price, qty) {
              return { sum: state.sum + price * qty, weight: state.weight + qty }
            },
            accumulateArgs: ["$price", "$qty"],
            merge: function(a, b) {
              return { sum: a.sum + b.sum, weight: a.weight + b.weight }
            },
            finalize: function(state) {
              return state.weight ? state.sum / state.weight : 0
            },
            lang: "js"
          },
          window: { documents: ["unbounded", "current"] }
        }
      }
  } }
])

重要:$accumulator 与 $function 都用 JavaScript 引擎执行,开销远高于原生算子,且不支持分片集群的某些优化。能用内置算子表达的逻辑,绝不要写 JS。

3. $lookup 的 pipeline 形式与 $graphLookup

3.1 $lookup 的 pipeline 形式

$lookup 的简写形式只能做等值关联,pipeline 形式则允许在子管道里做任意过滤、投影与关联,并能引用外层字段(通过 let 定义变量)。

db.orders.aggregate([
  { $lookup: {
      from: "items",
      let: { orderId: "$_id", minQty: 2 },
      pipeline: [
        { $match: {
            $expr: { $and: [
              { $eq: ["$orderId", "$$orderId"] },
              { $gte: ["$qty", "$$minQty"] }
            ] }
        } },
        { $project: { name: 1, price: 1, qty: 1 } },
        { $sort: { price: -1 } }
      ],
      as: "items"
  } }
])
形式关联条件子管道能力性能
简写 { localField, foreignField }仅等值无高,可走索引
pipeline 形式任意 $expr完整管道中,需外键索引

3.2 $graphLookup 递归遍历

$graphLookup 做递归图遍历,适合组织架构、分类树、社交关系链:

db.employees.aggregate([
  { $match: { name: "Alice" } },
  { $graphLookup: {
      from: "employees",
      startWith: "$managerId",         // 起始值
      connectFromField: "managerId",   // 从当前文档取哪个字段继续连
      connectToField: "_id",           // 连到目标文档的哪个字段
      as: "chainOfCommand",
      maxDepth: 5,                     // 最大递归深度
      depthField: "level",             // 记录深度
      restrictSearchWithMatch: { status: "active" }
  } }
])

注意:$graphLookup 不做环检测,存在环形引用时会一直递归到 maxDepth 才停。务必设置 maxDepth 上限,并确保关联字段上有索引,否则每一层都会全表扫描。

4. $facet 与 $bucket 分面分析

4.1 $facet 一次扫描多份统计

$facet 把同一份输入喂给多个子管道,每个子管道独立输出,适合电商的"分类计数 + 价格区间 + 标签统计"这类分面导航。

db.products.aggregate([
  { $match: { status: "on_sale" } },
  { $facet: {
      byCategory: [
        { $group: { _id: "$category", n: { $sum: 1 } } },
        { $sort: { n: -1 } }, { $limit: 10 }
      ],
      priceStats: [
        { $group: { _id: null, avg: { $avg: "$price" }, max: { $max: "$price" } } }
      ],
      byTag: [
        { $unwind: "$tags" },
        { $group: { _id: "$tags", n: { $sum: 1 } } },
        { $sort: { n: -1 } }, { $limit: 10 }
      ]
  } }
])

输出是单个文档,含 byCategory、priceStats、byTag 三个数组。

4.2 $bucket 与 $bucketAuto 分桶

$bucket 按自定义边界分桶,$bucketAuto 自动均分桶:

db.products.aggregate([
  { $bucket: {
      groupBy: "$price",
      boundaries: [0, 50, 100, 200, 500],
      default: "other",
      output: { count: { $sum: 1 }, avgPrice: { $avg: "$price" } }
  } }
])
db.products.aggregate([
  { $bucketAuto: { groupBy: "$price", buckets: 5 } }
])
阶段边界来源桶数适用
$bucket手动指定 boundaries由边界决定固定区间统计
$bucketAuto自动计算由 buckets 指定数据分布未知

决策铁律:$facet 的子管道共享一次输入扫描,但如果前面没有 $match 收窄,等于对全集合做多次遍历。分面统计前务必先用索引把候选集缩小。

5. 自定义聚合与物化视图

5.1 $function 自定义表达式

$function 在表达式层执行 JS,适合正则替换、复杂字符串处理等内置算子无法表达的场景:

db.users.aggregate([
  { $project: {
      normalized: {
        $function: {
          body: function(name) { return name.trim().toLowerCase().replace(/\s+/g, "-") },
          args: ["$name"],
          lang: "js"
        }
      }
  } }
])

5.2 $merge 物化视图

$merge 把聚合结果写回集合,支持按条件插入、合并、覆盖,是构建增量物化视图的核心:

db.orders.aggregate([
  { $match: { status: "paid" } },
  { $group: { _id: { date: "$date", region: "$region" },
              total: { $sum: "$amount" }, orders: { $sum: 1 } } },
  { $merge: {
      into: "daily_summary",
      on: "_id",
      whenMatched: "merge",      // 命中则合并字段
      whenNotMatched: "insert"   // 未命中则插入
  } }
])
// $out 直接覆盖目标集合(全量重建)
db.orders.aggregate([
  { $group: { _id: "$region", total: { $sum: "$amount" } } },
  { $out: "region_totals" }
])
阶段行为是否可增量适用
$merge按 on 匹配,可 merge/replace/keepExisting是增量物化视图
$out原子替换整个集合否全量重建快照

重要:$out 会原子替换目标集合,写入期间目标集合短暂不可用;$merge 逐条写入,可增量但需保证 on 字段上有唯一索引,否则可能重复插入。

6. $unwind 与数组处理的性能陷阱

6.1 $unwind 的文档放大效应

$unwind 把数组元素拆成多条文档,展开后的文档数等于数组长度之和。若数组平均 20 个元素,后续每个阶段都要处理 20 倍文档量。

// 展开前:1 条订单含 20 个明细
db.orders.aggregate([
  { $unwind: "$items" },
  { $group: { _id: "$items.sku", qty: { $sum: "$items.qty" } } }
])
// 展开后:20 条文档参与 $group

6.2 保留空数组与索引字段

// 保留空数组或缺失字段的文档
db.orders.aggregate([
  { $unwind: { path: "$items", preserveNullAndEmptyArrays: true } }
])

// 记录数组下标
db.orders.aggregate([
  { $unwind: { path: "$items", includeArrayIndex: "idx" } }
])

6.3 用投影与索引前置裁剪

避免把整个大文档带进 $unwind,先用 $project 裁掉不需要的字段:

db.orders.aggregate([
  { $match: { createdAt: { $gte: ISODate("2026-10-01") } } },  // 先用索引收窄
  { $project: { items: 1 } },                                  // 只带必要字段
  { $unwind: "$items" },
  { $group: { _id: "$items.sku", qty: { $sum: "$items.qty" } } }
])
陷阱表现规避
文档放大展开后文档数暴增前置 $match 收窄候选集
大文档搬运展开时携带全部字段$project 只留必要字段
内存排序展开后 $sort 超 100MB加 allowDiskUse 或前置排序
空数组丢失无明细订单被过滤preserveNullAndEmptyArrays

决策铁律:数组展开是聚合里最容易失控的操作。展开前先问"能不能不展开"——用 $filter、$map、$reduce 在数组内部计算,往往比展开后再 $group 快一个数量级。

7. 总结与最佳实践

  • 窗口函数:$setWindowFields 用 partitionBy 切分、sortBy 排序,$rank/$denseRank 区分并列语义
  • 关联查询:简写 $lookup 只做等值,复杂条件用 pipeline 形式;$graphLookup 必须设 maxDepth
  • 分面统计:$facet 一次扫描多份输出,前置 $match 是性能关键
  • 自定义逻辑:$function 与 $accumulator 用 JS 引擎,能用内置算子就不要写 JS
  • 物化输出:$merge 增量、$out 全量,$merge 的 on 字段需唯一索引
  • 数组处理:优先用 $filter/$reduce 代替 $unwind,必须展开时先投影裁剪

决策铁律:高级算子的性能瓶颈往往不在算子本身,而在它前面有没有用索引收窄输入。任何 $setWindowFields、$facet、$unwind 之前,先确认 $match 是否走了索引;否则再优雅的管道也只是对全集合的暴力扫描。

延伸阅读

继续阅读

探索更多技术文章

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

全部文章 返回首页

「mongodb」更多文章

  1. 数据生命周期、TTL 与冷热归档
  2. $graphLookup 与层次结构建模
  3. GridFS 与大文件存储实践