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

C# Async Await优化SQL到MongoDB迁移:批量插入异步协同需求

这场景太常见了!你要的就是异步生产者-消费者模式:一边不停攒SQL数据批次,另一边串行往MongoDB插,攒数据不用等插入完成,但插入必须等上一批完事才开始下一批。下面给你具体的实现方案:

核心思路

用线程安全队列缓存数据批次:

  • 生产者(SQL数据收集):异步读取SQL数据,每攒够2000条就丢进队列,继续收集下一批,完全不等待插入操作
  • 消费者(MongoDB插入):循环从队列取批次,每次必须等当前批次插入完成后,再取下一个批次执行,严格保证插入的串行性
完整代码实现
using System.Collections.Concurrent;
using System.Data.SqlClient;
using MongoDB.Driver;

// 替换成你实际的文档类型
public class SqlSourceDocument { /* SQL表对应的字段 */ }
public class MongoTargetDocument { /* MongoDB集合对应的字段 */ }

public class SqlToMongoMigration
{
    private readonly ConcurrentQueue<List<MongoTargetDocument>> _batchQueue = new();
    private readonly IMongoCollection<MongoTargetDocument> _mongoCollection;
    private const int BatchSize = 2000;

    public SqlToMongoMigration(IMongoCollection<MongoTargetDocument> mongoCollection)
    {
        _mongoCollection = mongoCollection;
    }

    // 生产者:异步收集SQL数据,攒够批次就入队
    private async Task CollectSqlBatchesAsync()
    {
        // 替换成你的SQL连接字符串和查询语句
        await using var sqlConn = new SqlConnection("Your_SQL_Connection_String");
        await sqlConn.OpenAsync();
        await using var cmd = new SqlCommand("SELECT * FROM Your_SQL_Table", sqlConn);
        await using var reader = await cmd.ExecuteReaderAsync();

        var currentBatch = new List<MongoTargetDocument>(BatchSize);
        while (await reader.ReadAsync())
        {
            // 把SQL数据映射成Mongo文档(根据你的字段自行实现)
            var mongoDoc = MapSqlToMongo(reader);
            currentBatch.Add(mongoDoc);

            // 攒够2000条就入队,重置批次继续收集
            if (currentBatch.Count == BatchSize)
            {
                _batchQueue.Enqueue(currentBatch);
                currentBatch = new List<MongoTargetDocument>(BatchSize);
            }
        }

        // 处理最后一批不足2000条的数据
        if (currentBatch.Count > 0)
        {
            _batchQueue.Enqueue(currentBatch);
        }

        // 插入结束标记,告诉消费者没有更多批次了
        _batchQueue.Enqueue(null);
    }

    // 消费者:串行处理队列中的批次,上一批插入完成才处理下一批
    private async Task InsertToMongoAsync()
    {
        while (true)
        {
            // 队列空时短暂等待,避免空转浪费资源
            while (!_batchQueue.TryDequeue(out var batch))
            {
                await Task.Delay(100);
            }

            // 遇到结束标记,退出循环
            if (batch == null)
            {
                break;
            }

            // 等待当前批次插入完成,再处理下一个
            await SaveToCollection(batch);
            // 可选:加日志记录批次插入结果
            // Console.WriteLine($"Successfully inserted {batch.Count} documents");
        }
    }

    // 你的MongoDB插入方法(这里包装成异步,保持和原有逻辑兼容)
    private async Task SaveToCollection(List<MongoTargetDocument> batch)
    {
        await _mongoCollection.InsertManyAsync(batch);
    }

    // SQL到Mongo的字段映射方法(根据你的实际表结构实现)
    private MongoTargetDocument MapSqlToMongo(SqlDataReader reader)
    {
        return new MongoTargetDocument
        {
            Id = Guid.NewGuid(),
            // 示例:UserName = reader.GetString(reader.GetOrdinal("UserName")),
            // 其他字段映射...
        };
    }

    // 启动迁移的入口方法
    public async Task RunMigrationAsync()
    {
        // 同时启动生产者和消费者,等待两者都完成
        var collectTask = CollectSqlBatchesAsync();
        var insertTask = InsertToMongoAsync();

        await Task.WhenAll(collectTask, insertTask);
    }
}
关键细节说明
  • ConcurrentQueue:线程安全的队列,保证异步环境下的入队/出队操作不会出现线程冲突,不用自己手动加锁
  • 结束标记:生产者最后插入的null是告诉消费者所有数据都收集完了,避免消费者无限循环
  • 可选优化:如果担心SQL读取太快导致队列积压(占用过多内存),可以加SemaphoreSlim限制队列最大批次数量,比如允许同时存5个批次,避免内存暴涨:
    private readonly SemaphoreSlim _queueLimit = new SemaphoreSlim(5);
    
    // 生产者入队前先等待信号量
    await _queueLimit.WaitAsync();
    _batchQueue.Enqueue(currentBatch);
    
    // 消费者处理完批次后释放信号量
    await SaveToCollection(batch);
    _queueLimit.Release();
    

内容的提问来源于stack exchange,提问作者Krishna

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 09:46:10