📖 数据密集型设计

分布式事务与最终一致性

深入探讨分布式事务的实现方案与最终一致性设计

一、分布式事务概述

分布式事务是指涉及多个独立数据源或服务的事务。在分布式系统中,由于网络分区、节点故障等问题,实现强一致性的分布式事务非常困难,因此通常采用最终一致性的方案。

二、分布式事务挑战

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 手动补偿

对于无法自动恢复的失败事务,提供手动补偿机制。

九、总结

分布式事务是分布式系统的核心挑战,最终一致性是实际应用中的主流选择。根据业务需求选择合适的方案,实现幂等性和重试机制,能够构建可靠的分布式事务系统。