MongoDB .NET Driver聚合管道:关联集合并重组输出结构
解决方案
假设的集合文档结构
先明确三个集合的示例文档(如果你的结构有差异,可对应调整代码):
Portfolio集合
{ "_id": ObjectId("60d21b4667d0d8992e610c85"), "userId": "user123", "Stocks": [ { "stockId": "AAPL", "quantity": 10 }, { "stockId": "MSFT", "quantity": 5 } ], "MutualFunds": [ { "fundId": "VTSAX", "units": 20 }, { "fundId": "VFIAX", "units": 15 } ] }
Stocks集合
{ "_id": ObjectId("60d21b8967d0d8992e610c86"), "stockId": "AAPL", "name": "Apple Inc.", "currentPrice": 175.50, "sector": "Technology" }
MutualFunds集合
{ "_id": ObjectId("60d21bc167d0d8992e610c88"), "fundId": "VTSAX", "name": "Vanguard Total Stock Market Index Fund", "nav": 85.10, "fundFamily": "Vanguard" }
目标输出格式
{ "userId": "user123", "Holdings": { "Stocks": [ { "stockId": "AAPL", "name": "Apple Inc.", "currentPrice": 175.50, "sector": "Technology", "quantity": 10, "totalValue": 1755.00 } ], "MutualFunds": [ { "fundId": "VTSAX", "name": "Vanguard Total Stock Market Index Fund", "nav": 85.10, "fundFamily": "Vanguard", "units": 20, "totalValue": 1702.00 } ] } }
.NET Core 6 代码实现
1. 定义实体类
using MongoDB.Bson; using MongoDB.Bson.Serialization.Attributes; // Portfolio集合实体 public class Portfolio { public ObjectId Id { get; set; } [BsonElement("userId")] public string UserId { get; set; } public List<PortfolioStock> Stocks { get; set; } public List<PortfolioMutualFund> MutualFunds { get; set; } } // Portfolio中的股票持仓项 public class PortfolioStock { [BsonElement("stockId")] public string StockId { get; set; } public int Quantity { get; set; } } // Portfolio中的基金持仓项 public class PortfolioMutualFund { [BsonElement("fundId")] public string FundId { get; set; } public int Units { get; set; } } // Stocks集合实体 public class Stock { public ObjectId Id { get; set; } [BsonElement("stockId")] public string StockId { get; set; } public string Name { get; set; } [BsonElement("currentPrice")] public decimal CurrentPrice { get; set; } public string Sector { get; set; } } // MutualFunds集合实体 public class MutualFund { public ObjectId Id { get; set; } [BsonElement("fundId")] public string FundId { get; set; } public string Name { get; set; } public decimal Nav { get; set; } [BsonElement("fundFamily")] public string FundFamily { get; set; } } // 最终输出DTO public class PortfolioHoldingsDto { public string UserId { get; set; } public HoldingsDto Holdings { get; set; } } public class HoldingsDto { public List<StockHoldingDto> Stocks { get; set; } public List<MutualFundHoldingDto> MutualFunds { get; set; } } public class StockHoldingDto { public string StockId { get; set; } public string Name { get; set; } public decimal CurrentPrice { get; set; } public string Sector { get; set; } public int Quantity { get; set; } public decimal TotalValue { get; set; } } public class MutualFundHoldingDto { public string FundId { get; set; } public string Name { get; set; } public decimal Nav { get; set; } public string FundFamily { get; set; } public int Units { get; set; } public decimal TotalValue { get; set; } }
2. 聚合管道实现代码
using MongoDB.Driver; public async Task<PortfolioHoldingsDto> GetUserPortfolioHoldings(string userId) { // 初始化MongoDB客户端和集合 var client = new MongoClient("mongodb://localhost:27017"); // 替换为你的连接字符串 var db = client.GetDatabase("YourDatabaseName"); // 替换为你的数据库名 var portfolioCol = db.GetCollection<Portfolio>("Portfolio"); // 构建聚合管道 var pipeline = new BsonDocument[] { // 筛选目标用户的Portfolio new("$match", new BsonDocument("userId", userId)), // 关联Stocks集合 new("$lookup", new BsonDocument { {"from", "Stocks"}, {"localField", "Stocks.stockId"}, {"foreignField", "stockId"}, {"as", "matchedStocks"} }), // 关联MutualFunds集合 new("$lookup", new BsonDocument { {"from", "MutualFunds"}, {"localField", "MutualFunds.fundId"}, {"foreignField", "fundId"}, {"as", "matchedFunds"} }), // 合并持仓数量与主集合数据 new("$addFields", new BsonDocument { "mergedStocks", new("$map", new BsonDocument { {"input", "$Stocks"}, {"as", "portStock"}, {"in", new("$mergeObjects", new BsonArray { new("$arrayElemAt", new BsonArray { "$matchedStocks", new("$indexOfArray", new BsonArray {"$matchedStocks.stockId", "$$portStock.stockId"}) }), "$$portStock" }) } }), "mergedFunds", new("$map", new BsonDocument { {"input", "$MutualFunds"}, {"as", "portFund"}, {"in", new("$mergeObjects", new BsonArray { new("$arrayElemAt", new BsonArray { "$matchedFunds", new("$indexOfArray", new BsonArray {"$matchedFunds.fundId", "$$portFund.fundId"}) }), "$$portFund" }) } }) }), // 重组结构并计算总价值 new("$project", new BsonDocument { {"_id", 0}, {"userId", 1}, {"Holdings", new BsonDocument { {"Stocks", new("$map", new BsonDocument { {"input", "$mergedStocks"}, {"as", "stock"}, {"in", new BsonDocument { {"stockId", "$$stock.stockId"}, {"name", "$$stock.name"}, {"currentPrice", "$$stock.currentPrice"}, {"sector", "$$stock.sector"}, {"quantity", "$$stock.quantity"}, {"totalValue", new("$multiply", new BsonArray {"$$stock.currentPrice", "$$stock.quantity"})} } } })}, {"MutualFunds", new("$map", new BsonDocument { {"input", "$mergedFunds"}, {"as", "fund"}, {"in", new BsonDocument { {"fundId", "$$fund.fundId"}, {"name", "$$fund.name"}, {"nav", "$$fund.nav"}, {"fundFamily", "$$fund.fundFamily"}, {"units", "$$fund.units"}, {"totalValue", new("$multiply", new BsonArray {"$$fund.nav", "$$fund.units"})} } } })} }} }) }; // 执行聚合并返回结果 return await portfolioCol.Aggregate<PortfolioHoldingsDto>(pipeline).FirstOrDefaultAsync(); }
关键步骤说明
- $match:先筛选目标用户的Portfolio,减少后续处理的数据量。
- $lookup:通过
stockId/fundId关联主集合,将匹配结果存入临时数组。 - $map + $mergeObjects:将Portfolio中的持仓数量/单位信息,与关联到的主集合数据合并,确保每个持仓项信息完整。
- $project:重组输出结构,计算持仓总价值,移除不需要的字段。
注意事项
- 确保集合字段名大小写与代码一致(MongoDB默认区分大小写)。
- 若存在持仓ID在主集合中不存在的情况,可添加
$ifNull处理空值。 - 强类型实体需通过
[BsonElement]属性映射MongoDB中的字段名。
内容的提问来源于stack exchange,提问作者Rithik Banerjee
相关产品推荐
相关产品推荐

