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

.NET Core WebAPI多数据库跨API分布式事务回滚实现咨询

跨API双MySQL库事务回滚方案(基于ApiTransactionScope思路)

核心原理:编排式分布式事务(TCC模型)

ApiTransactionScope本质是对编排式TCC事务的封装,通过一个协调角色(网关或独立协调服务)控制各个API的三个阶段:

  1. Try(预执行):执行数据修改但不提交,保留回滚能力
  2. Confirm(提交):确认所有预执行成功后,正式提交事务
  3. Cancel(回滚):任意环节失败,撤销所有已完成的预执行操作

具体实现步骤(适配你的.NET Core + MySQL场景)

1. 为两个API扩展事务控制接口

由于你直接操作MySQL(未用EF),需要给每个API新增三组接口,手动管理数据库事务:

  • 预执行接口:执行数据修改,开启事务但不提交,关联全局事务ID存储
  • 提交接口:根据全局事务ID提交指定事务
  • 回滚接口:根据全局事务ID回滚指定事务

示例代码(API1的MySQL事务操作):

// 预执行:修改DB1数据,挂起事务
[HttpPost("prepare-update")]
public IActionResult PrepareUpdate([FromBody] UpdateDto req, [FromQuery] string txId)
{
    using var conn = new MySqlConnection(Configuration.GetConnectionString("DB1"));
    conn.Open();
    var transaction = conn.BeginTransaction();
    try
    {
        // 执行你的数据修改SQL
        var cmd = new MySqlCommand("UPDATE user SET balance = balance - @amount WHERE id = @userId", conn, transaction);
        cmd.Parameters.AddWithValue("@amount", req.Amount);
        cmd.Parameters.AddWithValue("@userId", req.UserId);
        cmd.ExecuteNonQuery();

        // 用分布式缓存(如Redis)存储事务关联信息,替代内存缓存避免API重启丢失
        _cache.Set(txId, transaction, TimeSpan.FromMinutes(10));
        return Ok();
    }
    catch (Exception ex)
    {
        transaction.Rollback();
        return BadRequest($"预执行失败:{ex.Message}");
    }
}

// 提交事务
[HttpPost("commit")]
public IActionResult Commit([FromQuery] string txId)
{
    if (_cache.TryGetValue(txId, out MySqlTransaction transaction))
    {
        transaction.Commit();
        _cache.Remove(txId);
        return Ok();
    }
    return BadRequest("无效的事务ID");
}

// 回滚事务
[HttpPost("rollback")]
public IActionResult Rollback([FromQuery] string txId)
{
    if (_cache.TryGetValue(txId, out MySqlTransaction transaction))
    {
        try
        {
            transaction.Rollback();
        }
        finally
        {
            _cache.Remove(txId);
        }
    }
    return Ok(); // 即使事务已不存在,返回成功避免重试报错
}

2. 实现协调者逻辑(ApiTransactionScope的核心)

可以在Ocelot网关层新增中间件,或者单独写一个协调API,来完成全局事务的编排:

  • 生成全局唯一事务ID(txId)
  • 依次调用两个API的预执行接口,记录成功执行的API
  • 全部预执行成功则依次调用提交接口
  • 任意环节失败则调用所有已成功预执行API的回滚接口

示例协调代码:

private readonly HttpClient _httpClient;
private readonly ILogger<CoordinatorController> _logger;
private readonly IDistributedCache _cache;

public CoordinatorController(HttpClient httpClient, ILogger<CoordinatorController> logger, IDistributedCache cache)
{
    _httpClient = httpClient;
    _logger = logger;
    _cache = cache;
}

[HttpPost("cross-db-update")]
public async Task<IActionResult> CrossDbUpdate([FromBody] CrossUpdateRequest req)
{
    var txId = Guid.NewGuid().ToString();
    var succeededApis = new List<string>();

    try
    {
        // 调用API1预执行
        var api1Resp = await _httpClient.PostAsJsonAsync($"http://api1/prepare-update?txId={txId}", req.Api1Request);
        if (!api1Resp.IsSuccessStatusCode)
            throw new Exception(await api1Resp.Content.ReadAsStringAsync());
        succeededApis.Add("api1");

        // 调用API2预执行
        var api2Resp = await _httpClient.PostAsJsonAsync($"http://api2/prepare-update?txId={txId}", req.Api2Request);
        if (!api2Resp.IsSuccessStatusCode)
            throw new Exception(await api2Resp.Content.ReadAsStringAsync());
        succeededApis.Add("api2");

        // 提交API1
        var commitApi1 = await _httpClient.PostAsync($"http://api1/commit?txId={txId}", null);
        if (!commitApi1.IsSuccessStatusCode)
            throw new Exception("API1提交失败");

        // 提交API2
        var commitApi2 = await _httpClient.PostAsync($"http://api2/commit?txId={txId}", null);
        if (!commitApi2.IsSuccessStatusCode)
            throw new Exception("API2提交失败");

        return Ok("跨库更新成功");
    }
    catch (Exception ex)
    {
        _logger.LogError(ex, "跨库更新失败,开始回滚");
        // 回滚所有已成功的预执行
        foreach (var api in succeededApis)
        {
            try
            {
                var rollbackUrl = api == "api1" ? $"http://api1/rollback?txId={txId}" : $"http://api2/rollback?txId={txId}";
                await _httpClient.PostAsync(rollbackUrl, null);
            }
            catch (Exception rollbackEx)
            {
                _logger.LogError(rollbackEx, "回滚{Api}失败,txId:{TxId}", api, txId);
                // 这里需要添加人工补偿的告警或日志记录
            }
        }
        return BadRequest($"操作失败:{ex.Message}");
    }
}

3. 生产环境关键注意事项

  • 分布式缓存替代内存存储:示例中用分布式缓存(如Redis)存储事务对象,避免API重启后事务丢失,无法提交/回滚
  • 幂等性保障:所有接口要基于txId实现幂等,防止重复调用导致数据异常
  • 补偿机制:回滚失败时必须记录详细日志,触发人工补偿流程——分布式事务无法做到100%可靠,人工兜底是必要的
  • MySQL XA事务可选方案:如果不想手动管理事务挂起,可以启用MySQL的XA事务,预执行时开启XA事务,协调者最终调用XA的全局提交/回滚,不过XA性能略低,适合对一致性要求极高的场景

ApiTransactionScope的语法糖本质

你看到的ApiTransactionScope只是把上述协调逻辑封装成了类似本地TransactionScope的简洁语法,比如:

using(var scope = new ApiTransactionScope())
{
    await _apiClient.UpdateApi1(req1);
    await _apiClient.UpdateApi2(req2);
    scope.Complete(); // 触发提交流程
}

底层自动完成txId生成、预执行、提交/回滚的全流程,和我们手动实现的协调逻辑完全一致。

内容的提问来源于stack exchange,提问作者Pankaj Rawat

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 08:57:54