You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

MongoDB聚合查询CPU占用过高致无响应的优化方案问询

问题

我有一个gRPC方法,接收约100个friendIds,查询MongoDB并返回20条按groupId分组的记录。但执行过程中CPU占用大幅飙升,导致MongoDB出现无响应情况。目前数据库内约有134000条记录,且数据量预计会持续增长。

当前实现代码:

public override async Task<GrpcActivityFeed.ActivityFeeds> Get(GrpcActivityFeed.ActivityFeedGetRequest request, ServerCallContext context)
{
    var activityFeedCollection = _dbContext.GetCollection<ActivityFeedItem>("ActivityFeedItems");
    var friendIdList = request.UserFriendIds.ToList();

    DateTime? maxDate = null;
    if (!string.IsNullOrEmpty(request.MaxDate))
    {
        maxDate = DateTime.Parse(request.MaxDate);
    }

    var pipeline = new List<BsonDocument>
    {
        new BsonDocument("$match", new BsonDocument("UserId", new BsonDocument("$in", new BsonArray(friendIdList)))),
        new BsonDocument("$sort", new BsonDocument("CreatedOnUtc", -1)),
        new BsonDocument("$group", new BsonDocument
        {
            { "_id", "$GroupId" },
            { "latestItem", new BsonDocument("$first", "$$ROOT") }
        }),
        new BsonDocument("$lookup", new BsonDocument
        {
            { "from", "ActivityFeedItems" },
            { "localField", "_id" },
            { "foreignField", "GroupId" },
            { "as", "groupedItems" }
        }),
        new BsonDocument("$sort", new BsonDocument("latestItem.CreatedOnUtc", -1)),
        new BsonDocument("$skip", request.Skip),
        new BsonDocument("$limit", request.Take),
        new BsonDocument("$project", new BsonDocument
        {
            { "latestItem.CreatedOnUtc", 1 },
            { "latestItem.UserId", 1 },
            { "latestItem.Type", 1 },
            { "latestItem.User", 1 },
            { "latestItem.Collection", 1 },
            { "groupedItems", 1 }
        })
    };

    var aggregation = await activityFeedCollection.AggregateAsync<BsonDocument>(pipeline);
    var resultList = await aggregation.ToListAsync();
    var finalResult = new List<List<ActivityFeedItem>>();

    foreach (var document in resultList)
    {
        if (!document.Contains("groupedItems"))
        {
            continue; // groupedItems field not found, skip this document
        }

        var items = document["groupedItems"].AsBsonArray;
        var activityFeedItems = items.Select(item => BsonSerializer.Deserialize<ActivityFeedItem>(item.AsBsonDocument)).ToList();
        finalResult.Add(activityFeedItems);
    }

    var activityFeeds = new GrpcActivityFeed.ActivityFeeds();
    foreach (var result in finalResult.OrderByDescending(p => p.OrderByDescending(x => x.CreatedOnUtc).FirstOrDefault().CreatedOnUtc))
    {
        var groupedItems = result.ToList();
        var latestItem = groupedItems.OrderByDescending(item => item.CreatedOnUtc).FirstOrDefault();

        var activityFeed = new ActivityFeed
        {
            Id = latestItem.Id,
            UserId = latestItem.UserId,
            Type = latestItem.Type,
            CreatedOnUtc = latestItem?.CreatedOnUtc,
            User = latestItem?.User,
            GroupedItems = groupedItems,
            TotalItemCount = groupedItems.Count,
            CollectionType = latestItem?.Collection?.Entity?.CollectionType
        };

        var grpcActivity = _mapper.Map<GrpcActivityFeed.ActivityFeed>(activityFeed);
        activityFeeds.ActivityFeeds_.Add(grpcActivity);
    }

    return activityFeeds;
}

MongoDB查询日志:

{
  "type": "command",
  "ns": "***.ActivityFeedItems",
  "command": {
    "aggregate": "ActivityFeedItems",
    "pipeline": [
      {
        "$match": {
          "UserId": {
            "$in": [
              // userIds,
              212
            ]
          }
        }
      },
      {
        "$group": {
          "_id": "$GroupId",
          "latestItem": {
            "$first": "$$ROOT"
          }
        }
      },
      {
        "$sort": {
          "latestItem.CreatedOnUtc": -1
        }
      },
      {
        "$skip": 0
      },
      {
        "$limit": 20
      },
      {
        "$lookup": {
          "from": "ActivityFeedItems",
          "localField": "_id",
          "foreignField": "GroupId",
          "as": "groupedItems"
        }
      },
      {
        "$project": {
          "latestItem.CreatedOnUtc": 1,
          "latestItem.UserId": 1,
          "latestItem.Type": 1,
          "latestItem.User": 1,
          "latestItem.Collection": 1,
          "groupedItems": 1
        }
      }
    ],
    "cursor": {
      
    },
    "$db": "***",
    "lsid": {
      "id": {
        "$binary": {
          "base64": "B3dht3N+QI2rPdH0oqcLdQ==",
          "subType": "04"
        }
      }
    },
    "$clusterTime": {
      "clusterTime": {
        "$timestamp": {
          "t": 1716875987,
          "i": 1
        }
      },
      "signature": {
        "hash": {
          "$binary": {
            "base64": "2oLIMkZ7L6RpuGikgdfKDONMkY0=",
            "subType": "00"
          }
        },
        "keyId": 7340668007347126000
      }
    }
  },
  "planSummary": "IXSCAN { UserId: 1 }",
  "keysExamined": 64152,
  "docsExamined": 64152,
  "hasSortStage": true,
  "cursorExhausted": true,
  "numYields": 64,
  "nreturned": 20,
  "queryHash": "54B9C077",
  "planCacheKey": "FD407F29",
  "queryFramework": "sbe",
  "reslen": 40127,
  "locks": {
    "FeatureCompatibilityVersion": {
      "acquireCount": {
        "r": 116
      }
    },
    "Global": {
      "acquireCount": {
        "r": 116
      }
    },
    "Mutex": {
      "acquireCount": {
        "r": 52
      }
    }
  },
  "readConcern": {
    "level": "local",
    "provenance": "implicitDefault"
  },
  "writeConcern": {
    "w": "majority",
    "wtimeout": 0,
    "provenance": "implicitDefault"
  },
  "storage": {
    
  },
  "remote": "***",
  "protocol": "op_msg",
  "durationMillis": 426,
  "v": "6.0.15",
  "isTruncated": false
}

我正在寻找优化该查询以降低CPU占用、避免MongoDB无响应的最佳方案,恳请提供聚合管道或整体实现方案的改进建议。


优化方案

1. 调整聚合管道顺序,减少数据处理量

当前管道先对所有符合条件的文档分组,再做分页,最后才关联查询,导致大量不必要的计算。优化顺序后能将$lookup的操作范围缩小到仅目标20个分组:

  • 先通过$match过滤数据(务必加上之前未使用的maxDate条件,避免扫描全量历史数据)
  • 按CreatedOnUtc降序排序,确保$group能直接拿到每组最新项
  • 对分组结果排序后执行$skip/$limit,仅保留目标20组
  • 最后对这20组做$lookup关联查询

优化后的管道示例:

var pipeline = new List<BsonDocument>
{
    new BsonDocument("$match", new BsonDocument
    {
        { "UserId", new BsonDocument("$in", new BsonArray(friendIdList)) },
        { "CreatedOnUtc", maxDate.HasValue ? new BsonDocument("$lte", maxDate) : new BsonDocument() }
    }),
    new BsonDocument("$sort", new BsonDocument("CreatedOnUtc", -1)),
    new BsonDocument("$group", new BsonDocument
    {
        { "_id", "$GroupId" },
        { "latestItem", new BsonDocument("$first", "$$ROOT") }
    }),
    new BsonDocument("$sort", new BsonDocument("latestItem.CreatedOnUtc", -1)),
    new BsonDocument("$skip", request.Skip),
    new BsonDocument("$limit", request.Take),
    new BsonDocument("$lookup", new BsonDocument
    {
        { "from", "ActivityFeedItems" },
        { "localField", "_id" },
        { "foreignField", "GroupId" },
        { "as", "groupedItems" },
        // 子管道提前排序,避免客户端重复操作
        { "pipeline", new BsonArray
            {
                new BsonDocument("$sort", new BsonDocument("CreatedOnUtc", -1))
            }
        }
    }),
    new BsonDocument("$project", new BsonDocument
    {
        { "latestItem.CreatedOnUtc", 1 },
        { "latestItem.UserId", 1 },
        { "latestItem.Type", 1 },
        { "latestItem.User", 1 },
        { "latestItem.Collection", 1 },
        { "groupedItems", 1 }
    })
};

2. 添加复合索引,消除内存排序

当前仅使用UserId单字段索引,$match后仍需内存排序,CPU消耗极高。创建以下两个索引:

  • 覆盖$match和$sort的复合索引:
    db.ActivityFeedItems.createIndex({ UserId: 1, CreatedOnUtc: -1 })
    
    该索引让MongoDB直接从索引获取有序数据,避免内存排序。
  • 加快$lookup的GroupId索引:
    db.ActivityFeedItems.createIndex({ GroupId: 1 })
    

3. 简化客户端逻辑,减少内存开销

当前客户端存在重复排序、重复反序列化的冗余操作,可直接利用聚合管道的结果简化代码:

var aggregation = await activityFeedCollection.AggregateAsync<BsonDocument>(pipeline);
var resultList = await aggregation.ToListAsync();

var activityFeeds = new GrpcActivityFeed.ActivityFeeds();
foreach (var document in resultList)
{
    if (!document.Contains("groupedItems")) continue;

    var groupedItems = document["groupedItems"].AsBsonArray
        .Select(item => BsonSerializer.Deserialize<ActivityFeedItem>(item.AsBsonDocument))
        .ToList();
    var latestItem = BsonSerializer.Deserialize<ActivityFeedItem>(document["latestItem"].AsBsonDocument);

    var activityFeed = new ActivityFeed
    {
        Id = latestItem.Id,
        UserId = latestItem.UserId,
        Type = latestItem.Type,
        CreatedOnUtc = latestItem.CreatedOnUtc,
        User = latestItem.User,
        GroupedItems = groupedItems,
        TotalItemCount = groupedItems.Count,
        CollectionType = latestItem?.Collection?.Entity?.CollectionType
    };

    var grpcActivity = _mapper.Map<GrpcActivityFeed.ActivityFeed>(activityFeed);
    activityFeeds.ActivityFeeds_.Add(grpcActivity);
}

4. 限制单批次friendIds数量

一次传入100个friendIds会导致$in扫描大量文档,可拆分请求为每次20-30个friendIds,再合并结果;或优先查询近期有活动的friendIds,缩小扫描范围。

5. 预聚合数据(长期优化方案)

若数据量持续增长,实时聚合的压力会越来越大。可通过MongoDB变更流或定时任务,预聚合每个Group的最新活动和组内数据到单独集合(如ActivityFeedGroups),查询时直接读取该集合,彻底避免实时聚合开销。

预聚合集合示例结构:

{
  "_id": "groupId",
  "latestItem": { /* 最新活动项 */ },
  "groupedItems": [ /* 组内最近N条活动项 */ ],
  "totalItemCount": 10,
  "updatedAt": ISODate("2024-06-01T00:00:00Z")
}

内容的提问来源于stack exchange,提问作者İbrahim Halil Saçlı

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.23 07:02:03