-
C#网络编程之分布式事务(2PC、3PC、TCC、本地消息表)
第58章 分布式事务实战
58.1 分布式事务(2PC、3PC、TCC、本地消息表)
一、我踩过的分布式事务坑:从“用户转账钱丢了”到“本地消息表救了我”
做支付系统的时候,遇到过这辈子最严重的线上事故:用户从A账户转1000元到B账户,结果A账户的钱扣了,B账户的钱没到账,用户直接投诉到监管部门。查了半天发现是2PC协调者挂了,A银行已经提交了扣钱操作,B银行还没收到提交指令,导致数据不一致。后来换成本地消息表方案,订单服务创建订单后写本地消息,异步通知库存服务扣减库存,即使消息发送失败,定时任务也会重试,再也没出现过数据不一致的问题。这节我把这些血泪经验揉进去,用大白话讲透四种分布式事务方案的核心原理,结合C#实战代码逐行拆解,拓展生产级优化技巧,让你的分布式系统再也不出现“钱丢了”“库存不一致”的问题!
二、分布式事务核心问题:让多个服务的操作“要么全成,要么全败”
分布式事务是指多个服务操作不同数据库时,要保证所有操作要么全部成功,要么全部失败,不能出现部分成功部分失败的情况。
大白话解释:把分布式事务比作“聚餐AA制”
1.服务实例:聚餐的每个人;
2.数据库操作:每个人付自己的AA钱;
3.分布式事务:要么所有人都付了钱,要么所有人都没付,不能出现有人付了有人没付的情况(不然没付钱的人吃了霸王餐,付钱的人亏了)。
我踩过的坑:一开始用本地事务处理跨服务操作,结果订单服务创建订单成功了,库存服务扣减库存失败,导致订单有了但库存没扣,超卖了500单——本地事务只能管单个服务的数据库操作,跨服务必须用分布式事务!
三、2PC(两阶段提交):强一致性的“会议投票”方案,适合金融场景
2PC(Two-Phase Commit)是最经典的分布式事务协议,核心是“协调者+参与者”,分两个阶段完成事务:准备阶段和提交阶段。
大白话解释:把2PC比作“公司开会投票”
1.协调者(Coordinator):会议主持人;
2.参与者(Participant):参会的部门经理;
3.准备阶段:主持人问每个经理“是否同意这个方案?”,每个经理回去准备,回复“同意”或“不同意”;
4.提交阶段:如果所有经理都同意,主持人说“方案通过,执行!”;如果有一个经理不同意,主持人说“方案取消,回滚!”。
我踩过的坑:用2PC做银行转账,结果协调者挂了,A银行已经准备好扣钱,B银行还没收到提交指令,导致A银行的钱扣了,B银行的钱没到账——2PC的单点故障问题太致命!
实战1:C#模拟2PC流程(协调者+参与者)
csharp
using System;
using System.Collections.Generic;
using System.Threading.Tasks;
namespace DistributedTransaction.TwoPhaseCommit;
// 参与者接口:每个服务实现这个接口,处理准备和提交/回滚
public interface IParticipant
{
/// <summary>
/// 准备阶段:检查是否可以执行操作,锁定资源
/// </summary>
Task<bool> PrepareAsync();
/// <summary>
/// 提交阶段:执行最终操作
/// </summary>
Task CommitAsync();
/// <summary>
/// 回滚阶段:释放锁定的资源,恢复到之前的状态
/// </summary>
Task RollbackAsync();
}
// 银行A参与者:模拟扣钱操作
public class BankAParticipant : IParticipant
{
private decimal _originalBalance;
private readonly string _userId;
private readonly decimal _amount;
public BankAParticipant(string userId, decimal amount)
{
_userId = userId;
_amount = amount;
}
public async Task<bool> PrepareAsync()
{
// 1. 查询用户余额,检查是否足够扣钱
_originalBalance = await GetUserBalanceAsync(_userId);
if (_originalBalance < _amount)
{
Console.WriteLine($"银行A:用户 {_userId} 余额不足,准备失败");
return false;
}
// 2. 锁定余额(比如加行锁,或者写入临时表)
await LockUserBalanceAsync(_userId);
Console.WriteLine($"银行A:用户 {_userId} 余额充足,准备成功,锁定余额 {_originalBalance}");
return true;
}
public async Task CommitAsync()
{
// 3. 执行扣钱操作
await DeductUserBalanceAsync(_userId, _amount);
Console.WriteLine($"银行A:用户 {_userId} 扣钱成功,金额 {_amount},剩余余额 {_originalBalance - _amount}");
}
public async Task RollbackAsync()
{
// 4. 释放锁定的余额,恢复到之前的状态
await ReleaseUserBalanceLockAsync(_userId);
Console.WriteLine($"银行A:用户 {_userId} 回滚成功,释放余额锁定");
}
// 模拟查询用户余额
private async Task<decimal> GetUserBalanceAsync(string userId)
{
await Task.Delay(100);
return 2000; // 模拟用户余额2000元
}
// 模拟锁定用户余额
private async Task LockUserBalanceAsync(string userId)
{
await Task.Delay(50);
}
// 模拟扣减用户余额
private async Task DeductUserBalanceAsync(string userId, decimal amount)
{
await Task.Delay(50);
}
// 模拟释放余额锁定
private async Task ReleaseUserBalanceLockAsync(string userId)
{
await Task.Delay(50);
}
}
// 银行B参与者:模拟加钱操作
public class BankBParticipant : IParticipant
{
private decimal _originalBalance;
private readonly string _userId;
private readonly decimal _amount;
public BankBParticipant(string userId, decimal amount)
{
_userId = userId;
_amount = amount;
}
public async Task<bool> PrepareAsync()
{
// 1. 查询用户余额(不需要检查,加钱总是可以的)
_originalBalance = await GetUserBalanceAsync(_userId);
// 2. 锁定余额,防止其他操作修改
await LockUserBalanceAsync(_userId);
Console.WriteLine($"银行B:用户 {_userId} 准备成功,锁定余额 {_originalBalance}");
return true;
}
public async Task CommitAsync()
{
// 3. 执行加钱操作
await AddUserBalanceAsync(_userId, _amount);
Console.WriteLine($"银行B:用户 {_userId} 加钱成功,金额 {_amount},剩余余额 {_originalBalance + _amount}");
}
public async Task RollbackAsync()
{
// 4. 释放锁定的余额
await ReleaseUserBalanceLockAsync(_userId);
Console.WriteLine($"银行B:用户 {_userId} 回滚成功,释放余额锁定");
}
private async Task<decimal> GetUserBalanceAsync(string userId)
{
await Task.Delay(100);
return 1000; // 模拟用户余额1000元
}
private async Task LockUserBalanceAsync(string userId)
{
await Task.Delay(50);
}
private async Task AddUserBalanceAsync(string userId, decimal amount)
{
await Task.Delay(50);
}
private async Task ReleaseUserBalanceLockAsync(string userId)
{
await Task.Delay(50);
}
}
// 协调者:负责协调所有参与者完成两阶段提交
public class TwoPhaseCommitCoordinator
{
private readonly List<IParticipant> _participants = new List<IParticipant>();
public void AddParticipant(IParticipant participant)
{
_participants.Add(participant);
}
public async Task<bool> ExecuteTransactionAsync()
{
Console.WriteLine("=== 开始2PC事务:准备阶段 ===");
bool allPrepared = true;
// 1. 准备阶段:询问所有参与者是否准备好
foreach (var participant in _participants)
{
try
{
bool prepared = await participant.PrepareAsync();
if (!prepared)
{
allPrepared = false;
break;
}
}
catch (Exception ex)
{
Console.WriteLine($"参与者准备失败:{ex.Message}");
allPrepared = false;
break;
}
}
Console.WriteLine("=== 准备阶段结束 ===");
if (allPrepared)
{
Console.WriteLine("=== 所有参与者准备成功,开始提交阶段 ===");
// 2. 提交阶段:通知所有参与者提交事务
foreach (var participant in _participants)
{
try
{
await participant.CommitAsync();
}
catch (Exception ex)
{
Console.WriteLine($"参与者提交失败:{ex.Message}");
// 提交失败需要人工介入,因为已经有参与者提交成功了
throw new Exception("事务提交失败,需要人工修复数据不一致问题");
}
}
Console.WriteLine("=== 提交阶段结束,事务成功 ===");
return true;
}
else
{
Console.WriteLine("=== 有参与者准备失败,开始回滚阶段 ===");
// 3. 回滚阶段:通知所有参与者回滚事务
foreach (var participant in _participants)
{
try
{
await participant.RollbackAsync();
}
catch (Exception ex)
{
Console.WriteLine($"参与者回滚失败:{ex.Message}");
// 回滚失败需要人工介入
}
}
Console.WriteLine("=== 回滚阶段结束,事务失败 ===");
return false;
}
}
}
// 测试代码:模拟银行转账
class Program
{
static async Task Main(string[] args)
{
// 1. 创建参与者:银行A扣钱,银行B加钱
var bankA = new BankAParticipant("user_123", 1000);
var bankB = new BankBParticipant("user_456", 1000);
// 2. 创建协调者,添加参与者
var coordinator = new TwoPhaseCommitCoordinator();
coordinator.AddParticipant(bankA);
coordinator.AddParticipant(bankB);
// 3. 执行分布式事务
await coordinator.ExecuteTransactionAsync();
}
}
逐行讲解:
1.参与者接口:定义准备、提交、回滚三个方法,每个服务实现自己的业务逻辑;
2.准备阶段:参与者检查资源是否充足,锁定资源(比如加行锁),返回是否准备成功;
3.提交阶段:所有参与者准备成功后,协调者通知所有参与者执行最终操作;
4.回滚阶段:有参与者准备失败,协调者通知所有参与者释放锁定的资源,恢复到之前的状态;
5.人工介入:提交阶段如果有参与者失败,需要人工修复数据不一致(因为已经有参与者提交成功了)。
2PC的优缺点与适用场景
| 优点 | 缺点 | 适用场景 |
|---|---|---|
| 强一致性 | 协调者单点故障,容易导致数据不一致 | 金融系统、银行转账、支付系统(强一致性要求高) |
| 实现简单 | 准备阶段锁定资源,性能差,容易阻塞 | 低并发场景,跨服务操作少的场景 |
| 事务原子性保证 | 提交阶段如果协调者挂了,部分参与者提交部分未提交 | 对性能要求不高,对一致性要求极高的场景 |
我踩过的坑:用2PC做电商库存扣减,结果准备阶段锁定了库存,提交阶段因为网络延迟超时,导致库存一直被锁定,用户下单失败——2PC的阻塞问题太严重,高并发场景根本用不了!
四、3PC(三阶段提交):2PC的“改进版会议”,解决阻塞问题
3PC(Three-Phase Commit)是2PC的改进版,核心是把2PC的准备阶段分成两个阶段:CanCommit阶段和PreCommit阶段,解决2PC的阻塞问题。
大白话解释:把3PC比作“公司开会分两步投票”
1.CanCommit阶段:主持人问每个经理“是否有能力执行这个方案?”,经理不需要准备,直接回复“是”或“否”;
2.PreCommit阶段:如果所有经理都回复“是”,主持人让每个经理准备执行方案,回复“准备成功”或“准备失败”;
3.DoCommit阶段:如果所有经理都准备成功,主持人通知执行;如果有一个经理准备失败,主持人通知回滚。
3PC vs 2PC的核心改进
1.减少阻塞时间:CanCommit阶段不需要锁定资源,只有PreCommit阶段才锁定资源,减少了资源锁定的时间;
2.超时自动提交:如果参与者在DoCommit阶段没收到协调者的指令,会自动提交事务,避免因为协调者挂了导致的阻塞;
3.单点故障影响降低:协调者挂了之后,参与者会根据超时时间自动提交或回滚,不会一直阻塞。
3PC的优缺点与适用场景
优点 缺点 适用场景
减少阻塞时间 实现复杂,需要处理更多的超时和异常情况 对性能要求稍高的强一致性场景
超时自动提交 仍然可能出现数据不一致(比如参与者自动提交了,协调者其实要回滚) 很少用,大部分场景用2PC或TCC代替
单点故障影响降低 性能还是不如最终一致性方案 金融系统的部分场景,但是实际用得少
我踩过的坑:用3PC做支付系统,结果参与者在DoCommit阶段超时自动提交了,但是协调者其实要回滚,导致数据不一致——3PC的超时自动提交机制也有坑,实际用起来还是复杂!
五、TCC(Try-Confirm-Cancel):补偿型事务,适合业务场景的“手动AA制”
TCC(Try-Confirm-Cancel)是补偿型分布式事务协议,核心是“业务逻辑拆分+补偿机制”,分三个阶段完成事务:Try阶段、Confirm阶段、Cancel阶段。
大白话解释:把TCC比作“聚餐AA制手动操作”
1.Try阶段:每个人先把钱准备好,放在桌子上(锁定资源);
2.Confirm阶段:每个人把钱交给服务员(执行最终操作);
3.Cancel阶段:每个人把钱拿回去(释放锁定的资源)。
我踩过的坑:一开始用TCC做电商库存扣减,Try阶段锁定了库存,Confirm阶段因为网络故障失败,导致库存一直被锁定,用户下单失败——后来加了定时补偿任务,扫描锁定超过10分钟的库存,自动执行Cancel阶段。
实战2:C#实现TCC事务(电商库存扣减)
csharp
using System;
using System.Threading.Tasks;
namespace DistributedTransaction.Tcc;
// TCC事务接口:每个服务实现这三个方法
public interface ITccTransaction
{
/// <summary>
/// Try阶段:锁定资源,检查业务条件
/// </summary>
Task<string> TryAsync(string orderId, string productId, int quantity);
/// <summary>
/// Confirm阶段:执行最终业务逻辑,必须是幂等的
/// </summary>
Task<bool> ConfirmAsync(string tryToken);
/// <summary>
/// Cancel阶段:释放锁定的资源,必须是幂等的
/// </summary>
Task<bool> CancelAsync(string tryToken);
}
// 库存服务的TCC实现
public class InventoryTccService : ITccTransaction
{
// 模拟库存数据库:商品ID -> 库存数量
private readonly Dictionary<string, int> _inventory = new Dictionary<string, int>
{
{ "product_123", 100 }
};
// 模拟锁定的库存:TryToken -> (商品ID, 数量)
private readonly Dictionary<string, (string ProductId, int Quantity)> _lockedInventory = new Dictionary<string, (string, int)>();
public async Task<string> TryAsync(string orderId, string productId, int quantity)
{
// 1. 检查库存是否充足
if (!_inventory.ContainsKey(productId) || _inventory[productId] < quantity)
{
Console.WriteLine($"Try阶段失败:商品 {productId} 库存不足,当前库存 {_inventory.GetValueOrDefault(productId, 0)},需要 {quantity}");
return null;
}
// 2. 锁定库存,生成唯一的TryToken(用订单ID+商品ID作为Token)
string tryToken = $"{orderId}_{productId}";
if (_lockedInventory.ContainsKey(tryToken))
{
Console.WriteLine($"Try阶段失败:订单 {orderId} 商品 {productId} 已经被锁定");
return null;
}
// 3. 扣减可用库存,增加锁定库存
_inventory[productId] -= quantity;
_lockedInventory.Add(tryToken, (productId, quantity));
Console.WriteLine($"Try阶段成功:订单 {orderId} 商品 {productId} 锁定库存 {quantity},剩余可用库存 {_inventory[productId]}");
await Task.Delay(100);
return tryToken;
}
public async Task<bool> ConfirmAsync(string tryToken)
{
// 1. 幂等性检查:如果已经Confirm过,直接返回成功
if (!_lockedInventory.ContainsKey(tryToken))
{
Console.WriteLine($"Confirm阶段成功:TryToken {tryToken} 已经处理过");
return true;
}
try
{
// 2. 执行最终操作:删除锁定的库存(相当于确认扣减)
var (productId, quantity) = _lockedInventory[tryToken];
_lockedInventory.Remove(tryToken);
Console.WriteLine($"Confirm阶段成功:订单 {tryToken.Split('_')[0]} 商品 {productId} 确认扣减库存 {quantity}");
await Task.Delay(100);
return true;
}
catch (Exception ex)
{
Console.WriteLine($"Confirm阶段失败:TryToken {tryToken} 异常 {ex.Message}");
return false;
}
}
public async Task<bool> CancelAsync(string tryToken)
{
// 1. 幂等性检查:如果已经Cancel过,直接返回成功
if (!_lockedInventory.ContainsKey(tryToken))
{
Console.WriteLine($"Cancel阶段成功:TryToken {tryToken} 已经处理过");
return true;
}
try
{
// 2. 释放锁定的库存:把锁定的库存加回可用库存
var (productId, quantity) = _lockedInventory[tryToken];
_inventory[productId] += quantity;
_lockedInventory.Remove(tryToken);
Console.WriteLine($"Cancel阶段成功:订单 {tryToken.Split('_')[0]} 商品 {productId} 释放库存 {quantity},可用库存恢复为 {_inventory[productId]}");
await Task.Delay(100);
return true;
}
catch (Exception ex)
{
Console.WriteLine($"Cancel阶段失败:TryToken {tryToken} 异常 {ex.Message}");
return false;
}
}
}
// TCC事务管理器:负责协调TCC的三个阶段
public class TccTransactionManager
{
public async Task<bool> ExecuteTransactionAsync(ITccTransaction tccService, string orderId, string productId, int quantity)
{
Console.WriteLine("=== 开始TCC事务:Try阶段 ===");
string tryToken = await tccService.TryAsync(orderId, productId, quantity);
if (string.IsNullOrEmpty(tryToken))
{
Console.WriteLine("=== Try阶段失败,事务终止 ===");
return false;
}
Console.WriteLine("=== Try阶段成功,开始Confirm阶段 ===");
bool confirmSuccess = await tccService.ConfirmAsync(tryToken);
if (confirmSuccess)
{
Console.WriteLine("=== Confirm阶段成功,事务完成 ===");
return true;
}
Console.WriteLine("=== Confirm阶段失败,开始Cancel阶段 ===");
bool cancelSuccess = await tccService.CancelAsync(tryToken);
if (cancelSuccess)
{
Console.WriteLine("=== Cancel阶段成功,事务回滚 ===");
return false;
}
Console.WriteLine("=== Cancel阶段失败,需要人工介入修复数据 ===");
return false;
}
}
// 测试代码:模拟电商下单库存扣减
class Program
{
static async Task Main(string[] args)
{
var tccService = new InventoryTccService();
var transactionManager = new TccTransactionManager();
// 1. 正常下单:库存充足
await transactionManager.ExecuteTransactionAsync(tccService, "order_456", "product_123", 10);
// 2. 库存不足下单:Try阶段失败
await transactionManager.ExecuteTransactionAsync(tccService, "order_789", "product_123", 200);
}
}
逐行讲解:
1.Try阶段:检查库存是否充足,锁定库存,生成唯一的TryToken(用于幂等性);
2.Confirm阶段:确认扣减库存,删除锁定记录,必须是幂等的(多次调用结果一样);
3.Cancel阶段:释放锁定的库存,把库存加回可用库存,必须是幂等的;
4.事务管理器:负责协调三个阶段,Try成功后执行Confirm,Confirm失败执行Cancel;
5.幂等性:用TryToken作为唯一标识,避免重复处理同一个事务。
TCC的优缺点与适用场景
| 优点 | 缺点 | 适用场景 |
|---|---|---|
| 性能高,无阻塞 | 业务代码侵入性强,需要把业务逻辑拆分成三个阶段 | 电商、社交、出行等业务场景,对性能要求高 |
| 补偿机制完善 | 需要处理幂等性、重试、补偿等复杂逻辑 | 跨服务操作多,对一致性要求较高但不是强一致的场景 |
| 适合高并发场景 | 开发成本高,每个服务都要实现TCC接口 | 互联网公司的大部分分布式事务场景 |
我踩过的坑:TCC的Confirm阶段失败后,Cancel阶段也失败了,导致库存一直被锁定——后来加了定时任务,每隔5分钟扫描一次锁定的库存,如果锁定时间超过10分钟,自动执行Cancel阶段!
六、本地消息表:最终一致性的“异步AA制”方案,适合高并发场景
本地消息表是基于最终一致性的分布式事务方案,核心是“本地事务+异步消息+补偿机制”,分三步完成事务:写本地消息、异步发送消息、消费消息并执行操作。
大白话解释:把本地消息表比作“聚餐AA制记账”
1.本地事务:每个人先在自己的记账本上记下要付的AA钱(写本地消息表),同时付自己的钱(执行本地操作);
2.异步消息:每个人把记账本上的记录拍照发给服务员(异步发送消息);
3.消费消息:服务员收到照片后,确认每个人都付了钱(执行跨服务操作),然后在记账本上打勾(更新消息状态)。
我踩过的坑:一开始本地消息表没加重试机制,结果消息发送失败,导致库存没扣减——后来加了定时任务,每隔1分钟扫描一次未处理的消息,重新发送!
实战3:C#实现本地消息表(电商订单+库存)
步骤1:数据库表设计(订单表+本地消息表)
sql
-- 订单表
CREATE TABLE `orders` (
`id` varchar(32) NOT NULL COMMENT '订单ID',
`product_id` varchar(32) NOT NULL COMMENT '商品ID',
`quantity` int NOT NULL COMMENT '数量',
`status` tinyint NOT NULL COMMENT '订单状态:0-创建中,1-已创建,2-已取消',
`create_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
PRIMARY KEY (`id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
-- 本地消息表
CREATE TABLE `local_message` (
`id` varchar(32) NOT NULL COMMENT '消息ID',
`order_id` varchar(32) NOT NULL COMMENT '订单ID',
`message_content` text NOT NULL COMMENT '消息内容(JSON格式)',
`status` tinyint NOT NULL COMMENT '消息状态:0-待发送,1-已发送,2-已消费,3-消费失败',
`retry_count` int NOT NULL DEFAULT 0 COMMENT '重试次数',
`max_retry_count` int NOT NULL DEFAULT 3 COMMENT '最大重试次数',
`create_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`update_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
PRIMARY KEY (`id`),
KEY `idx_order_id` (`order_id`),
KEY `idx_status` (`status`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
步骤2:订单服务代码(创建订单+写本地消息表)
csharp
using Microsoft.EntityFrameworkCore;
using System;
using System.Threading.Tasks;
namespace DistributedTransaction.LocalMessageTable;
// 数据库上下文
public class OrderDbContext : DbContext
{
public DbSet<Order> Orders { get; set; }
public DbSet<LocalMessage> LocalMessages { get; set; }
protected override void OnConfiguring(DbContextOptionsBuilder optionsBuilder)
{
optionsBuilder.UseMySql("server=localhost;database=order_db;user=root;password=123456;",
new MySqlServerVersion(new Version(8, 0, 29)));
}
protected override void OnModelCreating(ModelBuilder modelBuilder)
{
modelBuilder.Entity<Order>().HasKey(o => o.Id);
modelBuilder.Entity<LocalMessage>().HasKey(m => m.Id);
}
}
// 订单实体
public class Order
{
public string Id { get; set; }
public string ProductId { get; set; }
public int Quantity { get; set; }
public int Status { get; set; }
public DateTime CreateTime { get; set; } = DateTime.Now;
}
// 本地消息实体
public class LocalMessage
{
public string Id { get; set; }
public string OrderId { get; set; }
public string MessageContent { get; set; }
public int Status { get; set; }
public int RetryCount { get; set; }
public int MaxRetryCount { get; set; } = 3;
public DateTime CreateTime { get; set; } = DateTime.Now;
public DateTime UpdateTime { get; set; } = DateTime.Now;
}
// 订单服务:创建订单+写本地消息表
public class OrderService
{
public async Task<bool> CreateOrderAsync(string productId, int quantity)
{
using var dbContext = new OrderDbContext();
using var transaction = await dbContext.Database.BeginTransactionAsync(); // 本地事务,保证订单和消息要么都写成功要么都失败
try
{
// 1. 创建订单
string orderId = Guid.NewGuid().ToString("N");
var order = new Order
{
Id = orderId,
ProductId = productId,
Quantity = quantity,
Status = 1 // 已创建
};
dbContext.Orders.Add(order);
await dbContext.SaveChangesAsync();
Console.WriteLine($"订单 {orderId} 创建成功");
// 2. 写本地消息表:消息内容是库存扣减的指令
string messageId = Guid.NewGuid().ToString("N");
var messageContent = new
{
OrderId = orderId,
ProductId = productId,
Quantity = quantity
};
string messageJson = System.Text.Json.JsonSerializer.Serialize(messageContent);
var localMessage = new LocalMessage
{
Id = messageId,
OrderId = orderId,
MessageContent = messageJson,
Status = 0 // 待发送
};
dbContext.LocalMessages.Add(localMessage);
await dbContext.SaveChangesAsync();
Console.WriteLine($"本地消息 {messageId} 写入成功,内容:{messageJson}");
// 3. 提交本地事务
await transaction.CommitAsync();
Console.WriteLine("=== 本地事务提交成功 ===");
return true;
}
catch (Exception ex)
{
// 4. 回滚本地事务
await transaction.RollbackAsync();
Console.WriteLine($"创建订单失败:{ex.Message},本地事务回滚");
return false;
}
}
}
// 消息发送器:异步发送本地消息到消息队列(比如RabbitMQ)
public class MessageSender
{
private readonly OrderDbContext _dbContext = new OrderDbContext();
public async Task SendPendingMessagesAsync()
{
// 1. 查询待发送的消息(状态0,重试次数<最大重试次数)
var pendingMessages = await _dbContext.LocalMessages
.Where(m => m.Status == 0 && m.RetryCount < m.MaxRetryCount)
.ToListAsync();
foreach (var message in pendingMessages)
{
try
{
Console.WriteLine($"开始发送消息 {message.Id},内容:{message.MessageContent}");
// 2. 模拟发送消息到RabbitMQ
await SendToRabbitMqAsync(message.MessageContent);
// 3. 更新消息状态为已发送
message.Status = 1;
message.RetryCount++;
await _dbContext.SaveChangesAsync();
Console.WriteLine($"消息 {message.Id} 发送成功,状态更新为已发送");
}
catch (Exception ex)
{
// 4. 发送失败,增加重试次数
message.RetryCount++;
if (message.RetryCount >= message.MaxRetryCount)
{
message.Status = 3; // 消费失败,需要人工介入
Console.WriteLine($"消息 {message.Id} 重试次数超过最大值,状态更新为消费失败");
}
await _dbContext.SaveChangesAsync();
Console.WriteLine($"消息 {message.Id} 发送失败:{ex.Message},重试次数增加到 {message.RetryCount}");
}
}
}
private async Task SendToRabbitMqAsync(string messageContent)
{
// 实际代码:发送消息到RabbitMQ
await Task.Delay(100);
// 模拟发送失败:每3条消息失败1次
// if (new Random().Next(3) == 0) throw new Exception("RabbitMQ连接失败");
}
}
// 测试代码:模拟创建订单
class Program
{
static async Task Main(string[] args)
{
var orderService = new OrderService();
await orderService.CreateOrderAsync("product_123", 10);
// 模拟定时任务发送消息(生产环境用Quartz.NET或Hangfire)
var messageSender = new MessageSender();
await messageSender.SendPendingMessagesAsync();
}
}
逐行讲解:
1.本地事务:创建订单和写本地消息表在同一个本地事务中,保证要么都成功要么都失败;
2.本地消息表:记录要发送的跨服务操作指令,状态字段跟踪消息处理情况;
3.消息发送器:定时扫描待发送的消息,发送到消息队列,发送失败重试最多3次;
4.重试机制:发送失败时增加重试次数,超过最大值标记为消费失败,人工介入;
5.幂等性:用订单ID作为唯一标识,库存服务消费消息时检查是否已经处理过该订单。
步骤3:库存服务代码(消费消息+扣减库存)
csharp
using Microsoft.EntityFrameworkCore;
using RabbitMQ.Client;
using RabbitMQ.Client.Events;
using System;
using System.Text;
using System.Threading.Tasks;
namespace DistributedTransaction.LocalMessageTable.InventoryService;
// 库存数据库上下文
public class InventoryDbContext : DbContext
{
public DbSet<Inventory> Inventories { get; set; }
public DbSet<ProcessedOrder> ProcessedOrders { get; set; } // 记录已处理的订单,实现幂等性
protected override void OnConfiguring(DbContextOptionsBuilder optionsBuilder)
{
optionsBuilder.UseMySql("server=localhost;database=inventory_db;user=root;password=123456;",
new MySqlServerVersion(new Version(8, 0, 29)));
}
protected override void OnModelCreating(ModelBuilder modelBuilder)
{
modelBuilder.Entity<Inventory>().HasKey(i => i.ProductId);
modelBuilder.Entity<ProcessedOrder>().HasKey(p => p.OrderId);
}
}
// 库存实体
public class Inventory
{
public string ProductId { get; set; }
public int Quantity { get; set; }
public DateTime UpdateTime { get; set; } = DateTime.Now;
}
// 已处理订单实体(幂等性)
public class ProcessedOrder
{
public string OrderId { get; set; }
public DateTime ProcessTime { get; set; } = DateTime.Now;
}
// 库存服务:消费消息扣减库存
public class InventoryService
{
private readonly InventoryDbContext _dbContext = new InventoryDbContext();
public async Task ConsumeMessageAsync(string messageJson)
{
using var transaction = await _dbContext.Database.BeginTransactionAsync();
try
{
// 1. 解析消息
var message = System.Text.Json.JsonSerializer.Deserialize<InventoryDeductMessage>(messageJson);
// 2. 幂等性检查:检查订单是否已经处理过
if (await _dbContext.ProcessedOrders.AnyAsync(p => p.OrderId == message.OrderId))
{
Console.WriteLine($"订单 {message.OrderId} 已经处理过,跳过");
await transaction.CommitAsync();
return;
}
// 3. 扣减库存
var inventory = await _dbContext.Inventories.FindAsync(message.ProductId);
if (inventory == null || inventory.Quantity < message.Quantity)
{
Console.WriteLine($"订单 {message.OrderId} 扣减库存失败:商品 {message.ProductId} 库存不足");
await transaction.RollbackAsync();
return;
}
inventory.Quantity -= message.Quantity;
inventory.UpdateTime = DateTime.Now;
_dbContext.Inventories.Update(inventory);
// 4. 记录已处理订单(幂等性)
_dbContext.ProcessedOrders.Add(new ProcessedOrder { OrderId = message.OrderId });
await _dbContext.SaveChangesAsync();
await transaction.CommitAsync();
Console.WriteLine($"订单 {message.OrderId} 扣减库存成功:商品 {message.ProductId} 数量 {message.Quantity},剩余库存 {inventory.Quantity}");
}
catch (Exception ex)
{
await transaction.RollbackAsync();
Console.WriteLine($"消费消息失败:{ex.Message}");
// 生产环境:可以把消息重新放回队列,或者死信队列
}
}
}
// 消息消费模型
public class InventoryDeductMessage
{
public string OrderId { get; set; }
public string ProductId { get; set; }
public int Quantity { get; set; }
}
// 测试代码:模拟消费RabbitMQ消息
class Program
{
static async Task Main(string[] args)
{
var inventoryService = new InventoryService();
// 模拟消费消息
await inventoryService.ConsumeMessageAsync(@"{""OrderId"":""order_456"",""ProductId"":""product_123"",""Quantity"":10}");
}
}
逐行讲解:
1.幂等性检查:用已处理订单表记录已经处理过的订单,避免重复扣减库存;
2.本地事务:扣减库存和记录已处理订单在同一个本地事务中,保证原子性;
3.异常处理:消费失败时回滚事务,消息可以重新放回队列或死信队列,等待重试;
4.状态更新:库存扣减成功后更新库存数量和时间,保证数据一致性。
本地消息表的优缺点与适用场景
| 优点 | 缺点 | 适用场景 |
|---|---|---|
| 性能高,无阻塞 | 最终一致性,有延迟(消息发送和消费需要时间) | 电商、社交、出行等高并发场景 |
| 实现简单,侵入性低 | 需要处理消息积压、重试、幂等性等问题 | 跨服务操作多,对一致性要求不是强一致的场景 |
| 容错性强,消息不丢失 | 本地消息表会增加数据库的存储压力 | 互联网公司的大部分分布式事务场景 |
| 不需要依赖中间件(可以用定时任务代替消息队列) | 定时任务扫描会增加数据库的查询压力 | 中小团队,不想依赖复杂中间件的场景 |
我踩过的坑:本地消息表的定时任务扫描频率太高,导致数据库查询压力大——后来改成每1分钟扫描一次,同时加了索引优化查询!
七、四种分布式事务方案对比与选择指南
| 方案 | 一致性级别 | 性能 | 实现复杂度 | 适用场景 |
|---|---|---|---|---|
| 2PC | 强一致 | 低 | 中 | 金融系统、银行转账、强一致性要求高的场景 |
| 3PC | 强一致 | 中 | 高 | 很少用,大部分场景用2PC或TCC代替 |
| TCC | 最终一致/强一致 | 高 | 高 | 电商、社交、出行等业务场景,对性能要求高 |
| 本地消息表 | 最终一致 | 高 | 中 | 高并发场景,对一致性要求不是强一致的场景 |
生产级选择指南
1.优先选最终一致性方案:除非是金融系统等强一致性场景,否则优先选本地消息表或TCC,性能高,适合高并发;
2.避免2PC:2PC的阻塞问题和单点故障问题太严重,高并发场景根本用不了;
3.中小团队选本地消息表:实现简单,不需要依赖复杂的中间件,容错性强;
4.复杂业务场景选TCC:可以灵活处理各种业务逻辑,补偿机制完善,适合复杂的跨服务操作;
5.强一致性场景选2PC:比如银行转账、支付系统,必须保证强一致性,只能用2PC(或者分布式数据库的XA事务)。
八、生产级分布式事务最佳实践与踩坑总结
-
生产级最佳实践
优先最终一致:强一致性方案性能差,尽量用最终一致性代替;
幂等性处理:每个操作必须是幂等的,避免重复处理;
重试机制:消息发送和消费失败时重试最多3次,超过则人工介入;
补偿机制:定时检查数据一致性,发现不一致时通过补偿任务修复;
监控告警:监控消息积压、消费延迟、数据不一致的情况,设置告警阈值;
避免分布式事务:能不用就不用,尽量把跨服务操作改成单服务操作,或者用领域事件、CQRS等架构模式减少分布式事务的使用。 -
我踩过的坑总结
1.用2PC做电商库存:系统在高并发下直接阻塞,QPS只有200——换成本地消息表后,QPS提升到1000;
2.TCC的Confirm阶段失败:库存一直被锁定,用户下单失败——加了定时补偿任务,扫描锁定超过10分钟的库存自动释放;
3.本地消息表没加重试:消息发送失败,导致库存没扣减——加了定时任务扫描未处理的消息;
4.没处理幂等性:消息重复消费,导致库存被重复扣减——用订单ID作为唯一标识,记录已处理订单;
5.消息积压:消息发送速度慢于消费速度,导致消息积压——增加消费实例数量,优化消费逻辑。
九、总结
分布式事务是分布式系统中最复杂的问题之一,四种方案各有优缺点:2PC适合强一致性场景但性能差,3PC是2PC的改进但用得少,TCC适合复杂业务场景但开发成本高,本地消息表适合高并发场景且实现简单。实战中要根据业务场景选择合适的方案,优先用最终一致性方案,避免强一致性方案,同时注意幂等性、重试、补偿、监控等问题,减少踩坑。
下一节我们会学习分布式缓存(Redis、Memcached),解决分布式系统中的数据缓存问题,提升系统性能。
转载请注明出处:










