.NET Core WebAPI多数据库跨API分布式事务回滚实现咨询
跨API双MySQL库事务回滚方案(基于ApiTransactionScope思路)
核心原理:编排式分布式事务(TCC模型)
ApiTransactionScope本质是对编排式TCC事务的封装,通过一个协调角色(网关或独立协调服务)控制各个API的三个阶段:
- Try(预执行):执行数据修改但不提交,保留回滚能力
- Confirm(提交):确认所有预执行成功后,正式提交事务
- 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
相关产品推荐
相关产品推荐

