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的复合索引:
该索引让MongoDB直接从索引获取有序数据,避免内存排序。db.ActivityFeedItems.createIndex({ UserId: 1, CreatedOnUtc: -1 }) - 加快
$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ı
相关产品推荐
相关产品推荐

