Clean Architecture中Mediatr+Dapper事务管理异常排查
Clean Architecture下Mediatr+Dapper事务管理异常排查与修复
异常详情
{ "messages": [], "source": "Microsoft.Data.SqlClient.SqlCommand", "exception": "BeginExecuteReader requires the command to have a transaction when the connection assigned to the command is in a pending local transaction. The Transaction property of the command has not been initialized.", "errorId": "6d4329bc-c24f-4852-ba88-35caf4b9ad6d", "supportMessage": "Provide the ErrorId 6d4329bc-c24f-4852-ba88-35caf4b9ad6d to the support team for further analysis.", "statusCode": 500 }
问题背景
在Clean Architecture架构中,结合Mediatr与Dapper实现事务管理时触发上述SQLClient异常。已尝试实现UnitOfWork模式,但问题未解决。期望达成:返回200状态码,提示“challan updated successfully.”,完成Challan的更新及关联表插入操作。
现有实现代码
请求处理类(UpdateChallan相关)
using Mapster; using Microsoft.Data.SqlClient; using API.Application.Common.Custom.IUnitOfWork; using API.Domain.MyAPI.Entities; using System.Data; using System.Transactions; namespace API.Application.MyAPI.Challan.Commands.Update; public class UpdateChallanRequest : IRequest<ChallanDto> { public int Id { get; set; } public decimal Property1 { get; set; } public DateTime Property2 { get; set; } } public class UpdateChallanRequestValidator : CustomValidator { public UpdateChallanRequestValidator() { } } internal class UpdateChallanRequestHandler : IRequestHandler<UpdateChallanRequest, ChallanDto> { private readonly IRepositoryWithEvents<someClass1> _repos; private readonly IDapperRepository _repository; private readonly IUnitOfWork _unitOfWork; private IDbTransaction _transaction; private IStringLocalizer<(someClass1, someClass2)> _localizer; public UpdateChallanRequestHandler(IDapperRepository repository, IUnitOfWork unitOfWork, IStringLocalizer<(someClass1, someClass2)> localizer, IRepositoryWithEvents<someClass1> repos, IDbTransaction transaction) { _repository = repository; _unitOfWork = unitOfWork; _localizer = localizer; _repos = repos; _transaction = transaction; } public async Task<ChallanDto> Handle(UpdateChallanRequest request, CancellationToken cancellationToken) { _unitOfWork.BeginTransaction(); var challan = _repos.GetByIdAsync(request.Id, cancellationToken); string query; string query1; var parameters = new { ID = request.Id, param1 = request.Property1, param2 = request.Property2 }; query = "BEGIN" + "UPDATE dbo.someTable1" + "SET someColumn1=@param1, Active=0, someColumn2=@param2, flagPaid=true" + "WHERE dbo.someTable1.id = @ID;" + "END"; var result = await _repository.QueryAsync<someClass1>(query, parameters, cancellationToken: cancellationToken); await _unitOfWork.SaveChangesAsync(); try { if(challan.Id == null) { throw new ArgumentNullException(); } else if (request.Id != null && challan.Id != null) { query1 = "BEGIN" + "INSERT INTO dbo.someTable2" + "(someColumn1, someColumn2, someColumn3)" + "VALUES (" + "SELECT someColumn1 from dbo.someTable1 WHERE dbo.someTable1.id = @ID;,", "SELECT someColumn2 from dbo.someTable1 WHERE dbo.someTable1.id = @ID;,", "SELECT someColumn3 from dbo.someTable1 WHERE dbo.someTable1.id = @ID;,", "END"; await _repository.QueryAsync<someClass1>(query1, parameters, cancellationToken: cancellationToken); await _unitOfWork.SaveChangesAsync(); _unitOfWork.Commit(); } } catch(Exception ex) { _unitOfWork.Rollback(); throw new ArgumentNullException("Challan Not Found."); } return result.Adapt<ChallanDto>(); } }
UnitOfWork实现
public class UnitOfWork : IUnitOfWork { private readonly ApplicationDbContext _context; private IDbContextTransaction? _transaction; public UnitOfWork(ApplicationDbContext context) { _context = context; } public void BeginTransaction() { if (_transaction != null) { return; } _transaction = _context.Database.BeginTransaction(); } public Task<int> SaveChangesAsync() { return _context.SaveChangesAsync(); } public void Commit() { if (_transaction == null) { return; } _transaction.Commit(); _transaction.Dispose(); _transaction = null; } public async Task SaveAndCommitAsync() { await SaveChangesAsync(); Commit(); } public void Rollback() { if (_transaction == null) { return; } _transaction.Rollback(); _transaction.Dispose(); _transaction = null; } public void Dispose() { if (_transaction == null) { return; } _transaction.Dispose(); _transaction = null; } } public interface IUnitOfWork : IDisposable, ITransientService { void BeginTransaction(); void Commit(); void Rollback(); Task<int> SaveChangesAsync(); Task SaveAndCommitAsync(); }
问题分析
- 事务未绑定Dapper命令:当前UnitOfWork基于EF Core的
ApplicationDbContext开启事务,但Dapper Repository未关联该事务,导致Dapper执行命令时,连接处于事务中但命令未绑定事务,触发异常。 - 异步调用未await:
var challan = _repos.GetByIdAsync(request.Id, cancellationToken);未使用await,导致后续判断challan.Id会出现空引用或逻辑错误。 - SQL语法错误:query1中的INSERT语句格式错误(VALUES内嵌套SELECT的写法违规、多余逗号),且手动添加的BEGIN/END块与UnitOfWork事务冲突。
- 事务逻辑混乱:多次调用
SaveChangesAsync,Commit仅在分支中执行,异常处理范围覆盖不全,可能导致事务未正确回滚。
修复方案
1. 改造UnitOfWork,暴露事务供Dapper使用
修改UnitOfWork类,添加获取当前事务的属性:
public class UnitOfWork : IUnitOfWork { // 原有代码不变 public IDbTransaction? CurrentTransaction => _transaction?.GetDbTransaction(); }
2. 调整Dapper Repository支持事务参数
确保IDapperRepository的QueryAsync方法接受事务参数,内部执行时绑定事务:
public interface IDapperRepository { Task<IEnumerable<T>> QueryAsync<T>(string sql, object param = null, IDbTransaction transaction = null, CancellationToken cancellationToken = default); // 其他方法... } public class DapperRepository : IDapperRepository { private readonly IDbConnectionFactory _connectionFactory; public DapperRepository(IDbConnectionFactory connectionFactory) { _connectionFactory = connectionFactory; } public async Task<IEnumerable<T>> QueryAsync<T>(string sql, object param = null, IDbTransaction transaction = null, CancellationToken cancellationToken = default) { using var connection = _connectionFactory.GetConnection(); return await connection.QueryAsync<T>(sql, param, transaction: transaction, cancellationToken: cancellationToken); } }
3. 修复Handler逻辑
internal class UpdateChallanRequestHandler : IRequestHandler<UpdateChallanRequest, ChallanDto> { private readonly IRepositoryWithEvents<someClass1> _repos; private readonly IDapperRepository _repository; private readonly IUnitOfWork _unitOfWork; private readonly IStringLocalizer<(someClass1, someClass2)> _localizer; // 移除不必要的IDbTransaction构造注入 public UpdateChallanRequestHandler(IDapperRepository repository, IUnitOfWork unitOfWork, IStringLocalizer<(someClass1, someClass2)> localizer, IRepositoryWithEvents<someClass1> repos) { _repository = repository; _unitOfWork = unitOfWork; _localizer = localizer; _repos = repos; } public async Task<ChallanDto> Handle(UpdateChallanRequest request, CancellationToken cancellationToken) { _unitOfWork.BeginTransaction(); try { // 补全await var challan = await _repos.GetByIdAsync(request.Id, cancellationToken); if(challan == null || challan.Id == null) { throw new ArgumentNullException(nameof(challan), "Challan Not Found."); } var parameters = new { ID = request.Id, param1 = request.Property1, param2 = request.Property2 }; // 修正更新SQL,移除BEGIN/END var updateQuery = @" UPDATE dbo.someTable1 SET someColumn1=@param1, Active=0, someColumn2=@param2, flagPaid=true WHERE dbo.someTable1.id = @ID;"; // 传入UnitOfWork事务 await _repository.QueryAsync<someClass1>(updateQuery, parameters, _unitOfWork.CurrentTransaction, cancellationToken); // 修正插入SQL,改用INSERT...SELECT语法 var insertQuery = @" INSERT INTO dbo.someTable2(someColumn1, someColumn2, someColumn3) SELECT someColumn1, someColumn2, someColumn3 FROM dbo.someTable1 WHERE id = @ID;"; await _repository.QueryAsync<someClass1>(insertQuery, parameters, _unitOfWork.CurrentTransaction, cancellationToken); // 统一保存变更并提交事务 await _unitOfWork.SaveChangesAsync(); _unitOfWork.Commit(); // 获取最新数据返回 var updatedChallan = await _repos.GetByIdAsync(request.Id, cancellationToken); return updatedChallan.Adapt<ChallanDto>(); } catch(Exception ex) { _unitOfWork.Rollback(); throw new InvalidOperationException("Challan update failed.", ex); } } }
验证结果
- Dapper命令会绑定UnitOfWork开启的事务,解决原异常
- 所有数据库操作原子性执行:要么全部成功提交,要么异常时全部回滚
- 存在的Challan会完成更新和关联表插入,返回200状态码及成功提示
- 不存在的Challan会触发异常,回滚事务并返回对应错误信息
内容的提问来源于stack exchange,提问作者SJI
相关产品推荐
相关产品推荐

