如何用C#多线程处理SQL Server十万条数据并调用外部API更新表?
多线程改造方案
针对你的场景,下面提供两种易上手的多线程改造方案,同时给出关键注意事项,避免踩坑:
方案一:用Parallel.For快速改造(适合新手)
直接替换原有的foreach为Parallel.ForEach,同时控制并发数,避免触发API限流或数据库压力过大:
public void MainProcess() { try { // 注意:如果10万条数据一次性加载内存压力大,建议改成分页读取(见方案二) List<rowResult> rowResults = (List<rowResult>)GetRowsFromTable(); // 配置并行参数,MaxDegreeOfParallelism根据API限流和服务器性能调整,比如设为8 var parallelOpts = new ParallelOptions { MaxDegreeOfParallelism = 8 }; Parallel.ForEach(rowResults, parallelOpts, row => { try { callExternalAPI(row); } catch (Exception ex) { // 单独处理单条记录的异常,比如写日志,不要抛出导致整个并行流程中断 Console.WriteLine($"处理记录ID[{row.Id}]失败:{ex.Message}"); } }); } catch (Exception ex) { Console.WriteLine($"主流程初始化异常:{ex.Message}"); throw; } }
方案二:分页读取+异步Task(更灵活,内存友好)
一次性加载10万条数据会占用大量内存,建议分页读取分批处理,同时用异步API调用提升效率:
// 改成异步方法,提升整体吞吐量 public async Task MainProcessAsync() { try { const int pageSize = 1000; // 每批处理1000条,可调整 int pageIndex = 1; List<rowResult> currentBatch; do { // 改造GetRowsFromTable,支持分页查询(传入页码和每页条数) currentBatch = (List<rowResult>)GetRowsFromTable(pageIndex, pageSize); if (!currentBatch.Any()) break; // 为当前批次的每条记录创建异步任务 var taskList = new List<Task>(); foreach (var row in currentBatch) { taskList.Add(ProcessSingleRowAsync(row)); } // 等待当前批次所有任务完成后,再取下一批 await Task.WhenAll(taskList); pageIndex++; } while (currentBatch.Count == pageSize); } catch (Exception ex) { Console.WriteLine($"主流程异常:{ex.Message}"); throw; } } // 单独封装单条记录的异步处理逻辑 private async Task ProcessSingleRowAsync(rowResult row) { try { // 把原有的callExternalAPI改成异步版本(如果API支持异步) var apiResult = await CallExternalAPIAsync(row); // 异步更新数据库,每个任务用独立的数据库连接 using (var conn = new SqlConnection("你的数据库连接字符串")) { await conn.OpenAsync(); var updateCmd = new SqlCommand("UPDATE 你的表名 SET 结果字段 = @result WHERE Id = @id", conn); updateCmd.Parameters.AddWithValue("@result", apiResult); updateCmd.Parameters.AddWithValue("@id", row.Id); await updateCmd.ExecuteNonQueryAsync(); } } catch (Exception ex) { Console.WriteLine($"处理记录ID[{row.Id}]失败:{ex.Message}"); } } // 异步调用外部API的示例 private async Task<string> CallExternalAPIAsync(rowResult row) { using (var httpClient = new HttpClient()) { // 替换为你的API地址和参数 var response = await httpClient.GetAsync($"https://your-api-url.com/{row.Id}"); response.EnsureSuccessStatusCode(); return await response.Content.ReadAsStringAsync(); } }
关键注意事项
- 控制并发数:不管用哪种方案,都要限制并发量,否则会触发API限流,或者耗尽数据库连接池。比如Parallel.ForEach的MaxDegreeOfParallelism,或者分批处理时控制每批的任务数。
- 数据库连接隔离:绝对不要在多个线程之间共享SqlConnection实例,每个任务必须创建自己的连接(用using语句自动释放)。
- 异常隔离处理:每条记录的处理异常要单独捕获,不要因为一条记录失败导致整个流程终止。
- 异步优先:如果外部API支持异步请求,一定要用异步方法,比同步阻塞的多线程效率高很多。
- API限流应对:如果外部API有调用频率限制,要调整并发数,或者添加重试/延迟逻辑(比如用Polly库实现重试)。
内容的提问来源于stack exchange,提问作者vambat
相关产品推荐
相关产品推荐

