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

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();
}

问题分析

  1. 事务未绑定Dapper命令:当前UnitOfWork基于EF Core的ApplicationDbContext开启事务,但Dapper Repository未关联该事务,导致Dapper执行命令时,连接处于事务中但命令未绑定事务,触发异常。
  2. 异步调用未await:var challan = _repos.GetByIdAsync(request.Id, cancellationToken);未使用await,导致后续判断challan.Id会出现空引用或逻辑错误。
  3. SQL语法错误:query1中的INSERT语句格式错误(VALUES内嵌套SELECT的写法违规、多余逗号),且手动添加的BEGIN/END块与UnitOfWork事务冲突。
  4. 事务逻辑混乱:多次调用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 10:07:04