MongoDB聚合管道实战:分组统计与数据转换操作符详解

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/

(0)
小编小编
上一篇 11小时前
下一篇 11小时前

相关推荐