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

无法用Parallel库并行化代码,DbContext多线程冲突求助

并行化EF Core数据库插入操作时遇到DbContext并发错误的解决求助

业务逻辑

  • 若待插入记录已存在且处于错误状态:将旧记录移至历史表后插入新记录
  • 若已存在但非错误状态:丢弃新记录
  • 若不存在:直接插入新记录

问题现象

为实现并行化,我已为每个线程创建独立的DbContext和UnitOfWork,但仍持续报错:
A second operation was started on this context instance before a previous operation completed. This is usually caused by different threads concurrently using the same instance of DbContext.

无法定位问题根源,求解决方案。

相关代码

ParallelOptions options = new ParallelOptions();
options.MaxDegreeOfParallelism = 2;

Parallel.ForEach(excelRawList.SuccessRows, options, (excelRow) =>
{
    using (var dbContext = new MyContext(
    new DbContextOptionsBuilder<MyContext>()
.UseSqlServer(_configuration.GetConnectionString("MyDB"))
            .Options))
    {
        using (var unitOfWork = new UnitOfWork(dbContext))
        {
            try
            {
                int? canBeInserted = await service3.CanBeInserted(excelRow.Id);

                // canBeIserted 有三种取值:null(无前置记录)、0(有前置记录且非错误状态)、Id值(有前置记录且为错误状态)

                if (canBeInserted != 0)
                {
                    if (!excelRow.Prorroga)
                    {
                        using (var cts = new CancellationTokenSource(TimeSpan.FromMilliseconds(_timeoutTransaccion)))
                        {
                            unitOfWork.BeginTransaction(cts);

                            if (canBeInserted != null)
                            {
                                service4.MoveToHistoric(unitOfWork, canBeInserted.Value, cts);
                            }

                            Entity2 entity2 = service3.Entity1ToEntity2(excelRow);
                            excelRow.entity2 = entity2;

                            unitOfWork.Context.Set<Entity2>().Add(entity2);
                            unitOfWork.Context.Set<Entity1>().Add(excelRow);

                            unitOfWork.SaveChanges(cts);
                            unitOfWork.CommitTransaction(cts);

                            if (entity2.Age == null)
                            {
                               _logger.LogWarning("未找到Age字段");
                            }
                        }
                    }
                    catch (OperationCanceledException ex)
                    {
                        unitOfWork.Rollback();
                    }
                    catch (Exception ex)
                    {
                        unitOfWork.Rollback();
                    }
              }
              else
              {
                   excelRow.FK2= null;
                   _service1.Add(excelRow);
              }
         }
         else
         {
               excelRawList.ErrorRows.Add(new ExcelError() { 
         }
      }
      catch (Exception ex)
      {
          excelRawList.ErrorRows.Add(new ExcelError()
          {
               Message = ex.InnerException.Message
          });
      }
    }
  }
});

public void MoveToHistoric(UnitOfWork unitOfWork, int Id, CancellationTokenSource cts)
{
    try
    {
        Entity1 entity1 = service1.Get(id);
        Entity2 entity2 = service2.Get(entity1.FK);
        Entity3 entity3 = _mapper.Map<Entity3>(entity1);
        Entity4 entity4 = _mapper.Map<Entity4>(entity2);

        unitOfWork.Context.Set<Entity1>().Remove(entity1);
        unitOfWork.Context.Set<Entity2>().Remove(entity2);

        await unitOfWork.SaveChangesAsyncLimite(cts);

        entity3.Entity4 = entity4;
        unitOfWork.Context.Set<Entity4>().Add(entity4);
        unitOfWork.Context.Set<Entity3>().Add(entity3);

        unitOfWork.SaveChanges(cts);

        return unitOfWork;
    }
    catch (Exception ex)
    {
        throw;
    }
}

错误根源

  1. 异步与同步并行框架冲突:Parallel.ForEach是同步并行模型,但代码中大量使用await异步调用,会导致线程上下文切换,同一个DbContext可能被多个线程复用,触发并发操作错误。
  2. 异步void方法隐患:MoveToHistoric声明为void但内部使用await,属于异步void方法,调用方无法等待其执行完成,会导致DbContext后续操作与未完成的异步操作冲突。
  3. 共享Service的DbContext问题:service1、service2等若为全局共享实例,且内部持有独立DbContext,会导致多个并行线程共用同一个DbContext实例。
  4. 线程不安全集合操作:excelRawList.ErrorRows.Add是线程不安全的,多线程同时操作会引发集合异常,也可能干扰DbContext的状态判断。

修复方案

1. 替换为异步并行模型

用Task.WhenAll替代Parallel.ForEach,适配异步操作逻辑:

var tasks = excelRawList.SuccessRows.Select(async excelRow =>
{
    using (var dbContext = new MyContext(
        new DbContextOptionsBuilder<MyContext>()
            .UseSqlServer(_configuration.GetConnectionString("MyDB"))
            .Options))
    {
        using (var unitOfWork = new UnitOfWork(dbContext))
        {
            try
            {
                int? canBeInserted = await service3.CanBeInserted(excelRow.Id);

                if (canBeInserted != 0)
                {
                    if (!excelRow.Prorroga)
                    {
                        using (var cts = new CancellationTokenSource(TimeSpan.FromMilliseconds(_timeoutTransaccion)))
                        {
                            unitOfWork.BeginTransaction(cts.Token);

                            if (canBeInserted != null)
                            {
                                await service4.MoveToHistoric(unitOfWork, canBeInserted.Value, cts.Token);
                            }

                            Entity2 entity2 = service3.Entity1ToEntity2(excelRow);
                            excelRow.entity2 = entity2;

                            unitOfWork.Context.Set<Entity2>().Add(entity2);
                            unitOfWork.Context.Set<Entity1>().Add(excelRow);

                            await unitOfWork.SaveChangesAsync(cts.Token);
                            unitOfWork.CommitTransaction(cts.Token);

                            if (entity2.Age == null)
                            {
                                _logger.LogWarning("未找到Age字段");
                            }
                        }
                    }
                }
                else
                {
                    excelRow.FK2 = null;
                    // 若_service1是共享实例,需确保Add方法线程安全,或为每个任务创建独立实例
                    await _service1.AddAsync(excelRow);
                }
            }
            catch (OperationCanceledException ex)
            {
                unitOfWork.Rollback();
            }
            catch (Exception ex)
            {
                unitOfWork.Rollback();
                // 线程安全地添加错误记录
                lock (excelRawList.ErrorRows)
                {
                    excelRawList.ErrorRows.Add(new ExcelError()
                    {
                        Message = ex.InnerException?.Message ?? ex.Message
                    });
                }
            }
        }
    }
});

await Task.WhenAll(tasks);

2. 修复MoveToHistoric方法

将其改为异步非void方法,确保调用方能等待操作完成:

public async Task MoveToHistoric(UnitOfWork unitOfWork, int id, CancellationToken token)
{
    try
    {
        // 确保service的Get方法使用当前UnitOfWork的DbContext,而非独立实例
        Entity1 entity1 = await service1.GetAsync(id, token);
        Entity2 entity2 = await service2.GetAsync(entity1.FK, token);
        
        Entity3 entity3 = _mapper.Map<Entity3>(entity1);
        Entity4 entity4 = _mapper.Map<Entity4>(entity2);

        unitOfWork.Context.Set<Entity1>().Remove(entity1);
        unitOfWork.Context.Set<Entity2>().Remove(entity2);

        await unitOfWork.SaveChangesAsync(token);

        entity3.Entity4 = entity4;
        unitOfWork.Context.Set<Entity4>().Add(entity4);
        unitOfWork.Context.Set<Entity3>().Add(entity3);

        await unitOfWork.SaveChangesAsync(token);
    }
    catch (Exception ex)
    {
        throw;
    }
}

3. 确保Service线程安全

  • 若service1、service2为单例,需移除内部持有的DbContext,所有数据库操作均使用传入的UnitOfWork的DbContext。
  • 在依赖注入场景(如ASP.NET Core)中,将Service设置为Scoped生命周期,确保每个任务获取独立的Service实例。

4. 线程安全处理错误集合

将excelRawList.ErrorRows替换为ConcurrentBag<ExcelError>,或在添加时使用锁保证线程安全。

内容的提问来源于stack exchange,提问作者Kenzo_Gilead

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 22:08:12