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

多线程并行使用Dapper从Azure SQL取数的报错解决方法

解决Azure SQL并行取数时的TCP传输错误

咱们先来拆解你遇到的问题:部署到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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.09 08:37:28