如何流式导入归档MongoDB集合至MongoDB?解决大集合导入报错
MongoDB 流式归档后大集合导入报错及解决方案
归档流程(高效实现)
通过IAsyncCursor遍历MongoDB集合,将原始Bson字节直接写入Azure Blob流,代码如下:
var cursor = await clientDb.GetCollection<RawBsonDocument>(collectionPath).Find(new BsonDocument()).ToCursorAsync(); while (cursor.MoveNext()) foreach (var document in cursor.Current) { var bytes = new byte[document.Slice.Length]; document.Slice.GetBytes(0, bytes, 0, document.Slice.Length); blobStream.Write(bytes, 0, bytes.Length); }
导入问题描述
导入归档数据时,目前需将全量字节载入内存流再插入,代码如下:
var rawRef = clientDb.GetCollection<RawBsonDocument>(collectionPath); using (var ms = new MemoryStream()) { await stream.CopyToAsync(ms); var bytes = ms.ToArray(); var rawBson = new RawBsonDocument(bytes); await rawRef.InsertOneAsync(rawBson); }
该方式仅适用于小集合,大集合导入时触发报错:
MongoDB.Driver.MongoConnectionException : An exception occurred while sending a message to the server. ---- System.IO.IOException : Unable to write data to the transport connection: An established connection was aborted by the software in your host machine.. -------- System.Net.Sockets.SocketException : An established connection was aborted by the software in your host machine.
提问:是否支持流式将原始字节数据导入MongoDB,或使用类似读取时的游标方式?
解决方案
问题根源在于一次性将整个大集合的Bson字节读入内存,既造成内存过载,又触发了MongoDB的请求大小/连接超时限制。可通过流式逐个解析归档的Bson文档+批量插入的方式解决,与归档时的游标思路一致:
实现逻辑
- 流式解析单个Bson文档:Bson文档以4字节小端序整数开头,表示文档总长度。利用这一特性,逐个读取每个文档的完整字节,无需加载全量数据到内存。
- 批量插入优化:使用
InsertManyAsync批量提交文档,减少网络请求次数,降低连接中断风险。
示例代码
var rawRef = clientDb.GetCollection<RawBsonDocument>(collectionPath); const int batchSize = 100; // 根据单文档大小调整,避免单次请求超MongoDB消息上限 var documentsBatch = new List<RawBsonDocument>(batchSize); // 确保流处于起始位置 stream.Seek(0, SeekOrigin.Begin); while (true) { // 读取Bson文档长度前缀(4字节) var lengthBuffer = new byte[4]; var readCount = await stream.ReadAsync(lengthBuffer, 0, 4); if (readCount < 4) break; // 流读取完毕 var documentLength = BitConverter.ToInt32(lengthBuffer, 0); if (documentLength <= 4) break; // 无效文档,终止解析 // 读取完整文档字节 var documentBytes = new byte[documentLength]; Buffer.BlockCopy(lengthBuffer, 0, documentBytes, 0, 4); readCount = await stream.ReadAsync(documentBytes, 4, documentLength - 4); if (readCount != documentLength - 4) break; // 文档不完整,终止 // 添加到批量列表 documentsBatch.Add(new RawBsonDocument(documentBytes)); // 达到批量阈值则插入 if (documentsBatch.Count >= batchSize) { await rawRef.InsertManyAsync(documentsBatch); documentsBatch.Clear(); } } // 插入剩余文档 if (documentsBatch.Count > 0) { await rawRef.InsertManyAsync(documentsBatch); }
关键注意点
- 批量大小调整:MongoDB默认单条消息最大为48MB(
maxMessageSizeBytes),需根据单文档平均大小计算batchSize,确保批量总大小不超过该限制。 - 异常处理:可添加重试逻辑(针对临时连接中断)、文档校验逻辑,提升导入稳定性。
- 流状态检查:开始解析前需确保Azure Blob流处于起始位置,避免读取偏移错误。
内容的提问来源于stack exchange,提问作者Ryan Langton
相关产品推荐
相关产品推荐

