-
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.最终一致性:消息可能延迟,但最终会被处理,数据最终一致。
拓展知识:本地消息表的优缺点
| 优点 | 缺点 | 适用场景 |
|---|---|---|
| 最终一致性,数据最终一致 | 消息延迟,数据可能在一段时间后不一致 | 跨系统数据同步、订单→库存、订单→物流 |
| 业务侵入性低,只需插入消息 | 依赖数据库和消息队列,部署复杂 | 对一致性要求不高的场景 |
| 实现简单,适合中小团队 | 性能一般,不适合超高并发场景 |
五、分布式事务最佳实践与踩坑总结
-
尽量避免分布式事务
能用本地事务解决的问题,别用分布式事务;能用最终一致性解决的问题,别用强一致性。比如电商订单和库存同步,用本地消息表比2PC好太多。
3.选择合适的分布式事务方案
| 场景 | 推荐方案 |
|---|---|
| 强一致性、低并发 | 2PC(XA协议) |
| 高并发、业务逻辑复杂 | TCC |
| 跨系统数据同步、最终一致 | 本地消息表、消息队列(RabbitMQ、Kafka) |
| 简单场景、最终一致 | 可靠消息服务(RocketMQ事务消息) |
-
处理分布式事务的核心原则
幂等性:每个操作必须有唯一ID,确保执行多次结果一致;
重试机制:失败操作最多重试3次,重试间隔递增(1s→3s→5s);
补偿机制:重试失败的操作必须人工补偿,定时任务扫描失败事务;
监控告警:监控事务成功率、失败率、重试次数,异常时及时告警;
超时机制:每个阶段设置超时时间,避免系统长时间挂起。 -
我踩过的坑总结
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










