一、分布式事务概述
分布式事务是指涉及多个独立数据源或服务的事务。在分布式系统中,由于网络分区、节点故障等问题,实现强一致性的分布式事务非常困难,因此通常采用最终一致性的方案。
二、分布式事务挑战
2.1 CAP理论
CAP理论指出,分布式系统无法同时满足一致性、可用性和分区容错性:
graph TD
A[CAP理论] --> B[一致性C]
A --> C[可用性A]
A --> D[分区容错P]
B --> B1[所有节点数据一致]
C --> C1[任何时候都能响应]
D --> D1[网络分区时继续工作]
E[分布式系统] --> D
E --> F{选择}
F --> G[CP: 牺牲可用性]
F --> H[AP: 牺牲一致性]
2.2 BASE理论
BASE理论是CAP理论的延伸,强调最终一致性:
- 基本可用(Basically Available):系统保证基本的可用性
- 软状态(Soft State):允许系统存在中间状态
- 最终一致性(Eventually Consistent):数据最终会达到一致
三、两阶段提交(2PC)
3.1 2PC流程
两阶段提交是传统的分布式事务协议:
sequenceDiagram
participant TM as 事务管理器
participant RM1 as 资源管理器1
participant RM2 as 资源管理器2
participant RM3 as 资源管理器3
TM->>RM1: Prepare
RM1-->>TM: Prepared
TM->>RM2: Prepare
RM2-->>TM: Prepared
TM->>RM3: Prepare
RM3-->>TM: Prepared
TM->>RM1: Commit
RM1-->>TM: Committed
TM->>RM2: Commit
RM2-->>TM: Committed
TM->>RM3: Commit
RM3-->>TM: Committed
3.2 2PC问题
| 问题 | 描述 | 影响 |
|---|---|---|
| 同步阻塞 | 所有参与者在等待时阻塞 | 性能差 |
| 单点故障 | 事务管理器故障导致阻塞 | 可用性低 |
| 数据不一致 | Commit阶段部分成功 | 数据不一致 |
| 脑裂 | 网络分区导致状态不一致 | 数据不一致 |
四、三阶段提交(3PC)
4.1 3PC流程
三阶段提交在2PC基础上增加了准备阶段:
sequenceDiagram
participant TM as 事务管理器
participant RM as 资源管理器
TM->>RM: CanCommit?
RM-->>TM: Yes
TM->>RM: PreCommit
RM-->>TM: Prepared
TM->>RM: DoCommit
RM-->>TM: Committed
4.2 3PC改进
- 超时机制:减少阻塞时间
- 预准备阶段:检查资源是否可用
- 独立决策:参与者可自行决定提交
4.3 3PC问题
3PC仍然存在数据不一致的问题,且实现复杂。
五、最终一致性方案
5.1 本地消息表
本地消息表是一种基于数据库的最终一致性方案:
flowchart TD
A[业务操作] --> B[写入本地消息表]
B --> C[提交事务]
C --> D[消息发送服务]
D --> E[轮询消息表]
E --> F[发送消息]
F --> G[消息队列]
G --> H[消费端处理]
H --> I{处理成功?}
I --> J[是]
I --> K[否]
J --> L[更新消息状态]
K --> M[重试/死信队列]
5.2 本地消息表实现
-- 创建消息表
CREATE TABLE local_message (
id BIGINT PRIMARY KEY,
message_type VARCHAR(50),
message_content TEXT,
status INT DEFAULT 0,
retry_count INT DEFAULT 0,
created_at DATETIME,
updated_at DATETIME
);
-- 业务操作
BEGIN TRANSACTION;
UPDATE orders SET status = 'paid' WHERE id = 1;
INSERT INTO local_message (id, message_type, message_content) VALUES (1, 'order_paid', '{"order_id": 1}');
COMMIT;
-- 消息发送服务
public class MessageSenderService
{
public void SendMessages()
{
var messages = _db.LocalMessages.Where(m => m.Status == 0).ToList();
foreach (var message in messages)
{
try
{
_messageQueue.Send(message.MessageContent);
message.Status = 1;
}
catch
{
message.RetryCount++;
if (message.RetryCount >= 3)
{
message.Status = 2; // 失败
}
}
_db.SaveChanges();
}
}
}
5.3 Saga模式
Saga模式将分布式事务拆分为多个本地事务:
flowchart TD
A[开始事务] --> B[步骤1: 服务A]
B --> C[步骤2: 服务B]
C --> D[步骤3: 服务C]
D --> E[完成]
B --> B1{失败}
C --> C1{失败}
D --> D1{失败}
B1 --> F[回滚步骤1]
C1 --> G[回滚步骤1和2]
D1 --> H[回滚步骤1、2和3]
5.4 Saga模式实现
public class OrderSaga
{
public async Task Execute(OrderCommand command)
{
// 步骤1: 创建订单
var orderId = await _orderService.CreateOrder(command);
try
{
// 步骤2: 扣减库存
await _inventoryService.DeductStock(command.ProductId, command.Quantity);
try
{
// 步骤3: 支付
await _paymentService.Pay(orderId, command.Amount);
}
catch
{
// 回滚步骤2
await _inventoryService.RevertStock(command.ProductId, command.Quantity);
throw;
}
}
catch
{
// 回滚步骤1
await _orderService.CancelOrder(orderId);
throw;
}
}
}
5.5 TCC模式
TCC模式分为Try、Confirm、Cancel三个阶段:
| 阶段 | 操作 | 说明 |
|---|---|---|
| Try | 预留资源 | 检查并锁定资源 |
| Confirm | 确认提交 | 执行实际操作 |
| Cancel | 取消操作 | 释放预留资源 |
5.6 TCC模式实现
public interface ITccService
{
void Try(string transactionId, object data);
void Confirm(string transactionId, object data);
void Cancel(string transactionId, object data);
}
public class InventoryTccService : ITccService
{
public void Try(string transactionId, object data)
{
// 预留库存
_db.Inventories.Where(i => i.ProductId == productId)
.Update(i => new { i.ReservedStock -= quantity });
}
public void Confirm(string transactionId, object data)
{
// 确认扣减
_db.Inventories.Where(i => i.ProductId == productId)
.Update(i => new { i.Stock -= quantity });
}
public void Cancel(string transactionId, object data)
{
// 释放预留
_db.Inventories.Where(i => i.ProductId == productId)
.Update(i => new { i.ReservedStock += quantity });
}
}
六、消息队列事务
6.1 事务消息
使用消息队列实现分布式事务:
sequenceDiagram
participant App as 应用
participant DB as 数据库
participant MQ as 消息队列
App->>DB: 业务操作
App->>MQ: 发送消息(半消息)
DB-->>App: 操作成功
App->>MQ: 确认消息
MQ-->>App: 消息发送成功
MQ->>Consumer: 投递消息
6.2 RocketMQ事务消息
// 发送事务消息
TransactionMQProducer producer = new TransactionMQProducer("transaction-group");
producer.setTransactionListener(new TransactionListener() {
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
// 执行本地事务
if (executeBusinessLogic()) {
return LocalTransactionState.COMMIT_MESSAGE;
}
return LocalTransactionState.ROLLBACK_MESSAGE;
}
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
// 回查事务状态
if (checkBusinessStatus()) {
return LocalTransactionState.COMMIT_MESSAGE;
}
return LocalTransactionState.UNKNOW;
}
});
producer.start();
七、分布式事务方案对比
| 方案 | 一致性 | 可用性 | 复杂度 | 适用场景 |
|---|---|---|---|---|
| 2PC | 强一致 | 低 | 高 | 传统分布式 |
| 3PC | 强一致 | 中 | 高 | 传统分布式 |
| 本地消息表 | 最终一致 | 高 | 低 | 中小型系统 |
| Saga | 最终一致 | 高 | 中 | 微服务 |
| TCC | 最终一致 | 高 | 高 | 复杂业务 |
| 事务消息 | 最终一致 | 高 | 中 | 异步场景 |
八、最终一致性最佳实践
8.1 选择合适的方案
根据业务需求和系统规模选择合适的分布式事务方案。
8.2 实现幂等性
确保消息消费的幂等性,避免重复处理:
public class MessageConsumer
{
public void Consume(Message message)
{
var exists = _db.MessageLogs.Any(m => m.MessageId == message.Id);
if (exists) return;
try
{
// 处理业务逻辑
ProcessMessage(message);
// 记录消费日志
_db.MessageLogs.Add(new MessageLog { MessageId = message.Id });
_db.SaveChanges();
}
catch
{
// 抛出异常,等待重试
throw;
}
}
}
8.3 设置重试机制
实现重试和死信队列机制:
public void ProcessWithRetry(Message message, int maxRetries = 3)
{
for (int i = 0; i < maxRetries; i++)
{
try
{
ProcessMessage(message);
return;
}
catch
{
Thread.Sleep((int)Math.Pow(2, i) * 1000);
}
}
// 发送到死信队列
_deadLetterQueue.Send(message);
}
8.4 监控事务状态
监控分布式事务的执行状态,及时发现问题。
8.5 手动补偿
对于无法自动恢复的失败事务,提供手动补偿机制。
九、总结
分布式事务是分布式系统的核心挑战,最终一致性是实际应用中的主流选择。根据业务需求选择合适的方案,实现幂等性和重试机制,能够构建可靠的分布式事务系统。