第22章:Mongo聚合性能优化——千万订单报表怎么跑?

1. 项目背景

业务场景:本地生活电商日均订单量突破 50 万,累计订单表超过 4000 万条。运营部门的周报需求越来越复杂——"我要按城市、类目、时间段(早/中/晚/夜)四个维度看订单分布,同时还要看客单价走势、好评率变化、新老客占比……" 开发用聚合管道实现后,发现报表接口耗时从 2 秒飙升至 45 秒,高峰时期直接超时。运维查看发现 $group 阶段内存使用超过 100MB 限制,触发了 allowDiskUse——数据被溢写到磁盘,速度进一步恶化。

痛点:聚合管道不是"一把梭哈写完就行的"。千万级数据集上,管道顺序、$match 下推、$project 裁剪、$lookup 索引、$group 内存限制、物化视图——每一个环节的疏忽都可能将秒级查询推向分钟级。更常见的是用应用层循环代替聚合管道优化,或者用多次 find + 代码合并代替 $facet 一次搞定。

2. 项目设计

小胖(盯着终端里的 45 秒倒计时):大师!运营周报要四维度交叉分析,我的聚合管道跑了 45 秒还没出来,CPU 100%!这怎么优化?

大师:聚合管道的优化有几个铁律。我们先从你管道的第一个阶段开始看——你写的第一个 stage 是什么?

小胖$group……我想先按城市和类目分组统计。

大师:这就是问题!你的 4000 万条数据全部涌进 $group 阶段做分组——这是在让一个人先把 4000 万颗豆子分类,然后再从里面挑红色的。正确做法是把 过滤放在第一个 $match,把数据量削减到最小,然后再 $group

技术映射:MongoDB 聚合优化器会自动把 $match 向管道前端"下推"(谓词下推),但不能 100% 依赖它——尤其是 $match 写在 $lookup 或者复杂的 $addFields 之后,优化器可能无法下推。

小胖:那我修正顺序——$match(最近 7 天 + 已支付)→ $group(城市+类目)。但这能快多少?

大师:你 4000 万条数据,最近 7 天的可能只有 350 万条(5% 不到)。数据量减了 20 倍,$group 的计算量也就减了 20 倍。这是聚合优化的第一原则——尽早削减数据量

小白:那 $project 也应该放在前面?我听说要先裁剪字段?

大师:对,用完就扔。如果后面的阶段只需要 citycategoryamount 三个字段,就在 $match 之后立刻 $project 把其他 30 个字段全剪掉——减少管道中的数据传输和内存开销。

技术映射$project 裁剪字段 = 列裁剪(Column Pruning)。在 MongoDB 中这意味着每个文档在管道中传递的体积变小,减少序列化开销和内存占用。

小胖:那 $lookup 呢?订单表要关联优惠券表,在 350 万条上做 $lookup 还是很慢。

大师$lookup 性能优化的关键在于外表(被关联的集合)必须有对应外键的索引。$lookup 的本质是对每条输入执行一次外表查询——如果没有索引,每次查询都是全表扫描,350 万次全表扫描 = 灾难。另外,用 pipeline 版本的 $lookup 可以在关联子查询里加 $match 进一步削减。

技术映射$lookup 的 pipeline 版本支持在关联子查询中做 $match$limit,相当于 SQL 的 LEFT JOIN (SELECT ... WHERE ... LIMIT ...),大幅减少关联出来的数据量。

大师(总结):聚合性能优化的五条铁律——① $match 最早;② $project 裁剪紧接其后;③ $lookup 外表必须有索引,pipeline 版本更佳;④ $group 内存不足用 allowDiskUse 保底但不作为首选;⑤ 数据量极大时考虑预聚合(物化视图/定时宽表)。

小胖:预聚合是啥?提前算好存起来?

大师:对。比如你每天的 GMV 不会变了——前一天结束后跑一次日报聚合,结果存入 daily_stats 集合。运营查过去 30 天的趋势时,直接从这个聚合结果集合查 30 条文档——毫秒级返回。$merge$out 阶段就是干这个的。

3. 项目实战

3.1 环境准备

docker compose -f mongodb-lab/docker-compose.yml ps

3.2 分步实现

步骤一:构造千万级订单数据

目标:生成 1000 万条订单(模拟中等规模数据)用于优化前后对比。

use local_life
db.orders_agg_opt.drop()const batchSize = 5000
const total = 500000    // 可根据机器性能调小到 10 万
const totalBatches = Math.ceil(total / batchSize)for (let batch = 0; batch < totalBatches; batch++) {const docs = []const actualSize = Math.min(batchSize, total - batch * batchSize)for (let i = 0; i < actualSize; i++) {const idx = batch * batchSize + iconst daysAgo = Math.floor(Math.random() * 180)docs.push({orderNo: "OPT" + String(idx).padStart(9, '0'),city: ["深圳","广州","北京","上海","杭州","成都","武汉","西安"][idx % 8],category: ["数码","家居","食品","美妆","服饰"][idx % 5],amount: NumberDecimal((Math.random() * 1000 + 10).toFixed(2)),quantity: Math.floor(Math.random() * 5) + 1,isPaid: Math.random() > 0.08,hour: Math.floor(Math.random() * 24),userId: "U" + Math.floor(Math.random() * 10000),couponId: idx % 10 === 0 ? "CP" + (idx % 5) : null,createdAt: new Date(Date.now() - daysAgo * 24 * 3600 * 1000)})}db.orders_agg_opt.insertMany(docs, { ordered: false })if ((batch + 1) % 20 === 0) print(`已插入 ${(batch + 1) * batchSize} 条...`)
}
print(`插入完成: ${db.orders_agg_opt.countDocuments()} 条`)// 建索引
db.orders_agg_opt.createIndex({ isPaid: 1, createdAt: -1 }, { name: "idx_paid_time" })
db.orders_agg_opt.createIndex({ city: 1, category: 1, createdAt: -1 },{ name: "idx_city_cat_time" })
db.orders_agg_opt.createIndex({ createdAt: -1 }, { name: "idx_time" })// 优惠券表
db.coupons_ref.drop()
for (let i = 1; i <= 5; i++) {db.coupons_ref.insertOne({_id: "CP" + i,name: "优惠券" + i,faceValue: NumberDecimal((i * 10).toString())})
}
db.coupons_ref.createIndex({ _id: 1 })

步骤二:优化前后性能对比

目标:用相同的聚合需求,对比"随手写"和"优化后"的性能差异。

const sevenDaysAgo = new Date(Date.now() - 7 * 24 * 3600 * 1000)// === 版本 A:未优化(慢) ===
const slowReport = db.orders_agg_opt.explain("executionStats").aggregate([// ❌ 错误1:$group 放在最前面——所有数据参与分组{$group: {_id: { city: "$city", category: "$category" },orderCount: { $sum: 1 },totalAmount: { $sum: "$amount" }}},// ❌ 错误2:$match 在 $group 后面——过滤掉已分组的结果,非削减输入{ $match: { orderCount: { $gt: 100 } } },{ $sort: { totalAmount: -1 } },{ $limit: 20 }
])
print("=== 未优化版本 ===")
print("  总耗时:", slowReport.executionStats?.executionTimeMillis || "N/A", "ms")// === 版本 B:优化后 ===
const fastReport = db.orders_agg_opt.explain("executionStats").aggregate([// ✓ 优化1:$match 最前面削减 95% 数据{ $match: { isPaid: true, createdAt: { $gte: sevenDaysAgo } } },// ✓ 优化2:$project 裁剪只保留需要的字段{$project: {city: 1, category: 1, amount: 1, _id: 0}},// ✓ 优化3:$group 面对的已是少量数据{$group: {_id: { city: "$city", category: "$category" },orderCount: { $sum: 1 },totalAmount: { $sum: "$amount" }}},{ $match: { orderCount: { $gt: 10 } } },{ $sort: { totalAmount: -1 } },{ $limit: 20 }
])
print("=== 优化版本 ===")
print("  总耗时:", fastReport.executionStats?.executionTimeMillis || "N/A", "ms")

步骤三:$lookup 优化对比

目标:对比无索引和有索引的 $lookup 性能。

// === 版本 A:无索引 $lookup ===
// 先删掉索引来模拟
db.coupons_ref.dropIndexes()  // 只剩下 _id 默认索引const lookupSlow = db.orders_agg_opt.explain("executionStats").aggregate([{ $match: { couponId: { $ne: null } } },{ $limit: 1000 },  // 只取 1000 条来模拟(否则太久){$lookup: {from: "coupons_ref",localField: "couponId",foreignField: "_id",as: "coupon"}},{ $unwind: { path: "$coupon", preserveNullAndEmptyArrays: true } }
])
print("=== $lookup 无索引 ===")// === 版本 B:重建索引 + pipeline 版本 ===
db.coupons_ref.createIndex({ _id: 1 })const lookupFast = db.orders_agg_opt.explain("executionStats").aggregate([{ $match: { couponId: { $ne: null } } },{ $limit: 1000 },{$lookup: {from: "coupons_ref",let: { cid: "$couponId" },pipeline: [{ $match: { $expr: { $eq: ["$_id", "$$cid"] } } },{ $project: { name: 1, faceValue: 1, _id: 0 } }],as: "coupon"}},{ $unwind: { path: "$coupon", preserveNullAndEmptyArrays: true } }
])
print("=== $lookup 有索引 + pipeline 裁剪 ===")
// pipeline 版本在关联时只取需要的字段,减少网络传输
print("  注意: explain 中的 executionStats 对比需在同样数据量下才有意义")

步骤四:$group 内存控制与 allowDiskUse

目标:演示 $group 分组过多时内存溢出及 allowDiskUse 的代价。

// 按城市+类目+小时三维度分组——分组数 = 8×5×24 = 960 组,不会溢出
const smallGroups = db.orders_agg_opt.aggregate([{ $match: { isPaid: true } },{$group: {_id: { city: "$city", category: "$category", hour: "$hour" },count: { $sum: 1 },total: { $sum: "$amount" }}},{ $sort: { total: -1 } },{ $limit: 10 }
], {allowDiskUse: false     // 禁用磁盘溢写,超过 100MB 即报错
}).toArray()print("多维分组结果数:", smallGroups.length)// 极端场景:按 orderNo 分组(每个订单独立一组 = 50 万组)
// 这会导致 $group 结果集超过内存限制
const largeGroups = db.orders_agg_opt.aggregate([{ $match: { isPaid: true } },{$group: {_id: "$orderNo",total: { $sum: "$amount" }}},{ $count: "groups" }
], {allowDiskUse: true     // 内存不够时溢写磁盘
}).toArray()print("大分组结果:", JSON.stringify(largeGroups))
print("  allowDiskUse 在分组过多时会显著增加延迟(写入临时文件)")

步骤五:预聚合——$merge 写入物化视图

目标:用 $merge 将日报结果存入独立的统计集合。

// 计算昨天的日报并写入 daily_stats
const yesterday = new Date(Date.now() - 24 * 3600 * 1000)
const yesterdayStart = new Date(yesterday.getFullYear(), yesterday.getMonth(), yesterday.getDate())
const todayStart = new Date(Date.now())
todayStart.setHours(0, 0, 0, 0)db.orders_agg_opt.aggregate([{ $match: { createdAt: { $gte: yesterdayStart, $lt: todayStart }, isPaid: true } },{$group: {_id: {date: { $dateToString: { format: "%Y-%m-%d", date: "$createdAt" } },city: "$city",category: "$category"},orderCount: { $sum: 1 },totalAmount: { $sum: "$amount" },avgAmount: { $avg: "$amount" },maxAmount: { $max: "$amount" }}},{ $project: { _id: 0, date: "$_id.date", city: "$_id.city",category: "$_id.category", orderCount: 1, totalAmount: 1,avgAmount: 1, maxAmount: 1 } },{$merge: {into: "daily_stats",on: ["date", "city", "category"],   // 按这三个字段 upsertwhenMatched: "replace",             // 存在就替换(同一天跑多次只保留最后一次)whenNotMatched: "insert"            // 不存在就插入}}
])
print("预聚合完成,查 daily_stats:")
const dailyData = db.daily_stats.find({ date: yesterdayStart.toISOString().slice(0,10) }).sort({ orderCount: -1 }).limit(5).toArray()
dailyData.forEach(r => {print(`  ${r.date} | ${r.city}|${r.category} | ${r.orderCount}单 | ¥${Number(r.totalAmount).toFixed(0)}`)
})// 运营查询最近 7 天数据——直接读聚合结果,毫秒级
const weeklyTrend = db.daily_stats.aggregate([{ $match: { date: { $gte: new Date(Date.now() - 7 * 24 * 3600 * 1000).toISOString().slice(0,10) } } },{ $group: { _id: "$date", dailyGMV: { $sum: "$totalAmount" } } },{ $sort: { _id: 1 } }
]).toArray()
print("\n最近 7 天 GMV 趋势(从预聚合表读,毫秒级):")
weeklyTrend.forEach(d => print(`  ${d._id}: ¥${Number(d.dailyGMV).toFixed(0)}`))

步骤六:聚合管道优化清单检查

目标:输出一份聚合优化的自检清单。

// 聚合优化自检脚本
function auditAggregation(pipeline, collectionName) {print(`=== 聚合管道优化审计: ${collectionName} ===`)const stages = pipeline// 1. 检查第一个阶段是否为 $matchconst firstStage = stages[0]const keys = Object.keys(firstStage)if (keys[0] !== '$match') {print("  ⚠ 首个阶段不是 $match,建议前置过滤条件")} else {print("  ✓ 首个阶段是 $match")}// 2. 检查 $match 是否能用索引// (启发式检查:简单字段查询可能走索引)print("  ℹ 请在 $match 字段上运行 explain 确认索引命中")// 3. 检查是否做了 $project 字段裁剪const hasProject = stages.some(s => Object.keys(s)[0] === '$project')if (!hasProject && stages.some(s => Object.keys(s)[0] === '$group')) {print("  ⚠ 有 $group 但无 $project——建议在 $group 前裁剪不需要的字段")} else if (hasProject) {print("  ✓ 有 $project 裁剪")}// 4. 检查 $lookup 是否用 pipeline 版本const hasLookup = stages.filter(s => Object.keys(s)[0] === '$lookup')hasLookup.forEach((l, i) => {if (l.$lookup.pipeline) {print(`  ✓ $lookup[${i}] 使用 pipeline 版本`)} else {print(`  ⚠ $lookup[${i}] 使用传统写法,建议改用 pipeline 版本做裁剪`)}})// 5. 检查是否有不必要的 $unwind → $group → $sort 链// (跳过,过于启发式)// 6. 检查是否设置了 allowDiskUse(只在预期大数据量时用)print("  ℹ 如果 $group 分组数 > 数千,建议评估 allowDiskUse 的必要性")print("=== 审计完成 ===")
}// 示例使用
auditAggregation([{ $match: { status: "在售" } },{ $project: { name: 1, category: 1, price: 1 } },{ $group: { _id: "$category", count: { $sum: 1 } } }
], "products")

3.3 完整代码清单

文件 用途
mongodb-lab/scripts/ch22-create-order-data.js 构造大量订单测试数据
mongodb-lab/scripts/ch22-agg-optimize.js 优化前/后对比
mongodb-lab/scripts/ch22-lookup-perf.js $lookup 性能测试
mongodb-lab/scripts/ch22-group-memory.js $group 内存与 allowDiskUse
mongodb-lab/scripts/ch22-preaggregate.js $merge 预聚合写物化视图
mongodb-lab/scripts/ch22-audit-pipeline.js 聚合优化自检清单

3.4 测试验证

use local_life// 1. 验证 $match 下推效果
const optQuery = db.orders_agg_opt.aggregate([{ $match: { isPaid: true, createdAt: { $gte: new Date(Date.now() - 7 * 24 * 3600 * 1000) } } },{ $group: { _id: "$city", count: { $sum: 1 } } }
]).explain("executionStats")
print("$match 后扫描文档数(应 << 总数):", optQuery.executionStats?.nReturned || "N/A")// 2. 验证 $merge 结果
const today = new Date().toISOString().slice(0, 10)
const merged = db.daily_stats.findOne({ date: today })
print("$merge 写入:", merged ? "PASS (有今天数据)" : "需先执行步骤五的预聚合")// 3. 验证 $lookup 外表索引
const lookupIndexes = db.coupons_ref.getIndexes()
print("coupons_ref 索引数:", lookupIndexes.length, lookupIndexes.length > 1 ? "PASS" : "仅_id索引")print("\n=== 聚合优化验证完成 ===")

4. 项目总结

4.1 聚合优化法则速查

法则 操作 收益
过滤前置 $match 放在第一个阶段 削减 80%-99% 数据量
列裁剪 $match 后紧跟 $project 减少 50%-80% 文档体积
索引 $match 在第一 $match 字段上建索引 IXSCAN 替代 COLLSCAN
索引 $sort 排序字段包含在索引中 消除内存 SORT
索引 $lookup 外表外键建索引 消除全表扫描
管道 $lookup pipeline 版本 + $limit 减少关联传输
预聚合 $merge / $out 写物化视图 将分钟查询变毫秒查询
allowDiskUse 在内存满时兜底 避免因 100MB 限制报错

4.2 适用场景

聚合优化适用

  1. 运营报表/仪表盘——按多维分组 + 时间范围 + 排序 Top N。
  2. 数据管道(ETL)——从业务表聚合后写入分析表。
  3. 实时大屏——预聚合的毫秒级读取配合 WebSocket 推送。
  4. 对账和清算——聚合计算总额与实际交易流水核对。

不适用场景

  1. 实时窗口计算(滑动窗口、漏斗分析)——MongoDB 聚合管道非流式计算引擎,用 Flink/Spark Streaming。
  2. 图关系计算(如社交二度关系)——不适合利用 MongoDB 聚合管道,用图数据库。

4.3 注意事项

注意事项 说明
allowDiskUse 不是性能优化 它是"内存不够"时的兜底,启用后写入磁盘,耗时会急剧增加
$lookup 不自动下推索引 外表查询不保证走索引,需要 explain 确认
$group 的 100MB 限制 是按单个 pipeline 执行器的单批内存算的,不是按 _id 分组数
$merge 写入可能失败 如果 on 字段的组合有唯一冲突,whenMatched/replace 比较好;否则考虑 whenNotMatched/insert
预聚合数据与源数据的延迟 预聚合是一次性快照,源数据变更后统计表不会自动更新——需定时任务刷新

4.4 常见踩坑经验

故障案例一:$group 数据倾斜导致内存溢出

某报表按 category 分组,其中一个类目占数据量的 80%,聚合过程变慢且内存占用集中在处理这个类目的片段上。根因:数据倾斜使得某个分组的聚合计算时间远超其他分组。解决:在 $group 前加 $match 将大类目拆成多个条件并行处理($facet),最后合并结果。

故障案例二:$lookup 外表在分片集群中无索引导致全分片扫描

某团队分片集群下单表在 shard_A,用户表在 shard_B。$lookup 的 userId 在用户表上没有分片键——mongos 无法定位到具体的 shard,于是向所有分片广播查询——每个分片都跑一次全表扫描。解决:在用户表的 userId 上建索引,且确保 $lookup 的外键能被下推到单个分片。

故障案例三:$sort 内存排序导致 OOM

某聚合在 $group 之后 $sort,排序字段不在分组键中——MongoDB 需要把所有分组结果加载到内存中排序,分组数超过 5 万时内存溢出。解决:在索引中预设排序字段;或者在 $group_id 中加入排序字段的前缀,让排序利用索引顺序。

4.5 思考题

  1. 为什么 $merge 不支持在分片集群中将结果写入分片集合的任意 shard?写入目标必须是非分片的集合或者指定 on 字段等于分片键,为什么?
  2. allowDiskUse: true 和直接增大 WiredTiger 缓存哪个更适合聚合优化?两者的区别和适用场景是什么?

(答案将在第 23 章末尾揭晓)


上一章思考题答案

  1. 优化器没有选择 {a:1,c:1,b:1} 索引去消除 SORT 的可能原因:① ESR 原则——等值条件 a, b 在范围/排序条件 c 之前,但 {a:1,b:1,c:1} 是正确 ESR 顺序,而 {a:1,c:1,b:1} 中 b 在排序字段 c 之后,b 作为等值条件无法利用——优化器可能认为该索引不如另一个在等值过滤上更高效的索引;② 候选计划竞速时该索引恰好在初期返回批量数据较慢因此被"淘汰"。用 hint 验证该索引能否消除 SORT 再对优化器行为进行推理是更直接的方法。

  2. $planCacheStats(MongoDB 6.0+)提供了每个查询形状的缓存统计(命中次数、占用空间、最后访问时间),而旧版 planCacheList() 只列出查询形状。评估"性价比"的方法:如果一个缓存的 cacheSize 很大但 accesses.ops 极少(比如一天命中 3 次的缓存占用了 10MB),那它就是低性价比的——可以通过 planCacheClear() 选择性清除,或者优化查询让其不生成 Plan Cache(如使用 hint 让优化器跳过竞速)。

延伸阅读与资源

python入门:Rquests从菜鸟脚本到企业级SDK的网络实战圣经
Milvus向量数据库实战修炼:从 0 到 1精通向量检索与生产落地
后端工程师的 AI 转型第一课:Ollama 与私有化大模型实战
10倍开发者的 Dify 魔法书:从零构建全栈 AI 应用
后端工程师转型AI第一课-Ollama 与私有化大模型实战

大型语言模型(LLM) vLLM 高性能推理落地实战

Agent开发之LlamaIndex 实战修炼与源码进阶

大语言模型Transformers 实战修炼与源码剖析