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
相关产品推荐
相关产品推荐

