.NET Core WebAPI如何使用多线程并行处理批量数据更新请求
实现方案说明
前置注意事项
- 不要用
Process.GetCurrentProcess().Threads.Count获取可用线程数:该属性返回的是进程当前所有运行的线程总数(包含系统后台线程、IO线程等),不是可用的CPU并行核心数,正确的获取方式是Environment.ProcessorCount,也可以手动指定并行度(比如你需求中的8) - EF Core的DbContext不是线程安全的:绝对不能在多线程/并行任务中共享同一个注入的DbContext实例,每个并行任务必须单独创建服务Scope,从Scope中获取独立的DbContext实例
- 跨线程事务注意事项:默认的数据库事务绑定单个DbContext连接,如果你需要所有数据更新全部成功才统一提交,并行场景下需要用到分布式事务,复杂度较高;如果业务允许单批次更新成功就提交,或者允许部分成功部分回滚,可以按批次单独处理事务
- 原有代码逐行调用
SaveChangesAsync性能极低:每次调用都会产生一次数据库交互,800条数据会产生800次数据库请求,建议单批次处理完统一调用一次SaveChangesAsync
并行实现代码
用.NET Core内置的Parallel.ForEachAsync实现可控并行度的异步处理,同时按需求拆分每个批次处理100条数据:
注意需要先在控制器构造函数中注入
IServiceProvider _serviceProvider
[HttpPatch("UpdateExams")] public async Task<ActionResult<ResponseDto2>> UpdateExam([FromBody] List<Exam> input) { // 指定并行度为8,也可以用Environment.ProcessorCount动态获取CPU核心数 const int parallelism = 8; // 计算每个批次处理的数据量 800/8=100,不足100的单独作为一个批次 int batchSize = (int)Math.Ceiling(input.Count / (double)parallelism); // 把输入数据拆分成指定大小的批次 var batches = input.Chunk(batchSize).ToList(); try { // 并行处理每个批次,指定最大并行度 await Parallel.ForEachAsync(batches, new ParallelOptions { MaxDegreeOfParallelism = parallelism }, async (batch, cancellationToken) => { // 每个并行任务单独创建服务Scope,获取独立的DbContext实例 using var scope = _serviceProvider.CreateScope(); var context = scope.ServiceProvider.GetRequiredService<替换为你的DbContext类名>(); // 每个批次单独开启事务 using var transaction = await context.Database.BeginTransactionAsync(cancellationToken); try { // 先批量查询当前批次所有需要更新的Exam,减少数据库查询次数 var formNos = batch.Select(x => x.FormNo).Distinct().ToList(); var regNos = batch.Select(x => x.RegNo).Distinct().ToList(); var existExams = await context.Exams .Where(e => formNos.Contains(e.FormNo) && regNos.Contains(e.RegNo)) .ToListAsync(cancellationToken); // 匹配数据更新字段 foreach (var item in batch) { var exam = existExams.FirstOrDefault(e => e.FormNo == item.FormNo && e.RegNo == item.RegNo); if (exam != null) { exam.ExamRollno = item.ExamRollno; context.Exams.Update(exam); } } // 批次内所有数据更新完统一提交一次 await context.SaveChangesAsync(cancellationToken); await transaction.CommitAsync(cancellationToken); } catch { await transaction.RollbackAsync(cancellationToken); // 异常往上抛出到外层统一处理 throw; } }); } catch (Exception e) { return StatusCode(StatusCodes.Status500InternalServerError, new ResponseDto2 { Message = "Something went wrong. Please try again later", Success = false, Payload = new { e.StackTrace, e.Message, e.InnerException, e.Source, e.Data } }); } return StatusCode(StatusCodes.Status200OK, new ResponseDto2 { Message = "Data Updated successfully", Success = true, Payload = null }); }
强一致性场景优化方案
如果要求所有数据要么全部更新成功要么全部回滚,不建议使用并行方案,优先优化单线程代码的性能,以下版本性能比你原有代码高数十倍:
[HttpPatch("UpdateExams")] public async Task<ActionResult<ResponseDto2>> UpdateExam([FromBody] List<Exam> input) { using var transaction = await _context.Database.BeginTransactionAsync(); try { // 一次性查询所有需要更新的记录,仅1次数据库查询 var formNos = input.Select(x => x.FormNo).Distinct().ToList(); var regNos = input.Select(x => x.RegNo).Distinct().ToList(); var existExams = await _context.Exams .Where(e => formNos.Contains(e.FormNo) && regNos.Contains(e.RegNo)) .ToListAsync(); foreach (var item in input) { var exam = existExams.FirstOrDefault(e => e.FormNo == item.FormNo && e.RegNo == item.RegNo); if (exam != null) { exam.ExamRollno = item.ExamRollno; _context.Exams.Update(exam); } } // 一次性提交所有更新,仅1次数据库写入 await _context.SaveChangesAsync(); await transaction.CommitAsync(); } catch (Exception e) { await transaction.RollbackAsync(); return StatusCode(StatusCodes.Status500InternalServerError, new ResponseDto2 { Message = "Something went wrong. Please try again later", Success = false, Payload = new { e.StackTrace, e.Message, e.InnerException, e.Source, e.Data } }); } return StatusCode(StatusCodes.Status200OK, new ResponseDto2 { Message = "Data Updated successfully", Success = true, Payload = null }); }
新手学习要点
- 先掌握EF Core的线程安全规则,禁止在并行场景直接共享DbContext
- 了解
Parallel类和Parallel.ForEachAsync的基础用法,掌握并行度参数的配置 - 学习数据库事务的ACID特性,理解本地事务和分布式事务的差异
内容的提问来源于stack exchange,提问作者Frost
相关产品推荐
相关产品推荐

