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

如何流式导入归档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文档+批量插入的方式解决,与归档时的游标思路一致:

实现逻辑

  1. 流式解析单个Bson文档:Bson文档以4字节小端序整数开头,表示文档总长度。利用这一特性,逐个读取每个文档的完整字节,无需加载全量数据到内存。
  2. 批量插入优化:使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 12:40:30