多线程并行使用Dapper从Azure SQL取数的报错解决方法
咱们先来拆解你遇到的问题:部署到K8S后出现的WriteAsync冲突错误,本质是多个异步操作共享了同一个底层网络连接资源,再加上你当前处理异步任务的方式有问题,甚至还踩了闭包陷阱和线程安全的坑。下面一步步给你梳理解决方案:
1. 先把异步任务的处理逻辑掰正
你当前的代码是把异步任务塞进列表后,通过foreach逐个调用task.Result——这会导致任务串行执行(Result会阻塞当前线程直到任务完成),完全没用到并行的优势,还可能引发线程阻塞相关的问题。正确的并行打开方式应该是这样:
// 1. 先创建所有异步任务(此时任务开始调度,但不会立即阻塞) var tasks = new List<Task<(IEnumerable<int> unprocessedData, IEnumerable<dynamic> rowData)>>(); for (int i = 0; i < numberOfPages; i++) { var currentStartRow = startRow; // 重点!捕获当前循环的行号,避免闭包陷阱 // 每个任务用独立的SqlBuilder,避免线程安全问题(后面会讲) var taskBuilder = new SqlBuilder(); // 把原builder的where/join/orderby等条件复制到新实例里(根据你的业务逻辑调整) tasks.Add(_tableviewRepository.GetTableviewRowsByPagination( tableviewExportCondition.TableviewName, modelMappingGroups, currentStartRow, taskBuilder, pageSize, appName, i)); startRow += pageSize; } // 2. 真正并行执行所有任务,等待全部完成 var results = await Task.WhenAll(tasks); // 3. 统一处理结果 foreach(var result in results) { if (result.rowData != null) { dataToExport.AddRange(result.rowData); } // 顺便收集失败的批次,方便后续重试 if (result.unprocessedData.Any()) { // 比如把这些行号存起来后续处理 } }
这里一定要注意闭包陷阱:如果直接在循环里用startRow,所有任务会捕获同一个变量的引用,最终所有任务的startRow都会变成循环最后一次的值,导致取数逻辑完全错误。
2. 解决SqlBuilder的线程安全问题
你遇到的WriteAsync错误,很大概率和共享SqlBuilder有关——SqlBuilder内部维护了SQL片段的集合,不是线程安全的!多个并行任务同时修改同一个builder(比如调用builder.Select()),会导致SQL拼接混乱,甚至引发底层的线程安全异常。
所以必须给每个并行任务分配独立的SqlBuilder实例,就像上面代码里那样,不要在循环外复用同一个builder。
3. 控制并行请求的数量,避免连接池耗尽
Azure SQL Server的连接池默认大小是100,你要取46万条数据(46个批次),如果直接并行46个请求,很容易超出连接池的承载能力,引发连接竞争和传输层错误。
可以用SemaphoreSlim来限制并行数量,比如最多同时跑10个任务:
var semaphore = new SemaphoreSlim(10); // 最多10个并行任务 var tasks = new List<Task<(IEnumerable<int> unprocessedData, IEnumerable<dynamic> rowData)>>(); for (int i = 0; i < numberOfPages; i++) { var currentStartRow = startRow; var taskBuilder = new SqlBuilder(); // 复制原builder的条件... tasks.Add(Task.Run(async () => { await semaphore.WaitAsync(); try { return await _tableviewRepository.GetTableviewRowsByPagination( tableviewExportCondition.TableviewName, modelMappingGroups, currentStartRow, taskBuilder, pageSize, appName, i); } finally { semaphore.Release(); // 不管成功失败都要释放信号量 } })); startRow += pageSize; } var results = await Task.WhenAll(tasks); // 后续处理...
这样能避免瞬间发起大量连接请求,减轻SQL Server的压力,也能避免连接池被耗尽。
4. 优化数据库连接的使用,彻底避免资源共享
你的GetTableviewRowsByPagination方法里,依赖_unitOfWork获取连接的方式可能有问题——如果_unitOfWorkServices.Build()返回的实例复用了同一个SqlConnection,就会导致多个任务共享连接,引发WriteAsync冲突。
建议直接在方法内创建独立的连接,更可控:
public async Task<(IEnumerable<int> unprocessedData, IEnumerable<dynamic> rowData)> GetTableviewRowsByPagination(...) { List<int> unprocessedData = new List<int>(); // 直接从配置获取对应应用的连接字符串 var connectionString = _configuration.GetConnectionString(appName.ToString()); // 每个任务用独立的连接,using块自动释放资源 using (var connection = new SqlConnection(connectionString)) { try { var columns = tableviewAttributeDetails.Select(c => $"{c.mapping_group_value} [{c.attribute}]"); var joinedColumn = string.Join(",", columns); builder.Select(joinedColumn); var selector = builder.AddTemplate($"SELECT /**select**/ FROM {tableName} with (nolock) /**innerjoin**/ /**where**/ /**orderby**/ OFFSET {startRow} ROWS FETCH NEXT {(pageSize == 0 ? 100 : pageSize)} ROWS ONLY"); // Dapper的QueryAsync会自动打开连接,不用手动调用Open() var data = await connection.QueryAsync(selector.RawSql, selector.Parameters); Console.WriteLine($"data completed for task{taskNumber}"); return (unprocessedData, data); } catch(Exception ex) { Console.WriteLine($"Exception: {ex.Message}"); if (ex.InnerException != null) Console.WriteLine($"InnerException: {ex.InnerException.Message}"); Console.WriteLine($"Error in fetching from row {startRow}"); unprocessedData.Add(startRow); return (unprocessedData, null); } } }
这样每个任务都用完全独立的连接,彻底杜绝了资源共享的问题。
5. 为什么加锁后系统崩溃?
你之前在方法内加锁,会导致所有任务变成串行执行,而且如果锁的范围太大(比如整个方法),很容易引发线程阻塞死锁,再加上K8S环境的资源限制(内存、CPU不足),最终导致系统崩溃。加锁不是解决这个问题的正确姿势,咱们要的是安全的并行,不是完全串行。
内容的提问来源于stack exchange,提问作者Purushothaman

