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实现订单与库存强一致性,结果QPS直接从1000跌到80,超时率飙到15%,被运营追着骂;后来换成TCC,结果没处理空回滚,导致库存多了50件,亏了5万;再后来用本地消息表,又因为重复消费导致库存vb.net教程C#教程python教程SQL教程access 2010教程扣减两次,超卖20单。这节我把这些血泪经验揉进去,用大白话讲透4种主流分布式事务方案,结合C#实战代码逐行拆解,拓展生产级优化技巧,让你在一致性、可用性、性能之间找到平衡。
二、2PC(两阶段提交):强一致性的“性能杀手”,只适合低并发
2PC(Two-Phase Commit)是最经典的强一致性协议,核心是“全员确认后再执行”,适合对一致性要求极高的场景(如银行转账),但性能极差,高并发场景千万别碰!
核心原理(大白话+银行转账例子)
把2PC比作银行转账的“双人复核制”:
1.第一阶段(准备):总行(协调者)通知A银行、B银行(参与者)准备转账,A银行冻结用户1000元,B银行准备接收1000元,都回复“准备好了”;
2.第二阶段(提交):总行看到所有人都准备好了,通知A银行扣钱、B银行加钱;如果有一家说“没准备好”,总行通知所有人回滚。
我踩过的坑:大促时协调者处理不过来,所有参与者都在等指令,系统直接卡成PPT——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 TwoPcOrderService
	{
	private readonly string _orderConnStr; // 订单库连接字符串
	private readonly string _stockConnStr; // 库存库连接字符串
	
	public TwoPcOrderService(string orderConnStr, string stockConnStr)
	{
	_orderConnStr = orderConnStr;
	_stockConnStr = stockConnStr;
	}
	
	/// <summary>
	/// 2PC实现订单创建+库存扣减强一致性
	/// </summary>
	public async Task<bool> CreateOrderWithStockDeductAsync(string orderId, string productId, int quantity)
	{
	MySqlXaTransaction xaTx = 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事务,关联两个连接
	xaTx = new MySqlXaTransaction();
	xaTx.Enlist(orderConn); // 把订单库加入XA事务
	xaTx.Enlist(stockConn); // 把库存库加入XA事务
	
	// 3. 第一阶段:准备事务(所有参与者锁定资源)
	await xaTx.PrepareAsync();
	Console.WriteLine("[2PC] 第一阶段:所有参与者准备完成,资源已锁定");
	
	// 4. 执行业务操作(必须在XA事务内)
	bool orderOk = await CreateOrderAsync(orderConn, orderId, productId, quantity);
	bool stockOk = await DeductStockAsync(stockConn, productId, quantity);
	
	if (!orderOk || !stockOk)
	{
	throw new Exception("业务操作失败,触发回滚");
	}
	
	// 5. 第二阶段:提交事务(所有参与者执行操作)
	await xaTx.CommitAsync();
	Console.WriteLine("[2PC] 第二阶段:事务提交成功,数据一致");
	return true;
	}
	catch (Exception ex)
	{
	// 6. 任何失败都回滚事务
	if (xaTx != null)
	{
	await xaTx.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, @pid, @qty, 'created')";
	using var cmd = new MySqlCommand(sql, conn);
	cmd.Parameters.AddWithValue("@id", orderId);
	cmd.Parameters.AddWithValue("@pid", productId);
	cmd.Parameters.AddWithValue("@qty", 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 - @qty WHERE product_id = @pid AND quantity >= @qty";
	using var cmd = new MySqlCommand(sql, conn);
	cmd.Parameters.AddWithValue("@pid", productId);
	cmd.Parameters.AddWithValue("@qty", quantity);
	int rows = await cmd.ExecuteNonQueryAsync();
	return rows > 0;
	}
	}
	
	// 测试代码
	class Program
	{
	static async Task Main(string[] args)
	{
	string orderConn = "server=localhost;user=root;password=123456;database=order_db";
	string stockConn = "server=localhost;user=root;password=123456;database=stock_db";
	var service = new TwoPcOrderService(orderConn, stockConn);
	
	bool success = await service.CreateOrderWithStockDeductAsync("order_123", "product_456", 10);
	Console.WriteLine(success ? "事务成功" : "事务失败");
	}
	}

核心代码拆解:
1.XA事务关联:用MySqlXaTransaction把两个数据库连接绑在一起,保证操作原子性;
2.两阶段提交:先PrepareAsync()锁定资源,再CommitAsync()执行操作,失败则RollbackAsync();
3.业务操作:创建订单和扣减库存必须在XA事务内执行,确保要么都成功,要么都失败;
4.性能问题:XA事务会锁定资源直到提交,大促时会导致资源长时间占用,性能暴跌。
拓展知识:2PC的优缺点与适用场景

优点 缺点 适用场景
强一致性,数据完全一致 同步阻塞,性能极差 银行转账、支付系统、低并发场景
实现简单,依赖数据库XA协议 协调者单点故障,导致系统挂起 对一致性要求极高的场景

适合关系型数据库 不适合高并发场景
三、3PC(三阶段提交):2PC的“鸡肋改进”,实际用得少
3PC(Three-Phase Commit)是2PC的改进版,核心是减少同步阻塞,引入超时机制,但还是存在一致性问题,实际生产中几乎没人用。
核心原理(大白话+婚礼筹备例子)
把3PC比作婚礼筹备的“三步确认制”:
1.CanCommit:策划师问酒店、婚庆“能来吗?”,对方只说“能/不能”,不预留资源;
2.PreCommit:所有人说“能”,策划师通知大家预留资源;有人说“不能”,直接取消;
3.DoCommit:所有人预留好资源,策划师通知执行婚礼;如果超时,酒店自动执行婚礼(默认协调者已通知)。
优缺点与适用场景
优点 缺点 适用场景
减少同步阻塞,CanCommit阶段不锁定资源 一致性问题:超时自动执行可能与协调者指令冲突 几乎不用,除非对一致性和性能都有要求
引入超时机制,解决协调者单点故障 实现复杂,很少有数据库支持
比2PC性能好 还是不适合高并发场景
注意:3PC解决了2PC的部分问题,但还是有一致性风险,实际生产中优先用TCC或本地消息表,别折腾3PC!
四、TCC(Try-Confirm-Cancel):高并发场景的“首选”,业务侵入性强
TCC(Try-Confirm-Cancel)是业务层面的分布式事务协议,核心是“预留资源+确认执行+取消预留”,适合高并发、业务逻辑复杂的场景(如电商订单),但需要修改业务代码,侵入性强。
核心原理(大白话+酒店预订例子)
把TCC比作酒店预订的“三步操作”:
1.Try:预订酒店房间,预留但不实际支付;
2.Confirm:确认预订,实际支付并锁定房间;
3.Cancel:取消预订,释放预留的房间。
我踩过的坑:空回滚问题——Try阶段没执行,Cancel阶段执行了,导致酒店房间多了一间,亏了500块——必须处理空回滚、悬挂问题和幂等性!
实战2:TCC实现订单与库存最终一致性(C#)
用C#实现TCC的三个阶段,处理空回滚、悬挂问题、幂等性。
步骤1:定义TCC接口
csharp

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

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

	using System;
	using System.Threading.Tasks;
	using StackExchange.Redis;
	
	namespace DistributedTransaction.TCC;
	
	public class OrderTccService : IOrderTccService
	{
	private readonly IDatabase _redis; // Redis用于幂等性检查
	private readonly IOrderRepo _orderRepo; // 订单仓储
	private readonly IStockRepo _stockRepo; // 库存仓储
	
	public OrderTccService(IConnectionMultiplexer redis, IOrderRepo orderRepo, IStockRepo stockRepo)
	{
	_redis = redis.GetDatabase();
	_orderRepo = orderRepo;
	_stockRepo = stockRepo;
	}
	
	/// <summary>
	/// Try阶段:预留订单和库存
	/// </summary>
	public async Task<bool> TryReserveAsync(string txId, string orderId, string productId, int quantity)
	{
	// 1. 幂等性检查:避免重复执行Try
	if (await _redis.StringGetAsync($"tcc:try:{txId}") == "1")
	{
	Console.WriteLine($"[TCC-Try] 事务 {txId} 已执行过Try,直接返回成功");
	return true;
	}
	
	// 2. 悬挂问题处理:如果Cancel已执行,Try不能再执行
	if (await _redis.StringGetAsync($"tcc:cancel:{txId}") == "1")
	{
	Console.WriteLine($"[TCC-Try] 事务 {txId} 已执行Cancel,Try拒绝执行");
	return false;
	}
	
	try
	{
	// 3. 预留订单:创建待支付订单
	bool orderReserve = await _orderRepo.ReserveOrderAsync(orderId, productId, quantity);
	if (!orderReserve) throw new Exception("订单预留失败");
	
	// 4. 预留库存:冻结库存
	bool stockReserve = await _stockRepo.ReserveStockAsync(productId, quantity);
	if (!stockReserve) throw new Exception("库存预留失败");
	
	// 5. 标记Try已执行
	await _redis.StringSetAsync($"tcc:try:{txId}", "1", TimeSpan.FromDays(1));
	Console.WriteLine($"[TCC-Try] 事务 {txId} 预留成功");
	return true;
	}
	catch (Exception ex)
	{
	// 6. 任何失败都回滚已执行的操作
	await _orderRepo.CancelReserveOrderAsync(orderId);
	await _stockRepo.CancelReserveStockAsync(productId);
	Console.WriteLine($"[TCC-Try] 事务 {txId} 预留失败:{ex.Message}");
	return false;
	}
	}
	
	/// <summary>
	/// Confirm阶段:确认订单和库存扣减
	/// </summary>
	public async Task<bool> ConfirmAsync(string txId)
	{
	// 1. 幂等性检查:避免重复执行Confirm
	if (await _redis.StringGetAsync($"tcc:confirm:{txId}") == "1")
	{
	Console.WriteLine($"[TCC-Confirm] 事务 {txId} 已执行过Confirm,直接返回成功");
	return true;
	}
	
	// 2. 空回滚处理:如果Try没执行,直接返回成功
	if (await _redis.StringGetAsync($"tcc:try:{txId}") != "1")
	{
	await _redis.StringSetAsync($"tcc:confirm:{txId}", "1", TimeSpan.FromDays(1));
	Console.WriteLine($"[TCC-Confirm] 事务 {txId} Try未执行,空回滚");
	return true;
	}
	
	try
	{
	// 3. 确认订单:待支付→已支付
	bool orderConfirm = await _orderRepo.ConfirmOrderAsync(txId);
	if (!orderConfirm) throw new Exception("订单确认失败");
	
	// 4. 确认库存:冻结→实际扣减
	bool stockConfirm = await _stockRepo.ConfirmStockAsync(txId);
	if (!stockConfirm) throw new Exception("库存确认失败");
	
	// 5. 标记Confirm已执行,清理Try标记
	await _redis.StringSetAsync($"tcc:confirm:{txId}", "1", TimeSpan.FromDays(1));
	await _redis.KeyDeleteAsync($"tcc:try:{txId}");
	Console.WriteLine($"[TCC-Confirm] 事务 {txId} 确认成功");
	return true;
	}
	catch (Exception ex)
	{
	// 6. 确认失败需要人工补偿,因为资源已锁定
	Console.WriteLine($"[TCC-Confirm] 事务 {txId} 确认失败,需人工补偿:{ex.Message}");
	return false;
	}
	}
	
	/// <summary>
	/// Cancel阶段:取消预留的订单和库存
	/// </summary>
	public async Task<bool> CancelAsync(string txId)
	{
	// 1. 幂等性检查:避免重复执行Cancel
	if (await _redis.StringGetAsync($"tcc:cancel:{txId}") == "1")
	{
	Console.WriteLine($"[TCC-Cancel] 事务 {txId} 已执行过Cancel,直接返回成功");
	return true;
	}
	
	// 2. 空回滚处理:如果Try没执行,直接返回成功
	if (await _redis.StringGetAsync($"tcc:try:{txId}") != "1")
	{
	await _redis.StringSetAsync($"tcc:cancel:{txId}", "1", TimeSpan.FromDays(1));
	Console.WriteLine($"[TCC-Cancel] 事务 {txId} Try未执行,空回滚");
	return true;
	}
	
	try
	{
	// 3. 取消订单:删除待支付订单
	bool orderCancel = await _orderRepo.CancelReserveOrderAsync(txId);
	if (!orderCancel) throw new Exception("订单取消失败");
	
	// 4. 取消库存:解冻冻结的库存
	bool stockCancel = await _stockRepo.CancelReserveStockAsync(txId);
	if (!stockCancel) throw new Exception("库存取消失败");
	
	// 5. 标记Cancel已执行,清理Try标记
	await _redis.StringSetAsync($"tcc:cancel:{txId}", "1", TimeSpan.FromDays(1));
	await _redis.KeyDeleteAsync($"tcc:try:{txId}");
	Console.WriteLine($"[TCC-Cancel] 事务 {txId} 取消成功");
	return true;
	}
	catch (Exception ex)
	{
	// 6. 取消失败需要人工补偿
	Console.WriteLine($"[TCC-Cancel] 事务 {txId} 取消失败,需人工补偿:{ex.Message}");
	return false;
	}
	}
	}
	
	// 模拟仓储接口
	public interface IOrderRepo
	{
	Task<bool> ReserveOrderAsync(string orderId, string productId, int quantity);
	Task<bool> CancelReserveOrderAsync(string orderId);
	Task<bool> ConfirmOrderAsync(string txId);
	}
	
	public interface IStockRepo
	{
	Task<bool> ReserveStockAsync(string productId, int quantity);
	Task<bool> CancelReserveStockAsync(string productId);
	Task<bool> ConfirmStockAsync(string txId);
	}

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

	using System;
	
	namespace DistributedTransaction.LocalMessage;
	
	public class LocalMessage
	{
	public long Id { get; set; }
	public string MsgId { 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 MySqlConnector;
	using System;
	using System.Text.Json;
	using System.Threading.Tasks;
	
	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 tx = await conn.BeginTransactionAsync();
	
	try
	{
	// 1. 创建订单
	bool orderOk = await InsertOrderAsync(conn, tx, orderId, productId, quantity);
	if (!orderOk) throw new Exception("创建订单失败");
	
	// 2. 插入本地消息(与订单同事务,确保一致性)
	string msgId = Guid.NewGuid().ToString("N");
	var msgContent = new { OrderId = orderId, ProductId = productId, Quantity = quantity };
	string content = JsonSerializer.Serialize(msgContent);
	bool msgOk = await InsertLocalMessageAsync(conn, tx, msgId, "stock.deduct", content);
	if (!msgOk) throw new Exception("插入消息失败");
	
	// 3. 提交事务
	await tx.CommitAsync();
	Console.WriteLine($"[订单服务] 订单 {orderId} 创建成功,消息已插入");
	return true;
	}
	catch (Exception ex)
	{
	// 4. 回滚事务
	await tx.RollbackAsync();
	Console.WriteLine($"[订单服务] 订单 {orderId} 创建失败:{ex.Message}");
	return false;
	}
	}
	
	private async Task<bool> InsertOrderAsync(MySqlConnection conn, MySqlTransaction tx, string orderId, string productId, int quantity)
	{
	string sql = "INSERT INTO orders (id, product_id, quantity, status) VALUES (@id, @pid, @qty, 'created')";
	using var cmd = new MySqlCommand(sql, conn, tx);
	cmd.Parameters.AddWithValue("@id", orderId);
	cmd.Parameters.AddWithValue("@pid", productId);
	cmd.Parameters.AddWithValue("@qty", quantity);
	return await cmd.ExecuteNonQueryAsync() > 0;
	}
	
	private async Task<bool> InsertLocalMessageAsync(MySqlConnection conn, MySqlTransaction tx, string msgId, string topic, string content)
	{
	string sql = @"INSERT INTO local_messages (msg_id, topic, content, status, create_time, update_time) 
	VALUES (@msgId, @topic, @content, 0, NOW(), NOW())";
	using var cmd = new MySqlCommand(sql, conn, tx);
	cmd.Parameters.AddWithValue("@msgId", msgId);
	cmd.Parameters.AddWithValue("@topic", topic);
	cmd.Parameters.AddWithValue("@content", content);
	return await cmd.ExecuteNonQueryAsync() > 0;
	}
	}

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

	using MySqlConnector;
	using RabbitMQ.Client;
	using System;
	using System.Collections.Generic;
	using System.Text;
	using System.Threading;
	using System.Threading.Tasks;
	
	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_msg_exchange", ExchangeType.Direct, durable: true);
	channel.QueueDeclare("stock_deduct_queue", durable: true, exclusive: false, autoDelete: false);
	channel.QueueBind("stock_deduct_queue", "local_msg_exchange", "stock.deduct");
	
	foreach (var msg in messages)
	{
	try
	{
	// 1. 推送消息到RabbitMQ
	var body = Encoding.UTF8.GetBytes(msg.Content);
	var props = channel.CreateBasicProperties();
	props.Persistent = true; // 持久化消息
	props.MessageId = msg.MsgId;
	
	channel.BasicPublish("local_msg_exchange", "stock.deduct", props, body);
	
	// 2. 标记消息为已发送
	await UpdateMessageStatusAsync(msg.Id, 1);
	Console.WriteLine($"[消息推送] 消息 {msg.MsgId} 推送成功");
	}
	catch (Exception ex)
	{
	Console.WriteLine($"[消息推送] 消息 {msg.MsgId} 推送失败:{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, msg_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),
	MsgId = reader.GetString(1),
	Content = reader.GetString(2)
	});
	}
	return messages;
	}
	
	private async Task<bool> UpdateMessageStatusAsync(long msgId, 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", msgId);
	cmd.Parameters.AddWithValue("@status", status);
	return await cmd.ExecuteNonQueryAsync() > 0;
	}
	}

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

	using RabbitMQ.Client;
	using RabbitMQ.Client.Events;
	using StackExchange.Redis;
	using System;
	using System.Text;
	using System.Text.Json;
	using System.Threading.Tasks;
	
	namespace DistributedTransaction.LocalMessage;
	
	public class StockConsumer
	{
	private readonly string _mqConnStr;
	private readonly IStockRepo _stockRepo;
	private readonly IDatabase _redis;
	
	public StockConsumer(string mqConnStr, IStockRepo stockRepo, IConnectionMultiplexer redis)
	{
	_mqConnStr = mqConnStr;
	_stockRepo = stockRepo;
	_redis = 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 msgId = ea.BasicProperties.MessageId;
	Console.WriteLine($"[库存消费] 收到消息 {msgId}:{content}");
	
	try
	{
	// 1. 幂等性检查:避免重复消费
	if (await _redis.StringGetAsync($"local:msg:{msgId}") == "1")
	{
	Console.WriteLine($"[库存消费] 消息 {msgId} 已处理过,直接确认");
	channel.BasicAck(ea.DeliveryTag, false);
	return;
	}
	
	// 2. 解析消息
	var msg = JsonSerializer.Deserialize<StockDeductMsg>(content);
	
	// 3. 扣减库存
	bool success = await _stockRepo.DeductStockAsync(msg.ProductId, msg.Quantity);
	if (!success) throw new Exception($"扣减库存失败:{msg.ProductId}");
	
	// 4. 标记消息为已处理
	await _redis.StringSetAsync($"local:msg:{msgId}", "1", TimeSpan.FromDays(7));
	
	// 5. 确认消息,从队列删除
	channel.BasicAck(ea.DeliveryTag, false);
	Console.WriteLine($"[库存消费] 消息 {msgId} 处理成功");
	}
	catch (Exception ex)
	{
	Console.WriteLine($"[库存消费] 消息 {msgId} 处理失败:{ex.Message}");
	// 6. 重试3次,失败则放到死信队列
	int retryCount = ea.BasicProperties.Headers?.ContainsKey("retry") == true ? (int)ea.BasicProperties.Headers["retry"] : 0;
	if (retryCount < 3)
	{
	retryCount++;
	ea.BasicProperties.Headers["retry"] = retryCount;
	channel.BasicPublish("", "stock_deduct_queue", ea.BasicProperties, body);
	Console.WriteLine($"[库存消费] 消息 {msgId} 重试 {retryCount} 次");
	}
	else
	{
	channel.BasicNack(ea.DeliveryTag, false, false); // 拒绝并放入死信队列
	Console.WriteLine($"[库存消费] 消息 {msgId} 重试3次失败,放入死信队列");
	}
	}
	};
	
	channel.BasicConsume("stock_deduct_queue", false, consumer);
	Console.WriteLine("库存服务消费已启动");
	Console.ReadLine();
	}
	}
	
	public class StockDeductMsg
	{
	public string OrderId { get; set; }
	public string ProductId { get; set; }
	public int Quantity { get; set; }
	}
	
	public interface IStockRepo
	{
	Task<bool> DeductStockAsync(string productId, int quantity);
	}

核心代码拆解:
1.订单服务:订单和消息同事务,确保订单创建成功则消息必存在;
2.定时任务:定期推送未发送的消息,失败则下次重试,保证消息最终会被推送;
3.库存服务:消费消息时处理幂等性,重试失败的消息,放入死信队列人工处理;
4.最终一致性:消息可能延迟,但最终会被处理,数据最终一致。
拓展知识:本地消息表的优缺点

优点 缺点 适用场景
最终一致性,数据最终一致 消息延迟,数据可能在一段时间后不一致 跨系统数据同步、订单→库存、订单→物流
业务侵入性低,只需插入消息 依赖数据库和消息队列,部署复杂 对一致性要求不高的场景
实现简单,适合中小团队 性能一般,不适合超高并发场景  

五、分布式事务最佳实践与踩坑总结

  1. 尽量避免分布式事务
    能用本地事务解决的问题,别用分布式事务;能用最终一致性解决的问题,别用强一致性。比如电商订单和库存同步,用本地消息表比2PC好太多。
    3.选择合适的分布式事务方案
场景 推荐方案
强一致性、低并发 2PC(XA协议)
高并发、业务逻辑复杂 TCC
跨系统数据同步、最终一致 本地消息表、消息队列(RabbitMQ、Kafka)
简单场景、最终一致 可靠消息服务(RocketMQ事务消息)
  1. 处理分布式事务的核心原则
    幂等性:每个操作必须有唯一ID,确保执行多次结果一致;
    重试机制:失败操作最多重试3次,重试间隔递增(1s→3s→5s);
    补偿机制:重试失败的操作必须人工补偿,定时任务扫描失败事务;
    监控告警:监控事务成功率、失败率、重试次数,异常时及时告警;
    超时机制:每个阶段设置超时时间,避免系统长时间挂起。
  2. 我踩过的坑总结
    1.2PC性能问题:大促时QPS从1000跌到100——高并发场景千万别用2PC;
    2.TCC空回滚:导致库存多了50件——必须在Confirm/Cancel阶段检查Try是否执行过;
    3.本地消息表重复消费:导致库存扣减两次——必须用消息ID做幂等性检查;
    4.协调者单点故障:2PC协调者挂了导致系统挂起——用ZooKeeper实现协调者高可用;
    5.消息推送失败:定时任务挂了导致消息没推送——用Quartz.NET实现分布式定时任务。
    六、总结
    分布式事务没有银弹,要根据场景选择合适的方案:强一致性低并发用2PC,高并发复杂业务用TCC,跨系统同步用本地消息表。尽量避免分布式事务,能用最终一致性就别用强一致性,处理好幂等性、重试、补偿机制,才能保证分布式系统的数据一致性和可用性。
    下一节我们会学习分布式追踪:OpenTelemetry的实战,解决分布式系统中的链路追踪问题。

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


相关教程