03. 分布式事务

分布式事务深度解析:2PC/3PC、TCC、Saga、本地消息表、Seata AT模式与最终一致性实践

单体架构使用数据库本地事务即可保证 ACID,但微服务架构下业务操作跨越多个服务和数据库,分布式事务成为必须解决的问题。本文系统讲解从强一致性到最终一致性的各种分布式事务方案。

1. 分布式事务分类

方案一致性性能复杂度适用场景
2PC强一致短事务、低并发
3PC强一致2PC 改进,较少使用
TCC最终一致高并发、核心资产业务
Saga最终一致长事务、业务流程
本地消息表最终一致异步场景
Seata AT最终一致较高侵入小、快速接入
最大努力通知最终一致对账场景

2. 两阶段提交 (2PC)

2PC (Two-Phase Commit) 是最经典的分布式事务协议,由协调者 (Coordinator) 和参与者 (Participants) 组成。

2.1 协议流程

Phase 1 (投票阶段):
  Coordinator → Prepare → Participant A
  Coordinator → Prepare → Participant B
  Coordinator → Prepare → Participant C

  A: Undo/Redo log written, return Yes
  B: Undo/Redo log written, return Yes
  C: Failed to prepare, return No

Phase 2 (提交/回滚):
  Coordinator (收到 No) → Rollback → A
  Coordinator (收到 No) → Rollback → B

2.2 2PC 的状态机

Coordinator 状态:
  INIT → PREPARING → PREPARED → COMMITTING → COMMITTED
                              → ABORTING → ABORTED

Participant 状态:
  INIT → READY → COMMITTED/ABORTED

2.3 2PC 的问题

问题原因后果
同步阻塞参与者需锁定资源等待协调者指令性能差
单点故障协调者宕机,参与者一直阻塞可用性低
数据不一致协调者发送 Commit 后宕机,部分参与者未收到脑裂
_TIMEOUT网络超时导致不确定性需人工干预

2.4 Java 实现示例

@Component
public class XAOrderService {
    
    @Autowired
    private JdbcTemplate orderJdbc;
    
    @Autowired
    private JdbcTemplate inventoryJdbc;
    
    @Transactional(rollbackFor = Exception.class)
    public void createOrder(OrderRequest request) {
        // XA 两阶段提交,通过 JTA 管理
        // atomikos 或 narayana 实现
        
        UserTransaction ut = userTransactionManager.getUserTransaction();
        try {
            ut.begin();
            
            // 操作订单库
            orderJdbc.update("INSERT INTO orders (user_id, amount) VALUES (?, ?)",
                request.getUserId(), request.getAmount());
            
            // 操作库存库
            inventoryJdbc.update("UPDATE inventory SET count = count - ? WHERE sku = ?",
                request.getQuantity(), request.getSku());
            
            ut.commit();
        } catch (Exception e) {
            try {
                ut.rollback();
            } catch (Exception ex) {
                log.error("Rollback failed", ex);
            }
            throw new BusinessException("Order creation failed", e);
        }
    }
}

2.5 3PC 改进

3PC 增加了一个 CanCommit 阶段,减少协调者宕机导致的阻塞时间,但实现复杂且网络开销更大,实际很少使用。

CanCommit → PreCommit → DoCommit

3. TCC (Try-Confirm-Cancel)

TCC 由阿里提出,将业务操作拆分为三个阶段,通过业务逻辑保证最终一致性。

3.1 TCC 三阶段

阶段操作资源状态
Try预留资源,执行业务检查资源被冻结/预扣
Confirm真正执行业务资源正式扣除
Cancel释放预留资源回滚到初始状态
用户下单:
  Try: 冻结库存 1,冻结余额 100
       库存: 10 → 可用 9, 冻结 1
       余额: 500 → 可用 400, 冻结 100
  
  Confirm: 确认订单
       库存: 可用 9, 冻结 0, 已售 +1
       余额: 可用 400, 冻结 0, 已扣 100
  
  Cancel: 取消订单
       库存: 可用 10, 冻结 0
       余额: 可用 500, 冻结 0

3.2 TCC 实现

public interface InventoryTccAction {
    
    @TwoPhaseBusinessAction(name = "inventoryTccAction", 
        tryMethod = "tryDeduct", 
        confirmMethod = "commit", 
        rollbackMethod = "rollback")
    boolean tryDeduct(@BusinessActionContextParameter(paramName = "sku") String sku,
                     @BusinessActionContextParameter(paramName = "count") int count);
    
    boolean commit(BusinessActionContext context);
    
    boolean rollback(BusinessActionContext context);
}

@Service
public class InventoryTccActionImpl implements InventoryTccAction {
    
    @Autowired
    private InventoryMapper inventoryMapper;
    
    @Override
    public boolean tryDeduct(String sku, int count) {
        // 检查库存
        Inventory inventory = inventoryMapper.selectBySku(sku);
        if (inventory.getAvailable() < count) {
            throw new BusinessException("Insufficient inventory");
        }
        // 冻结库存
        return inventoryMapper.freeze(sku, count) > 0;
    }
    
    @Override
    public boolean commit(BusinessActionContext context) {
        String sku = context.getActionContext("sku");
        int count = Integer.parseInt(context.getActionContext("count"));
        // 将冻结库存转为实际扣减
        return inventoryMapper.confirmDeduct(sku, count) > 0;
    }
    
    @Override
    public boolean rollback(BusinessActionContext context) {
        String sku = context.getActionContext("sku");
        int count = Integer.parseInt(context.getActionContext("count"));
        // 释放冻结库存
        return inventoryMapper.unfreeze(sku, count) > 0;
    }
}

3.3 TCC 注意事项

  • 幂等性:Confirm 和 Cancel 必须幂等,可能因网络重试而多次执行
  • 空回滚:Try 尚未执行就触发 Cancel,需要记录 Try 是否执行过
  • 悬挂:Cancel 先执行,Try 后执行(超时导致的乱序),需防悬挂
// 幂等性控制
@Override
public boolean commit(BusinessActionContext context) {
    String xid = context.getXid();
    // 检查是否已提交
    if (tccLogMapper.isCommitted(xid)) {
        return true;  // 已处理,直接返回
    }
    // ... 执行业务
    tccLogMapper.markCommitted(xid);
    return true;
}

4. Saga 模式

Saga 将长事务拆分为多个本地事务,每个本地事务提交后立即释放资源,通过补偿操作回滚。

4.1 两种 Saga 实现

类型机制代表框架
编排式 (Choreography)每个服务完成本地事务后发送事件触发下一个服务事件驱动
编排式 (Orchestration)中央协调器统一调度各服务的执行和补偿Camunda, Apache Camel

4.2 编排式 Saga (Orchestration)

@Service
public class OrderSagaOrchestrator {
    
    @Autowired
    private StateMachineFactory<OrderStatus, OrderEvent> stateMachineFactory;
    
    public void startOrderSaga(Order order) {
        StateMachine<OrderStatus, OrderEvent> sm = stateMachineFactory.getStateMachine();
        sm.start();
        
        // 状态流转:
        // CREATED → [create inventory reservation] → INVENTORY_RESERVED
        //         → [create payment] → PAYMENT_COMPLETED
        //         → [ship order] → SHIPPED
        //         → [complete] → COMPLETED
        
        // 补偿链(反向执行):
        // PAYMENT_COMPLETED → [refund payment] → INVENTORY_RESERVED
        //                 → [release inventory] → CREATED
    }
}

// Spring State Machine 配置
@Configuration
public class OrderSagaConfig extends StateMachineConfigurerAdapter<OrderStatus, OrderEvent> {
    
    @Override
    public void configure(StateMachineTransitionConfigurer<OrderStatus, OrderEvent> transitions) 
            throws Exception {
        transitions
            .withExternal()
                .source(OrderStatus.CREATED)
                .target(OrderStatus.INVENTORY_RESERVED)
                .event(OrderEvent.RESERVE_INVENTORY)
                .action(reserveInventoryAction())
                .and()
            .withExternal()
                .source(OrderStatus.INVENTORY_RESERVED)
                .target(OrderStatus.PAYMENT_COMPLETED)
                .event(OrderEvent.PROCESS_PAYMENT)
                .action(processPaymentAction())
                .and()
            // 补偿路径
            .withExternal()
                .source(OrderStatus.PAYMENT_COMPLETED)
                .target(OrderStatus.INVENTORY_RESERVED)
                .event(OrderEvent.REFUND_PAYMENT)
                .action(refundPaymentAction());
    }
}

4.3 事件编排式 Saga

// 订单服务
@Transactional
public void createOrder(OrderRequest request) {
    Order order = orderRepository.save(new Order(request));
    eventPublisher.publish(new OrderCreatedEvent(order.getId(), request));
}

// 库存服务监听
@EventListener
@Transactional
public void onOrderCreated(OrderCreatedEvent event) {
    inventoryService.reserve(event.getSku(), event.getQuantity());
    eventPublisher.publish(new InventoryReservedEvent(event.getOrderId()));
}

// 支付服务监听
@EventListener
@Transactional
public void onInventoryReserved(InventoryReservedEvent event) {
    try {
        paymentService.charge(event.getOrderId());
        eventPublisher.publish(new PaymentCompletedEvent(event.getOrderId()));
    } catch (Exception e) {
        eventPublisher.publish(new PaymentFailedEvent(event.getOrderId()));
    }
}

// 补偿:支付失败释放库存
@EventListener
@Transactional
public void onPaymentFailed(PaymentFailedEvent event) {
    inventoryService.release(event.getOrderId());
    orderService.cancel(event.getOrderId());
}

5. 本地消息表

基于可靠消息实现最终一致性,适用于异步场景。

5.1 核心思想

业务操作和消息记录在同一个本地事务中:
  BEGIN
    INSERT INTO orders (...)  -- 业务表
    INSERT INTO message_queue (topic, payload, status)  -- 消息表
  COMMIT

后台任务轮询消息表,发送到消息队列
消息消费方处理完毕后 ACK,消息表状态更新为 DONE

5.2 实现

@Service
public class ReliableMessageService {
    
    @Transactional
    public void createOrderWithMessage(OrderRequest request) {
        // 1. 保存订单
        Order order = orderRepository.save(new Order(request));
        
        // 2. 记录消息(同库同事务)
        OutboxMessage message = new OutboxMessage();
        message.setTopic("order_created");
        message.setPayload(JsonUtils.toJson(new OrderCreatedEvent(order)));
        message.setStatus(MessageStatus.PENDING);
        message.setRetryCount(0);
        outboxRepository.save(message);
    }
    
    // 定时任务:轮询消息表
    @Scheduled(fixedRate = 5000)
    public void pollOutboxMessages() {
        List<OutboxMessage> pending = outboxRepository
            .findByStatusAndRetryCountLessThan(MessageStatus.PENDING, 3);
        
        for (OutboxMessage msg : pending) {
            try {
                kafkaTemplate.send(msg.getTopic(), msg.getPayload()).get(5, TimeUnit.SECONDS);
                msg.setStatus(MessageStatus.SENT);
            } catch (Exception e) {
                msg.setRetryCount(msg.getRetryCount() + 1);
            }
            outboxRepository.save(msg);
        }
    }
    
    // 消息消费确认
    @KafkaListener(topics = "order_created")
    public void handleOrderCreated(String payload, Acknowledgment ack) {
        OrderCreatedEvent event = JsonUtils.fromJson(payload, OrderCreatedEvent.class);
        
        // 幂等处理
        if (processedMessageRepository.existsByMessageId(event.getMessageId())) {
            ack.acknowledge();
            return;
        }
        
        inventoryService.deduct(event.getSku(), event.getQuantity());
        processedMessageRepository.save(new ProcessedMessage(event.getMessageId()));
        ack.acknowledge();
    }
}

6. Seata AT 模式

Seata (Simple Extensible Autonomous Transaction Architecture) 是阿里开源的分布式事务解决方案,AT 模式对业务零侵入。

6.1 AT 模式原理

1. 一阶段:业务 SQL → 解析 SQL → 查询前镜像 → 执行业务 SQL → 查询后镜像 → 记录 UNDO_LOG
2. 二阶段成功:异步删除 UNDO_LOG
3. 二阶段回滚:用 UNDO_LOG 的前镜像生成反向 SQL 回滚

6.2 部署与配置

# application.yml
seata:
  enabled: true
  application-id: ${spring.application.name}
  tx-service-group: my_tx_group
  service:
    vgroup-mapping:
      my_tx_group: default
  client:
    rm:
      async-commit-buffer-limit: 10000
    tm:
      commit-retry-count: 5
      rollback-retry-count: 5
  datasource:
    proxy-datasource: true

6.3 业务代码(零侵入)

@Service
public class BusinessService {
    
    @Autowired
    private OrderService orderService;
    
    @Autowired
    private StorageService storageService;
    
    @Autowired
    private AccountService accountService;
    
    @GlobalTransactional(name = "create-order", rollbackFor = Exception.class)
    public void createOrder(Order order) {
        // 1. 创建订单
        orderService.create(order);
        
        // 2. 扣减库存
        storageService.deduct(order.getCommodityCode(), order.getCount());
        
        // 3. 扣减账户余额
        accountService.debit(order.getUserId(), order.getMoney());
        
        // 任意步骤异常,全局回滚
    }
}

Seata 代理数据源自动拦截 SQL,生成 UNDO_LOG。

6.4 Seata 四种模式对比

模式侵入性适用场景性能
AT零侵入CRUD 为主较高
TCC高(需实现三接口)复杂业务
Saga中(状态机配置)长事务、业务流程
XA兼容传统 2PC

7. 方案选择决策树

是否需要实时强一致性?
├── 是 → 2PC/XA(短事务、低并发)
│         或:业务层面通过单服务聚合避免分布式事务
└── 否 → 最终一致性
          ├── 是否需要高并发?
          │   ├── 是 → TCC(金融核心、库存扣减)
          │   └── 否 → Saga(业务流程、长事务)
          └── 是否需异步解耦?
              ├── 是 → 本地消息表 / Outbox
              └── 否 → Seata AT(快速接入、零侵入)

总结

维度2PC/XATCCSaga本地消息表Seata AT
一致性最终最终最终最终
性能较高
复杂度
侵入性
回滚能力自动业务补偿业务补偿无法自动回滚自动
适用传统系统核心资产业务流程异步通知通用

分布式事务选型建议:

  1. 能不用就不用:通过业务设计避免分布式事务(如将操作收敛到单一服务)
  2. 最终一致性优先:绝大多数场景最终一致性足够
  3. TCC 用于核心资产:资金、库存等高并发且需精确控制
  4. Seata AT 快速接入:已有系统改造,不想改动业务代码
  5. Saga 处理长事务:订单流程、审批流等跨多服务的业务流程

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「distributed-systems」更多文章

  1. 分布式高可用架构模式:多活、容灾、降级与 K8s 编排高可用
  2. 分布式链路追踪实战:OpenTelemetry、Jaeger 与 W3C Trace Context
  3. 分布式缓存深度策略:Redis Cluster、一致性哈希与多级缓存架构