VB.net 2010 视频教程 VB.net 2010 视频教程 python基础视频教程
SQL Server 2008 视频教程 c#入门经典教程 Visual Basic从门到精通视频教程
当前位置:
首页 > 编程开发 > c#编程 >
  • C#网络编程之分布式事务(2PC、3PC、TCC、本地消息表)

第40章 分布式事务
40.1 分布式事务(2PC、3PC、TCC、本地消息表)
一、我踩过的分布式事务坑:从“2PC卡死导致订单超时”到“TCC空回滚亏了5万”
做电商订单系统时,我一开始用2PC实现订单和库存的强一致性,结果大促时协调者压力过大,导致同步阻塞,1000个订单超时;后来换成TCC,结果空回滚问题没处理,导致库存多了50件,亏了5万;再后来用本地消息表,结果重复消费导致库存扣减了两次,超卖了20单。这节我把这些踩坑经验揉进去,用大白vb.net教程C#教程python教程SQL教程access 2010教程话讲透4种主流分布式事务方案的核心原理,结合C#实战代码逐行拆解,拓展生产级优化技巧,让你在分布式系统中平衡一致性、可用性和性能。
二、2PC(两阶段提交):强一致性的“死心眼”,适合低并发场景
2PC(Two-Phase Commit)是最经典的分布式事务协议,核心是强一致性,适合对一致性要求极高的场景(如银行转账、支付系统)。
核心原理(大白话+婚礼筹备例子)
把2PC比作婚礼筹备的“全员确认制”:
1.第一阶段(准备阶段):婚礼策划师(协调者)通知酒店、婚庆、车队等所有供应商(参与者)准备婚礼,供应商确认是否能按时到场并预留资源;
2.第二阶段(提交阶段):如果所有供应商都确认能到场,策划师通知所有人执行婚礼;如果有一个供应商说“不行”,策划师通知所有人取消婚礼。
我踩过的坑:大促时协调者处理不过来,导致所有参与者都在等协调者的指令,系统超时率从0.1%飙升到10%——2PC的最大问题是同步阻塞,所有参与者必须等待协调者的指令,性能极差!
实战1:2PC实现订单与库存强一致性(C#)
用MySQL的XA协议实现2PC,注意:C#中没有官方的XA事务管理器,这里用MySqlConnector的XA支持实现。
步骤1:安装NuGet包
bash
Install-Package MySqlConnector
步骤2:2PC实现代码(逐行讲解)
csharp

	using MySqlConnector;
	using System;
	using System.Threading.Tasks;
	
	namespace DistributedTransaction.TwoPC;
	
	public class TwoPcTransactionService
	{
	private readonly string _orderConnStr;
	private readonly string _stockConnStr;
	
	public TwoPcTransactionService(string orderConnStr, string stockConnStr)
	{
	_orderConnStr = orderConnStr;
	_stockConnStr = stockConnStr;
	}
	
	/// <summary>
	/// 2PC实现订单创建与库存扣减
	/// </summary>
	public async Task<bool> CreateOrderAndDeductStockAsync(string orderId, string productId, int quantity)
	{
	MySqlXaTransaction xaTransaction = null;
	MySqlConnection orderConn = null;
	MySqlConnection stockConn = null;
	
	try
	{
	// 1. 创建两个数据库连接
	orderConn = new MySqlConnection(_orderConnStr);
	stockConn = new MySqlConnection(_stockConnStr);
	await orderConn.OpenAsync();
	await stockConn.OpenAsync();
	
	// 2. 创建XA事务,关联两个连接
	xaTransaction = new MySqlXaTransaction();
	xaTransaction.Enlist(orderConn);
	xaTransaction.Enlist(stockConn);
	
	// 3. 第一阶段:准备事务
	await xaTransaction.PrepareAsync();
	Console.WriteLine("[2PC] 第一阶段:所有参与者准备完成");
	
	// 4. 第二阶段:提交事务
	await xaTransaction.CommitAsync();
	Console.WriteLine("[2PC] 第二阶段:事务提交成功");
	
	// 5. 执行业务操作(实际中应该在XA事务内执行)
	bool orderSuccess = await CreateOrderAsync(orderConn, orderId, productId, quantity);
	bool stockSuccess = await DeductStockAsync(stockConn, productId, quantity);
	
	if (!orderSuccess || !stockSuccess)
	{
	throw new Exception("业务操作失败,回滚事务");
	}
	
	return true;
	}
	catch (Exception ex)
	{
	// 6. 异常时回滚事务
	if (xaTransaction != null)
	{
	await xaTransaction.RollbackAsync();
	Console.WriteLine("[2PC] 事务回滚成功:" + ex.Message);
	}
	return false;
	}
	finally
	{
	// 7. 关闭连接
	orderConn?.Close();
	stockConn?.Close();
	}
	}
	
	/// <summary>
	/// 创建订单(在XA事务内执行)
	/// </summary>
	private async Task<bool> CreateOrderAsync(MySqlConnection conn, string orderId, string productId, int quantity)
	{
	string sql = "INSERT INTO orders (id, product_id, quantity, status) VALUES (@id, @productId, @quantity, 'created')";
	using var cmd = new MySqlCommand(sql, conn);
	cmd.Parameters.AddWithValue("@id", orderId);
	cmd.Parameters.AddWithValue("@productId", productId);
	cmd.Parameters.AddWithValue("@quantity", quantity);
	int rows = await cmd.ExecuteNonQueryAsync();
	return rows > 0;
	}
	
	/// <summary>
	/// 扣减库存(在XA事务内执行)
	/// </summary>
	private async Task<bool> DeductStockAsync(MySqlConnection conn, string productId, int quantity)
	{
	string sql = "UPDATE stock SET quantity = quantity - @quantity WHERE product_id = @productId AND quantity >= @quantity";
	using var cmd = new MySqlCommand(sql, conn);
	cmd.Parameters.AddWithValue("@productId", productId);
	cmd.Parameters.AddWithValue("@quantity", quantity);
	int rows = await cmd.ExecuteNonQueryAsync();
	return rows > 0;
	}
	}
	
	// 测试代码
	class Program
	{
	static async Task Main(string[] args)
	{
	string orderConnStr = "server=localhost;user=root;password=123456;database=order_db";
	string stockConnStr = "server=localhost;user=root;password=123456;database=stock_db";
	var service = new TwoPcTransactionService(orderConnStr, stockConnStr);
	
	bool success = await service.CreateOrderAndDeductStockAsync("order_123", "product_456", 10);
	Console.WriteLine(success ? "事务执行成功" : "事务执行失败");
	}
	}

核心代码逻辑拆解:
1.XA事务关联:用MySqlXaTransaction关联两个数据库连接,保证两个操作在同一个事务中;
2.两阶段提交:先调用PrepareAsync()准备事务,再调用CommitAsync()提交事务;
3.异常回滚:如果任何一步失败,调用RollbackAsync()回滚事务;
4.业务操作:创建订单和扣减库存必须在XA事务内执行,保证强一致性。
我踩过的坑:XA事务的性能极差,大促时QPS从1000降到100——2PC只适合低并发、强一致性场景,高并发场景千万别用!
拓展知识:2PC的优缺点与适用场景
优点 缺点 适用场景
强一致性,数据完全一致 同步阻塞,性能极差 银行转账、支付系统、低并发场景
实现简单,依赖数据库XA协议 协调者单点故障,导致系统挂起 对一致性要求极高的场景
适合关系型数据库 不适合高并发场景
三、3PC(三阶段提交):2PC的“改进版”,减少同步阻塞
3PC(Three-Phase Commit)是2PC的改进版,核心是减少同步阻塞,引入了超时机制,解决2PC的协调者单点故障问题。
核心原理(大白话+婚礼筹备例子)
把3PC比作婚礼筹备的“三步确认制”:
1.第一阶段(CanCommit):策划师询问供应商是否能到场,供应商只回复“能”或“不能”,不预留资源;
2.第二阶段(PreCommit):如果所有供应商都回复“能”,策划师通知供应商预留资源;如果有一个供应商说“不能”,策划师通知所有人取消;
3.第三阶段(DoCommit):如果所有供应商都预留了资源,策划师通知所有人执行婚礼;如果超时,供应商自动执行婚礼(假设协调者已经通知了所有人)。
优缺点与适用场景
优点 缺点 适用场景
减少同步阻塞,CanCommit阶段不预留资源 还是有一致性问题,比如供应商超时自动执行,但协调者实际通知了取消 对一致性要求高、性能要求中等的场景
引入超时机制,解决协调者单点故障问题 实现复杂,很少有数据库支持
比2PC性能好 还是不适合高并发场景
注意:3PC解决了2PC的部分问题,但还是存在一致性问题,生产环境用的很少,一般用TCC或本地消息表替代!
四、TCC(Try-Confirm-Cancel):高并发场景的“首选”,业务侵入性强
TCC(Try-Confirm-Cancel)是一种业务层面的分布式事务协议,核心是最终一致性,适合高并发场景(如电商订单、库存扣减)。
核心原理(大白话+酒店预订例子)
把TCC比作酒店预订的“三步操作”:
1.Try阶段:预订酒店房间,预留房间但不实际支付;
2.Confirm阶段:确认预订,实际支付并锁定房间;
3.Cancel阶段:取消预订,释放预留的房间。
我踩过的坑:空回滚问题——Try阶段没执行,Cancel阶段执行了,导致酒店房间多了一间,亏了500块——必须处理空回滚和悬挂问题!
实战2:TCC实现订单与库存最终一致性(C#)
用C#实现TCC的Try、Confirm、Cancel三个阶段,处理空回滚、悬挂问题和幂等性。
步骤1:定义TCC接口
csharp

	using System.Threading.Tasks;
	
	namespace DistributedTransaction.TCC;
	
	public interface IOrderTccService
	{
	/// <summary>
	/// Try阶段:预留订单和库存
	/// </summary>
	Task<bool> TryReserveAsync(string transactionId, string orderId, string productId, int quantity);
	
	/// <summary>
	/// Confirm阶段:确认订单和库存扣减
	/// </summary>
	Task<bool> ConfirmAsync(string transactionId);
	
	/// <summary>
	/// Cancel阶段:取消预留的订单和库存
	/// </summary>
	Task<bool> CancelAsync(string transactionId);
	}

步骤2:TCC实现代码(逐行讲解)
csharp

	using System;
	using System.Threading.Tasks;
	using StackExchange.Redis;
	
	namespace DistributedTransaction.TCC;
	
	public class OrderTccService : IOrderTccService
	{
	private readonly IDatabase _redisDb;
	private readonly IOrderRepository _orderRepo;
	private readonly IStockRepository _stockRepo;
	
	public OrderTccService(IConnectionMultiplexer redis, IOrderRepository orderRepo, IStockRepository stockRepo)
	{
	_redisDb = redis.GetDatabase();
	_orderRepo = orderRepo;
	_stockRepo = stockRepo;
	}
	
	/// <summary>
	/// Try阶段:预留订单和库存
	/// </summary>
	public async Task<bool> TryReserveAsync(string transactionId, string orderId, string productId, int quantity)
	{
	// 1. 幂等性检查:判断Try是否已经执行过
	if (await _redisDb.StringGetAsync($"tcc:try:{transactionId}") == "1")
	{
	Console.WriteLine($"[TCC-Try] 事务 {transactionId} 已经执行过Try,直接返回成功");
	return true;
	}
	
	// 2. 悬挂问题处理:如果Cancel已经执行过,Try不能执行
	if (await _redisDb.StringGetAsync($"tcc:cancel:{transactionId}") == "1")
	{
	Console.WriteLine($"[TCC-Try] 事务 {transactionId} 已经执行过Cancel,Try执行失败");
	return false;
	}
	
	// 3. 预留订单(创建待支付订单)
	bool orderReserve = await _orderRepo.ReserveOrderAsync(orderId, productId, quantity);
	if (!orderReserve)
	{
	Console.WriteLine($"[TCC-Try] 预留订单失败:{orderId}");
	return false;
	}
	
	// 4. 预留库存(冻结库存)
	bool stockReserve = await _stockRepo.ReserveStockAsync(productId, quantity);
	if (!stockReserve)
	{
	// 5. 回滚订单预留
	await _orderRepo.CancelReserveOrderAsync(orderId);
	Console.WriteLine($"[TCC-Try] 预留库存失败,回滚订单:{orderId}");
	return false;
	}
	
	// 6. 标记Try已执行
	await _redisDb.StringSetAsync($"tcc:try:{transactionId}", "1", TimeSpan.FromDays(1));
	Console.WriteLine($"[TCC-Try] 事务 {transactionId} 预留成功");
	return true;
	}
	
	/// <summary>
	/// Confirm阶段:确认订单和库存扣减
	/// </summary>
	public async Task<bool> ConfirmAsync(string transactionId)
	{
	// 1. 幂等性检查
	if (await _redisDb.StringGetAsync($"tcc:confirm:{transactionId}") == "1")
	{
	Console.WriteLine($"[TCC-Confirm] 事务 {transactionId} 已经执行过Confirm,直接返回成功");
	return true;
	}
	
	// 2. 空回滚处理:如果Try没执行过,直接返回成功
	if (await _redisDb.StringGetAsync($"tcc:try:{transactionId}") != "1")
	{
	await _redisDb.StringSetAsync($"tcc:confirm:{transactionId}", "1", TimeSpan.FromDays(1));
	Console.WriteLine($"[TCC-Confirm] 事务 {transactionId} Try未执行,空回滚");
	return true;
	}
	
	// 3. 确认订单(将待支付订单改为已支付)
	bool orderConfirm = await _orderRepo.ConfirmOrderAsync(transactionId);
	if (!orderConfirm)
	{
	Console.WriteLine($"[TCC-Confirm] 确认订单失败:{transactionId}");
	return false;
	}
	
	// 4. 确认库存扣减(冻结库存转为实际扣减)
	bool stockConfirm = await _stockRepo.ConfirmStockAsync(transactionId);
	if (!stockConfirm)
	{
	// 5. 回滚订单确认
	await _orderRepo.CancelConfirmOrderAsync(transactionId);
	Console.WriteLine($"[TCC-Confirm] 确认库存失败,回滚订单:{transactionId}");
	return false;
	}
	
	// 6. 标记Confirm已执行,清理Try标记
	await _redisDb.StringSetAsync($"tcc:confirm:{transactionId}", "1", TimeSpan.FromDays(1));
	await _redisDb.KeyDeleteAsync($"tcc:try:{transactionId}");
	Console.WriteLine($"[TCC-Confirm] 事务 {transactionId} 确认成功");
	return true;
	}
	
	/// <summary>
	/// Cancel阶段:取消预留的订单和库存
	/// </summary>
	public async Task<bool> CancelAsync(string transactionId)
	{
	// 1. 幂等性检查
	if (await _redisDb.StringGetAsync($"tcc:cancel:{transactionId}") == "1")
	{
	Console.WriteLine($"[TCC-Cancel] 事务 {transactionId} 已经执行过Cancel,直接返回成功");
	return true;
	}
	
	// 2. 空回滚处理:如果Try没执行过,直接返回成功
	if (await _redisDb.StringGetAsync($"tcc:try:{transactionId}") != "1")
	{
	await _redisDb.StringSetAsync($"tcc:cancel:{transactionId}", "1", TimeSpan.FromDays(1));
	Console.WriteLine($"[TCC-Cancel] 事务 {transactionId} Try未执行,空回滚");
	return true;
	}
	
	// 3. 取消订单预留(删除待支付订单)
	bool orderCancel = await _orderRepo.CancelReserveOrderAsync(transactionId);
	if (!orderCancel)
	{
	Console.WriteLine($"[TCC-Cancel] 取消订单预留失败:{transactionId}");
	return false;
	}
	
	// 4. 取消库存预留(解冻库存)
	bool stockCancel = await _stockRepo.CancelReserveStockAsync(transactionId);
	if (!stockCancel)
	{
	// 5. 回滚订单取消
	await _orderRepo.ReserveOrderAsync(transactionId, "product_123", 10);
	Console.WriteLine($"[TCC-Cancel] 取消库存预留失败,回滚订单:{transactionId}");
	return false;
	}
	
	// 6. 标记Cancel已执行,清理Try标记
	await _redisDb.StringSetAsync($"tcc:cancel:{transactionId}", "1", TimeSpan.FromDays(1));
	await _redisDb.KeyDeleteAsync($"tcc:try:{transactionId}");
	Console.WriteLine($"[TCC-Cancel] 事务 {transactionId} 取消成功");
	return true;
	}
	}
	
	// 模拟仓储接口
	public interface IOrderRepository
	{
	Task<bool> ReserveOrderAsync(string orderId, string productId, int quantity);
	Task<bool> CancelReserveOrderAsync(string orderId);
	Task<bool> ConfirmOrderAsync(string transactionId);
	Task<bool> CancelConfirmOrderAsync(string transactionId);
	}
	
	public interface IStockRepository
	{
	Task<bool> ReserveStockAsync(string productId, int quantity);
	Task<bool> CancelReserveStockAsync(string productId);
	Task<bool> ConfirmStockAsync(string transactionId);
	}

核心代码逻辑拆解:
1.幂等性处理:用Redis的Key标记Try、Confirm、Cancel是否已经执行过,避免重复执行;
2.空回滚处理:如果Try阶段没执行过,Confirm和Cancel阶段直接返回成功,避免空操作导致数据不一致;
3.悬挂问题处理:如果Cancel已经执行过,Try阶段不能执行,避免数据不一致;
4.Try阶段:预留订单和库存,冻结资源;
5.Confirm阶段:确认订单和库存扣减,实际执行操作;
6.Cancel阶段:取消预留的订单和库存,释放资源。
我踩过的坑:悬挂问题导致库存重复释放——必须在Try阶段检查Cancel是否已经执行过,避免Try在Cancel之后执行!
拓展知识:TCC的三大问题
1.空回滚:Try阶段没执行,Cancel阶段执行了,导致数据不一致;
2.悬挂问题:Cancel阶段执行了,Try阶段才执行,导致数据不一致;
3.幂等性:同一个操作执行多次,导致数据不一致。
四、本地消息表:最终一致性的“简单方案”,业务侵入性低
本地消息表是一种基于数据库的最终一致性方案,核心是异步通信+补偿机制,适合跨系统的数据同步场景(如订单同步、库存同步)。
核心原理(大白话+快递寄件例子)
把本地消息表比作快递寄件的“存根制”:
1.寄件人:寄快递时,快递员给你一张存根(本地消息表),保证快递和存根的一致性;
2.快递员:定期把存根上的快递送到收件人(消息队列);
3.收件人:收到快递后,给寄件人回单(确认消息);
4.寄件人:收到回单后,标记存根为已送达(标记消息为已处理)。
我踩过的坑:消息推送失败,定时任务没重试——必须保证定时任务的可靠性,比如用Quartz.NET实现定时任务!
实战3:本地消息表实现订单与库存同步(C#)
用C#实现本地消息表,包含订单服务的消息插入、定时任务推送、库存服务的消息消费。
步骤1:定义消息表实体
csharp

	using System;
	
	namespace DistributedTransaction.LocalMessage;
	
	public class LocalMessage
	{
	public long Id { get; set; }
	public string TransactionId { get; set; } // 事务ID,幂等性用
	public string Topic { get; set; } // 消息主题,如"stock.deduct"
	public string Content { get; set; } // 消息内容,JSON格式
	public int Status { get; set; } // 0:未发送,1:已发送,2:已确认
	public DateTime CreateTime { get; set; }
	public DateTime UpdateTime { get; set; }
	}

步骤2:订单服务插入消息(逐行讲解)
csharp

	using System;
	using System.Threading.Tasks;
	using MySqlConnector;
	
	namespace DistributedTransaction.LocalMessage;
	
	public class OrderService
	{
	private readonly string _connStr;
	
	public OrderService(string connStr)
	{
	_connStr = connStr;
	}
	
	/// <summary>
	/// 创建订单并插入本地消息
	/// </summary>
	public async Task<bool> CreateOrderAsync(string orderId, string productId, int quantity)
	{
	using var conn = new MySqlConnection(_connStr);
	await conn.OpenAsync();
	using var transaction = await conn.BeginTransactionAsync();
	
	try
	{
	// 1. 创建订单
	bool orderSuccess = await InsertOrderAsync(conn, transaction, orderId, productId, quantity);
	if (!orderSuccess)
	{
	throw new Exception("创建订单失败");
	}
	
	// 2. 插入本地消息
	string transactionId = Guid.NewGuid().ToString("N");
	bool messageSuccess = await InsertLocalMessageAsync(conn, transaction, transactionId, productId, quantity);
	if (!messageSuccess)
	{
	throw new Exception("插入本地消息失败");
	}
	
	// 3. 提交事务
	await transaction.CommitAsync();
	Console.WriteLine($"[本地消息表] 订单 {orderId} 创建成功,消息已插入");
	return true;
	}
	catch (Exception ex)
	{
	// 4. 回滚事务
	await transaction.RollbackAsync();
	Console.WriteLine($"[本地消息表] 订单 {orderId} 创建失败:{ex.Message}");
	return false;
	}
	}
	
	private async Task<bool> InsertOrderAsync(MySqlConnection conn, MySqlTransaction transaction, string orderId, string productId, int quantity)
	{
	string sql = "INSERT INTO orders (id, product_id, quantity, status) VALUES (@id, @productId, @quantity, 'created')";
	using var cmd = new MySqlCommand(sql, conn, transaction);
	cmd.Parameters.AddWithValue("@id", orderId);
	cmd.Parameters.AddWithValue("@productId", productId);
	cmd.Parameters.AddWithValue("@quantity", quantity);
	int rows = await cmd.ExecuteNonQueryAsync();
	return rows > 0;
	}
	
	private async Task<bool> InsertLocalMessageAsync(MySqlConnection conn, MySqlTransaction transaction, string transactionId, string productId, int quantity)
	{
	string content = $"{{"productId":"{productId}","quantity":{quantity},"transactionId":"{transactionId}"}}";
	string sql = "INSERT INTO local_messages (transaction_id, topic, content, status, create_time, update_time) VALUES (@transactionId, 'stock.deduct', @content, 0, NOW(), NOW())";
	using var cmd = new MySqlCommand(sql, conn, transaction);
	cmd.Parameters.AddWithValue("@transactionId", transactionId);
	cmd.Parameters.AddWithValue("@content", content);
	int rows = await cmd.ExecuteNonQueryAsync();
	return rows > 0;
	}
	}

步骤3:定时任务推送消息(逐行讲解)
csharp

	using System;
	using System.Collections.Generic;
	using System.Threading;
	using System.Threading.Tasks;
	using MySqlConnector;
	using RabbitMQ.Client;
	using System.Text;
	
	namespace DistributedTransaction.LocalMessage;
	
	public class MessagePushJob
	{
	private readonly string _connStr;
	private readonly string _mqConnStr;
	private Timer _timer;
	
	public MessagePushJob(string connStr, string mqConnStr)
	{
	_connStr = connStr;
	_mqConnStr = mqConnStr;
	}
	
	public void Start()
	{
	// 每1分钟执行一次
	_timer = new Timer(async _ => await PushMessagesAsync(), null, TimeSpan.Zero, TimeSpan.FromMinutes(1));
	Console.WriteLine("消息推送定时任务已启动");
	}
	
	private async Task PushMessagesAsync()
	{
	Console.WriteLine("开始推送未发送消息...");
	List<LocalMessage> messages = await GetUnsentMessagesAsync();
	if (messages.Count == 0)
	{
	Console.WriteLine("没有未发送的消息");
	return;
	}
	
	var factory = new ConnectionFactory { Uri = new Uri(_mqConnStr) };
	using var connection = factory.CreateConnection();
	using var channel = connection.CreateModel();
	channel.ExchangeDeclare("local_message_exchange", ExchangeType.Direct, durable: true);
	channel.QueueDeclare("stock_deduct_queue", durable: true, exclusive: false, autoDelete: false);
	channel.QueueBind("stock_deduct_queue", "local_message_exchange", "stock.deduct");
	
	foreach (var message in messages)
	{
	try
	{
	// 1. 推送消息到RabbitMQ
	var body = Encoding.UTF8.GetBytes(message.Content);
	var properties = channel.CreateBasicProperties();
	properties.Persistent = true;
	properties.MessageId = message.TransactionId;
	
	channel.BasicPublish("local_message_exchange", "stock.deduct", properties, body);
	
	// 2. 标记消息为已发送
	await UpdateMessageStatusAsync(message.Id, 1);
	Console.WriteLine($"[消息推送] 消息 {message.TransactionId} 推送成功");
	}
	catch (Exception ex)
	{
	Console.WriteLine($"[消息推送] 消息 {message.TransactionId} 推送失败:{ex.Message}");
	// 失败不标记,下次继续推送
	}
	}
	}
	
	private async Task<List<LocalMessage>> GetUnsentMessagesAsync()
	{
	List<LocalMessage> messages = new List<LocalMessage>();
	using var conn = new MySqlConnection(_connStr);
	await conn.OpenAsync();
	string sql = "SELECT id, transaction_id, content FROM local_messages WHERE status = 0 LIMIT 100";
	using var cmd = new MySqlCommand(sql, conn);
	using var reader = await cmd.ExecuteReaderAsync();
	while (await reader.ReadAsync())
	{
	messages.Add(new LocalMessage
	{
	Id = reader.GetInt64(0),
	TransactionId = reader.GetString(1),
	Content = reader.GetString(2)
	});
	}
	return messages;
	}
	
	private async Task<bool> UpdateMessageStatusAsync(long id, int status)
	{
	using var conn = new MySqlConnection(_connStr);
	await conn.OpenAsync();
	string sql = "UPDATE local_messages SET status = @status, update_time = NOW() WHERE id = @id";
	using var cmd = new MySqlCommand(sql, conn);
	cmd.Parameters.AddWithValue("@id", id);
	cmd.Parameters.AddWithValue("@status", status);
	int rows = await cmd.ExecuteNonQueryAsync();
	return rows > 0;
	}
	}

步骤4:库存服务消费消息(逐行讲解)
csharp

	using System;
	using System.Text;
	using System.Text.Json;
	using System.Threading.Tasks;
	using RabbitMQ.Client;
	using RabbitMQ.Client.Events;
	using StackExchange.Redis;
	
	namespace DistributedTransaction.LocalMessage;
	
	public class StockConsumer
	{
	private readonly string _mqConnStr;
	private readonly IStockRepository _stockRepo;
	private readonly IDatabase _redisDb;
	
	public StockConsumer(string mqConnStr, IStockRepository stockRepo, IConnectionMultiplexer redis)
	{
	_mqConnStr = mqConnStr;
	_stockRepo = stockRepo;
	_redisDb = redis.GetDatabase();
	}
	
	public void StartConsuming()
	{
	var factory = new ConnectionFactory { Uri = new Uri(_mqConnStr) };
	using var connection = factory.CreateConnection();
	using var channel = connection.CreateModel();
	channel.BasicQos(0, 1, false);
	
	var consumer = new EventingBasicConsumer(channel);
	consumer.Received += async (model, ea) =>
	{
	var body = ea.Body.ToArray();
	string content = Encoding.UTF8.GetString(body);
	string messageId = ea.BasicProperties.MessageId;
	Console.WriteLine($"[库存消费] 收到消息 {messageId}:{content}");
	
	try
	{
	// 1. 幂等性检查:判断消息是否已经处理过
	if (await _redisDb.StringGetAsync($"local:message:{messageId}") == "1")
	{
	Console.WriteLine($"[库存消费] 消息 {messageId} 已经处理过,直接返回");
	channel.BasicAck(ea.DeliveryTag, false);
	return;
	}
	
	// 2. 解析消息
	var message = JsonSerializer.Deserialize<StockDeductMessage>(content);
	
	// 3. 扣减库存
	bool success = await _stockRepo.DeductStockAsync(message.ProductId, message.Quantity);
	if (!success)
	{
	throw new Exception($"扣减库存失败:{message.ProductId}");
	}
	
	// 4. 标记消息为已处理
	await _redisDb.StringSetAsync($"local:message:{messageId}", "1", TimeSpan.FromDays(7));
	
	// 5. 确认消息
	channel.BasicAck(ea.DeliveryTag, false);
	Console.WriteLine($"[库存消费] 消息 {messageId} 处理成功");
	}
	catch (Exception ex)
	{
	Console.WriteLine($"[库存消费] 消息 {messageId} 处理失败:{ex.Message}");
	// 6. 重试3次,失败后放到死信队列
	int retryCount = ea.BasicProperties.Headers?.ContainsKey("retry_count") == true ? (int)ea.BasicProperties.Headers["retry_count"] : 0;
	if (retryCount < 3)
	{
	retryCount++;
	ea.BasicProperties.Headers["retry_count"] = retryCount;
	channel.BasicPublish("", "stock_deduct_queue", ea.BasicProperties, body);
	Console.WriteLine($"[库存消费] 消息 {messageId} 重试 {retryCount} 次");
	}
	else
	{
	channel.BasicNack(ea.DeliveryTag, false, false);
	Console.WriteLine($"[库存消费] 消息 {messageId} 重试3次失败,放到死信队列");
	}
	}
	};
	
	channel.BasicConsume("stock_deduct_queue", false, consumer);
	Console.WriteLine("库存服务消费已启动");
	Console.ReadLine();
	}
	}
	
	public class StockDeductMessage
	{
	public string ProductId { get; set; }
	public int Quantity { get; set; }
	public string TransactionId { get; set; }
	}
	
	public interface IStockRepository
	{
	Task<bool> DeductStockAsync(string productId, int quantity);
	}

核心代码逻辑拆解:
1.订单服务:开启数据库事务,创建订单和插入消息在同一个事务中,保证订单和消息的一致性;
2.定时任务:每1分钟扫描未发送的消息,推送到RabbitMQ,失败的消息下次继续推送;
3.库存服务:消费消息,处理幂等性,扣减库存,重试失败的消息,放到死信队列;
4.幂等性处理:用消息ID标记消息是否已经处理过,避免重复消费。
我踩过的坑:死信队列的消息没人处理——必须有监控告警,及时处理死信队列的消息!
拓展知识:本地消息表的优缺点
优点 缺点 适用场景
最终一致性,数据最终一致 消息延迟,数据可能在一段时间后不一致 跨系统数据同步、电商订单、库存同步
业务侵入性低,只需要插入消息 依赖数据库和消息队列,部署复杂 对一致性要求不高的场景
实现简单,适合中小团队 性能一般,不适合超高并发场景
五、分布式事务最佳实践与踩坑总结

  1. 尽量避免分布式事务
    能用本地事务解决的问题,不要用分布式事务;能用最终一致性解决的问题,不要用强一致性。比如电商订单和库存同步,用最终一致性的本地消息表比用强一致性的2PC好。
  2. 选择合适的分布式事务方案
    场景 推荐方案
    强一致性、低并发 2PC(XA协议)
    高并发、业务逻辑复杂 TCC
    跨系统数据同步、最终一致 本地消息表、消息队列(RabbitMQ、Kafka)
    简单场景、最终一致 可靠消息服务(RocketMQ的事务消息)
  3. 处理分布式事务的异常
    超时:设置合理的超时时间,避免系统挂起;
    重试:对失败的操作进行重试,最多重试3次,重试间隔递增;
    补偿:对重试失败的操作进行人工补偿,比如定时任务扫描失败的事务,通知运营人员处理;
    幂等性:每个操作必须有唯一ID,保证操作执行多次结果一致;
    监控告警:监控事务的成功率、失败率、重试次数,及时发现问题。
  4. 踩过的坑总结
    2PC的同步阻塞:大促时系统超时率飙升——高并发场景千万别用2PC;
    TCC的空回滚和悬挂问题:导致数据不一致——必须处理空回滚和悬挂问题;
    本地消息表的重复消费:导致库存扣减两次——必须处理幂等性;
    没有补偿机制:导致数据不一致——定时任务扫描失败的事务,进行人工补偿;
    依赖单点故障:定时任务挂了导致消息没推送——用Quartz.NET实现分布式定时任务,避免单点故障。
    六、总结
    2PC适合强一致性低并发场景,3PC是2PC的改进版但用的很少,TCC适合高并发业务逻辑复杂场景,本地消息表适合跨系统最终一致场景。选择合适的分布式事务方案,处理好异常和幂等性,能保证分布式系统的数据一致性。
    下一节我们会学习分布式追踪:OpenTelemetry的实战,解决分布式系统中的链路追踪问题。

 本站原创,转载请注明出处:https://www.xin3721.com/ArticlecSharp/c49565.html


相关教程