无法用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; } }
错误根源
- 异步与同步并行框架冲突:
Parallel.ForEach是同步并行模型,但代码中大量使用await异步调用,会导致线程上下文切换,同一个DbContext可能被多个线程复用,触发并发操作错误。 - 异步void方法隐患:
MoveToHistoric声明为void但内部使用await,属于异步void方法,调用方无法等待其执行完成,会导致DbContext后续操作与未完成的异步操作冲突。 - 共享Service的DbContext问题:
service1、service2等若为全局共享实例,且内部持有独立DbContext,会导致多个并行线程共用同一个DbContext实例。 - 线程不安全集合操作:
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
相关产品推荐
相关产品推荐

