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

.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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 12:45:09