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

并行数据库操作中如何恢复事务超时功能?

问题核心

为提升批量数据库操作性能,将原有异步代码改为同步并通过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"});
            }
        }
    }
});
尝试过的方案及问题
  1. 定时器、Task包装同步操作:实现复杂且性能损耗大,不满意。
  2. Parallel.ForEachAsync:外部创建CancellationTokenSource后,未在迭代内结合单超时逻辑,导致无法实现每个数据库操作独立超时。
解决方案思路

方案1:异步并行+单迭代独立超时(推荐)

放弃同步操作,改用异步并行,给每个迭代设置独立超时令牌,同时关联全局取消信号,确保每个数据库操作都响应超时。

关键步骤:

  1. 调整UnitOfWork实现异步事务方法,支持传入CancellationToken。
  2. 在Parallel.ForEachAsync的每个迭代内,创建独立超时令牌,与传入的全局令牌关联。
  3. 所有数据库操作(事务、保存、业务方法)改用异步版本并传入联合令牌。
  4. 使用线程安全集合存储错误行,避免并发冲突。

示例代码:

// 全局取消令牌(可选,用于手动终止整体操作)
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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 19:42:33