-
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标记消息是否已经处理过,避免重复消费。
我踩过的坑:死信队列的消息没人处理——必须有监控告警,及时处理死信队列的消息!
拓展知识:本地消息表的优缺点
优点 缺点 适用场景
最终一致性,数据最终一致 消息延迟,数据可能在一段时间后不一致 跨系统数据同步、电商订单、库存同步
业务侵入性低,只需要插入消息 依赖数据库和消息队列,部署复杂 对一致性要求不高的场景
实现简单,适合中小团队 性能一般,不适合超高并发场景
五、分布式事务最佳实践与踩坑总结
-
尽量避免分布式事务
能用本地事务解决的问题,不要用分布式事务;能用最终一致性解决的问题,不要用强一致性。比如电商订单和库存同步,用最终一致性的本地消息表比用强一致性的2PC好。 -
选择合适的分布式事务方案
场景 推荐方案
强一致性、低并发 2PC(XA协议)
高并发、业务逻辑复杂 TCC
跨系统数据同步、最终一致 本地消息表、消息队列(RabbitMQ、Kafka)
简单场景、最终一致 可靠消息服务(RocketMQ的事务消息) -
处理分布式事务的异常
超时:设置合理的超时时间,避免系统挂起;
重试:对失败的操作进行重试,最多重试3次,重试间隔递增;
补偿:对重试失败的操作进行人工补偿,比如定时任务扫描失败的事务,通知运营人员处理;
幂等性:每个操作必须有唯一ID,保证操作执行多次结果一致;
监控告警:监控事务的成功率、失败率、重试次数,及时发现问题。 -
踩过的坑总结
2PC的同步阻塞:大促时系统超时率飙升——高并发场景千万别用2PC;
TCC的空回滚和悬挂问题:导致数据不一致——必须处理空回滚和悬挂问题;
本地消息表的重复消费:导致库存扣减两次——必须处理幂等性;
没有补偿机制:导致数据不一致——定时任务扫描失败的事务,进行人工补偿;
依赖单点故障:定时任务挂了导致消息没推送——用Quartz.NET实现分布式定时任务,避免单点故障。
六、总结
2PC适合强一致性低并发场景,3PC是2PC的改进版但用的很少,TCC适合高并发业务逻辑复杂场景,本地消息表适合跨系统最终一致场景。选择合适的分布式事务方案,处理好异常和幂等性,能保证分布式系统的数据一致性。
下一节我们会学习分布式追踪:OpenTelemetry的实战,解决分布式系统中的链路追踪问题。
本站原创,转载请注明出处:https://www.xin3721.com/ArticlecSharp/c49565.html










