📖 数据密集型设计

分布式事务与最终一致性

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

一、分布式事务概述

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

二、分布式事务挑战

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

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

九、总结

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