如何将Cosmos中键值对格式的发票数据合并为单发票多明细结构?
发票行项目合并并存储到Cosmos DB的实现方案
针对将Cosmos DB中分散的发票行数据合并为「单发票+多行项目」结构并存储到新容器的需求,以下是几个实用的实现方案:
方案一:Cosmos DB内置查询+存储过程
1. 分组查询源数据
直接用Cosmos DB的SQL语句按InvoiceNumber分组聚合,生成符合目标结构的数据集:
SELECT c.InvoiceNumber, ARRAY_AGG({ "InvoiceLineAmount": c.InvoiceLineAmount, "InvoiceId": c.InvoiceId, "AccountingDate": c.AccountingDate, "SupplierNumber": c.SupplierNumber, "SupplierSite": c.SupplierSite }) AS InvoiceLineDeatils FROM c GROUP BY c.InvoiceNumber
2. 编写存储过程批量写入
将查询结果批量写入目标容器,减少网络请求开销。以下是JavaScript存储过程示例:
function bulkMergeInvoices() { const collection = getContext().getCollection(); const response = getContext().getResponse(); // 执行分组查询 const query = `SELECT c.InvoiceNumber, ARRAY_AGG({ "InvoiceLineAmount": c.InvoiceLineAmount, "InvoiceId": c.InvoiceId, "AccountingDate": c.AccountingDate, "SupplierNumber": c.SupplierNumber, "SupplierSite": c.SupplierSite }) AS InvoiceLineDeatils FROM c GROUP BY c.InvoiceNumber`; const isQueryAccepted = collection.queryDocuments( collection.getSelfLink(), query, {}, (err, feed) => { if (err) throw err; if (!feed.length) { response.setBody('未找到可合并的发票数据'); return; } // 替换为目标容器的链接 const targetCollLink = 'dbs/你的数据库名/colls/目标容器名'; let writeCount = 0; feed.forEach(doc => { const isWriteAccepted = collection.createDocument( targetCollLink, doc, err => { if (err) throw err; writeCount++; if (writeCount === feed.length) { response.setBody(`成功合并并写入 ${writeCount} 条发票数据`); } } ); if (!isWriteAccepted) return; }); } ); if (!isQueryAccepted) throw new Error('查询请求未被服务器接受'); }
方案二:代码端处理(以C#为例)
如果需要更灵活的逻辑控制,可通过Cosmos DB SDK在代码中完成数据读取、分组和写入:
1. 读取源容器数据
using Azure.Cosmos; var cosmosClient = new CosmosClient("你的Cosmos DB连接字符串"); var sourceContainer = cosmosClient.GetContainer("你的数据库名", "源容器名"); // 读取所有发票行数据 var query = new QueryDefinition("SELECT * FROM c"); var invoiceLines = new List<InvoiceLine>(); using var iterator = sourceContainer.GetItemQueryIterator<InvoiceLine>(query); while (iterator.HasMoreResults) { var response = await iterator.ReadNextAsync(); invoiceLines.AddRange(response.Resource); }
2. 分组构建目标结构
用LINQ按发票号分组,生成合并后的实体:
using System.Linq; // 定义实体类 public class InvoiceLine { public string InvoiceNumber { get; set; } public string InvoiceLineAmount { get; set; } public int InvoiceId { get; set; } public string AccountingDate { get; set; } public string SupplierNumber { get; set; } public string SupplierSite { get; set; } } public class MergedInvoice { public string InvoiceNumber { get; set; } public List<InvoiceLineDetail> InvoiceLineDeatils { get; set; } } public class InvoiceLineDetail { public string InvoiceLineAmount { get; set; } public int InvoiceId { get; set; } public string AccountingDate { get; set; } public string SupplierNumber { get; set; } public string SupplierSite { get; set; } } // 分组合并 var mergedInvoices = invoiceLines .GroupBy(line => line.InvoiceNumber) .Select(group => new MergedInvoice { InvoiceNumber = group.Key, InvoiceLineDeatils = group.Select(line => new InvoiceLineDetail { InvoiceLineAmount = line.InvoiceLineAmount, InvoiceId = line.InvoiceId, AccountingDate = line.AccountingDate, SupplierNumber = line.SupplierNumber, SupplierSite = line.SupplierSite }).ToList() }) .ToList();
3. 批量写入目标容器
var targetContainer = cosmosClient.GetContainer("你的数据库名", "目标容器名"); // 按分区键分组批量写入,避免跨分区事务问题 foreach (var batchGroup in mergedInvoices.GroupBy(m => m.InvoiceNumber)) { using var batch = targetContainer.CreateTransactionalBatch(new PartitionKey(batchGroup.Key)); foreach (var invoice in batchGroup) { batch.CreateItem(invoice); } var batchResponse = await batch.ExecuteAsync(); if (!batchResponse.IsSuccessStatusCode) { // 处理写入失败逻辑 throw new Exception($"批量写入失败:{batchResponse.StatusCode}"); } }
方案三:ETL工具自动化处理(如Azure Data Factory)
如果不需要编写代码,可通过ETL工具完成整个流程:
- 源数据集:配置连接到Cosmos DB源容器,读取所有发票行数据
- 转换活动:添加「聚合」转换,选择
InvoiceNumber作为分组键,将其他字段聚合为数组(对应InvoiceLineDeatils) - 目标数据集:配置连接到Cosmos DB目标容器,将转换后的结构数据写入
关键注意事项
- 源容器的分区键设计要合理,避免分组查询时的跨分区扫描,提升查询性能
- 批量写入时需关注Cosmos DB的吞吐量(RU/s)限制,避免触发限流
- 若需实时处理新增数据,可结合Cosmos DB的变更捕获(Change Feed)机制,实时合并新增的发票行
内容的提问来源于stack exchange,提问作者abdul kalam
相关产品推荐
相关产品推荐

