并行数据库操作中如何恢复事务超时功能?
问题核心
为提升批量数据库操作性能,将原有异步代码改为同步并通过Parallel.ForEach并行化,但同步版EF Core事务方法(如BeginTransaction)不支持传入CancellationToken,导致事务超时控制功能丢失。尝试过定时器、Task包装方案均不满意,改用Parallel.ForEachAsync后仍未正确实现单迭代的超时控制。
现有代码摘要
Parallel.ForEach(excelRawList.SuccessRows, options, (excelRow) => { using (var dbContext = new MyContext(new DbContextOptionsBuilder<MyContext>().UseSqlServer(_configuration.GetConnectionString ("connectionstring")).Options)) { using (var unitOfWork = new UnitOfWork(dbContext)) { unitOfWork.BeginTransaction(); try { int? decision = CanBeInserted(unitOfWork, excelRow.id); if (decision != 0) { if (!excelRow.Store) { try { if (decision != null) { service3.MoveToBackup(unitOfWork, decision.Value); } Entity2 entity2ToBeInserted = service1.changeExcelRowToEntity2(excelRow); excelRow.Entity2 = entity2ToBeInserted; unitOfWork.Context.Set<Entity2>().Add(entity2ToBeInserted); unitOfWork.Context.Set<Entity1>().Add(excelRow); unitOfWork.SaveChanges(); unitOfWork.CommitTransaction(); } catch (OperationCanceledException ex) { unitOfWork.Rollback(); excelRawList.ErrorRows.Add(new ExcelError() {"Timeout ended"}); } catch (Exception ex) { unitOfWork.Rollback(); excelRawList.ErrorRows.Add(new ExcelError() {"Exception"}); } } else { excelRow.FK_Entity1_Entity2 = null; unitOfWork.Context.Set<Entity1>().Add(excelRow); unitOfWork.SaveChanges(); unitOfWork.CommitTransaction(); } } else { excelRawList.ErrorRows.Add(new ExcelError() {"Repeated"}); } } catch (Exception ex) { excelRawList.ErrorRows.Add(new ExcelError() {"Exception"}); } } } });
尝试过的方案及问题
- 定时器、Task包装同步操作:实现复杂且性能损耗大,不满意。
Parallel.ForEachAsync:外部创建CancellationTokenSource后,未在迭代内结合单超时逻辑,导致无法实现每个数据库操作独立超时。
解决方案思路
方案1:异步并行+单迭代独立超时(推荐)
放弃同步操作,改用异步并行,给每个迭代设置独立超时令牌,同时关联全局取消信号,确保每个数据库操作都响应超时。
关键步骤:
- 调整
UnitOfWork实现异步事务方法,支持传入CancellationToken。 - 在
Parallel.ForEachAsync的每个迭代内,创建独立超时令牌,与传入的全局令牌关联。 - 所有数据库操作(事务、保存、业务方法)改用异步版本并传入联合令牌。
- 使用线程安全集合存储错误行,避免并发冲突。
示例代码:
// 全局取消令牌(可选,用于手动终止整体操作) var globalCts = new CancellationTokenSource(); // 每个迭代的超时时间 var perIterationTimeout = TimeSpan.FromSeconds(30); // 线程安全集合存储错误行,替代原非线程安全的ErrorRows var errorRows = new ConcurrentBag<ExcelError>(); await Parallel.ForEachAsync(excelRawList.SuccessRows, new ParallelOptions { MaxDegreeOfParallelism = Environment.ProcessorCount * 2 // 控制并行度,避免数据库连接耗尽 }, async (excelRow, globalToken) => { using (var dbContext = new MyContext(new DbContextOptionsBuilder<MyContext>() .UseSqlServer(_configuration.GetConnectionString("connectionstring")) .Options)) using (var unitOfWork = new UnitOfWork(dbContext)) { // 创建当前迭代的超时令牌,关联全局令牌 using var iterationCts = CancellationTokenSource.CreateLinkedTokenSource(globalToken); iterationCts.CancelAfter(perIterationTimeout); var combinedToken = iterationCts.Token; try { // 异步启动事务,传入令牌 await unitOfWork.BeginTransactionAsync(combinedToken); // 将同步判断改为异步(若涉及数据库操作) int? decision = await CanBeInsertedAsync(unitOfWork, excelRow.id, combinedToken); if (decision != 0) { if (!excelRow.Store) { if (decision != null) { // 异步执行业务操作,传入令牌 await service3.MoveToBackupAsync(unitOfWork, decision.Value, combinedToken); } Entity2 entity2ToBeInserted = service1.changeExcelRowToEntity2(excelRow); excelRow.Entity2 = entity2ToBeInserted; unitOfWork.Context.Set<Entity2>().Add(entity2ToBeInserted); unitOfWork.Context.Set<Entity1>().Add(excelRow); // 异步保存,传入令牌 await unitOfWork.SaveChangesAsync(combinedToken); await unitOfWork.CommitTransactionAsync(combinedToken); } else { excelRow.FK_Entity1_Entity2 = null; unitOfWork.Context.Set<Entity1>().Add(excelRow); await unitOfWork.SaveChangesAsync(combinedToken); await unitOfWork.CommitTransactionAsync(combinedToken); } } else { errorRows.Add(new ExcelError { Message = "Repeated" }); } } catch (OperationCanceledException) { await unitOfWork.RollbackAsync(); errorRows.Add(new ExcelError { Message = "Timeout ended" }); } catch (Exception ex) { await unitOfWork.RollbackAsync(); errorRows.Add(new ExcelError { Message = $"Exception: {ex.Message}" }); } } }); // 将错误行合并到原列表(如果需要) foreach (var error in errorRows) { excelRawList.ErrorRows.Add(error); }
配套UnitOfWork异步实现:
public class UnitOfWork { private readonly MyContext _context; private IDbContextTransaction? _transaction; public UnitOfWork(MyContext context) { _context = context; } public async Task BeginTransactionAsync(CancellationToken cancellationToken = default) { _transaction = await _context.Database.BeginTransactionAsync(cancellationToken); } public async Task CommitTransactionAsync(CancellationToken cancellationToken = default) { if (_transaction == null) return; await _transaction.CommitAsync(cancellationToken); await _transaction.DisposeAsync(); } public async Task RollbackAsync() { if (_transaction == null) return; await _transaction.RollbackAsync(); await _transaction.DisposeAsync(); } public async Task SaveChangesAsync(CancellationToken cancellationToken = default) { await _context.SaveChangesAsync(cancellationToken); } }
方案2:同步并行+Task包装超时(兼容旧代码)
若必须保留同步操作,可在每个迭代内用Task.Run包装同步代码,并设置超时令牌,通过Task.Wait等待并捕获超时异常。
示例代码:
Parallel.ForEach(excelRawList.SuccessRows, options, (excelRow) => { var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); try { Task.Run(() => { using (var dbContext = new MyContext(new DbContextOptionsBuilder<MyContext>().UseSqlServer(_configuration.GetConnectionString("connectionstring")).Options)) using (var unitOfWork = new UnitOfWork(dbContext)) { unitOfWork.BeginTransaction(); try { int? decision = CanBeInserted(unitOfWork, excelRow.id); if (decision != 0) { if (!excelRow.Store) { if (decision != null) { service3.MoveToBackup(unitOfWork, decision.Value); } Entity2 entity2ToBeInserted = service1.changeExcelRowToEntity2(excelRow); excelRow.Entity2 = entity2ToBeInserted; unitOfWork.Context.Set<Entity2>().Add(entity2ToBeInserted); unitOfWork.Context.Set<Entity1>().Add(excelRow); unitOfWork.SaveChanges(); unitOfWork.CommitTransaction(); } else { excelRow.FK_Entity1_Entity2 = null; unitOfWork.Context.Set<Entity1>().Add(excelRow); unitOfWork.SaveChanges(); unitOfWork.CommitTransaction(); } } else { // 加锁确保线程安全 lock (excelRawList.ErrorRows) { excelRawList.ErrorRows.Add(new ExcelError() { Message = "Repeated" }); } } } catch (Exception ex) { unitOfWork.Rollback(); lock (excelRawList.ErrorRows) { excelRawList.ErrorRows.Add(new ExcelError() { Message = "Exception" }); } } } }, cts.Token).Wait(cts.Token); } catch (AggregateException ae) { ae.Handle(ex => { if (ex is OperationCanceledException) { lock (excelRawList.ErrorRows) { excelRawList.ErrorRows.Add(new ExcelError() { Message = "Timeout ended" }); } return true; } lock (excelRawList.ErrorRows) { excelRawList.ErrorRows.Add(new ExcelError() { Message = "Exception" }); } return true; }); } });
注意点:
- 对共享集合
excelRawList.ErrorRows添加操作必须加锁,避免并发冲突。 - 此方案性能略低于异步并行,且线程调度开销更大,仅作为兼容旧代码的备选。
内容的提问来源于stack exchange,提问作者Kenzo_Gilead
相关产品推荐
相关产品推荐

