最终一致性方案:本地消息表、事务消息、Saga模式

引言

先讲一个我亲身踩过的坑。

2023年,我当时所在的团队负责一个电商中台的订单服务重构。业务逻辑并不复杂:用户下单后,需要扣减库存、锁定优惠券、增加积分。最初是单机应用,一个本地事务包住所有操作,干干净净。

微服务拆分后,订单、库存、优惠券、积分各成独立服务,各自拥有独立数据库。问题来了:用户下单成功后,如果库存扣减失败怎么办?如果优惠券锁定成功但积分增加失败怎么办?

当时我们天真地采用了先更新订单状态,再依次调用其他服务的方案。上线第一个月,就出现了大量数据不一致:有订单显示已支付但库存没扣、有优惠券被锁定但订单已取消。客服工单爆满,凌晨三点被叫起来排查数据。

后来我们调研并落地了最终一致性方案,才真正解决了这个问题。

本文将以订单场景为主线,深入剖析三种主流方案:本地消息表事务消息Saga模式。我会给出源码级分析和完整可运行的代码示例,并对比它们的优缺点和适用场景。

核心概念

生活类比:支付宝转账的"异步记账"

想象一个场景:你在餐厅吃完饭,用支付宝扫码付款。收银台系统需要做两件事:扣你的钱、记餐厅的账。

如果这两个操作必须同时成功或同时失败(强一致),那么支付宝和餐厅系统之间就需要一个分布式事务协调者,整个过程耗时可能从毫秒级变成秒级——你肯定不愿意在收银台等3秒。

现实中的做法是:支付宝先扣你的钱,生成一条转账记录(本地事务),然后异步通知餐厅的账户系统入账。如果通知失败,后台会有补偿任务不断重试,直到餐厅入账成功。你扫码后立刻能走人,餐厅晚几秒收到钱也无所谓。

这就是最终一致性:系统不保证任意时刻数据完全一致,但保证在一段时间后,所有副本最终达到一致状态。

技术定义

在分布式系统中,跨服务的业务操作无法用单一数据库的ACID事务来保证。最终一致性方案的核心思想是:

  1. 将一个大事务拆分为多个本地事务,每个服务在自己的数据库内完成一个本地事务
  2. 通过消息或事件驱动,将本地事务的结果异步传播给其他服务
  3. 引入补偿机制,当后续步骤失败时,回滚前面已成功的操作

三种主流方案的区别在于:如何保证"本地事务"与"消息发送"的原子性,以及如何设计补偿逻辑。

源码/原理深度分析

方案一:本地消息表

这是最经典、也是实现成本最低的方案,由 eBay 提出。核心思想是:将业务操作和消息写入放在同一个本地事务中

#### 核心流程

graph TD A[订单服务] -->|1. 开启本地事务| B[(订单表)] A -->|2. 写入| C[(本地消息表)] B --> D[事务提交成功] C --> D D -->|3. 异步发送消息| E[消息队列] E -->|4. 消费消息| F[库存服务] F -->|5. 扣减库存| G[(库存表)] G -->|6. 处理成功| H[确认消息] H -->|7. 删除本地消息记录| C C -.->|8. 定时任务扫描未确认消息| A A -.->|9. 重新发送| E

#### 关键源码

/**
 * 本地消息表方案 - 订单创建服务
 * 
 * 核心思想:订单数据和消息数据在同一个本地事务中写入,
 * 保证"业务操作"和"消息发送"的原子性。
 */
@Service
public class OrderService {
    
    @Autowired
    private OrderMapper orderMapper;
    
    @Autowired
    private LocalMessageMapper messageMapper;
    
    @Autowired
    private KafkaTemplate<String, String> kafkaTemplate;
    
    /**
     * 创建订单并写入本地消息表
     * 
     * @param order 订单信息
     * @return 订单ID
     */
    @Transactional(rollbackFor = Exception.class)
    public Long createOrder(Order order) {
        // 1. 插入订单数据(本地事务的一部分)
        orderMapper.insert(order);
        
        // 2. 构造本地消息记录
        LocalMessage message = LocalMessage.builder()
            .messageId(UUID.randomUUID().toString())
            .bizType("ORDER_CREATED")
            .bizId(order.getId())
            .payload(JSON.toJSONString(order))
            .status(MessageStatus.PENDING)  // 待发送状态
            .retryCount(0)
            .createTime(new Date())
            .build();
        
        // 3. 插入本地消息表(同一个本地事务!)
        messageMapper.insert(message);
        
        // 4. 注意:这里不直接发送MQ消息!
        // 事务提交后,由定时任务或事务回调来发送
        return order.getId();
    }
    
    /**
     * 事务提交后的回调 - 发送消息
     * 使用 TransactionSynchronizationManager 注册回调,
     * 确保只有在事务真正提交后才发送消息。
     */
    @TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)
    public void sendMessageOnCommit(OrderCreatedEvent event) {
        LocalMessage message = messageMapper.selectById(event.getMessageId());
        if (message != null && message.getStatus() == MessageStatus.PENDING) {
            sendMessageWithRetry(message);
        }
    }
    
    /**
     * 发送消息并更新状态
     */
    private void sendMessageWithRetry(LocalMessage message) {
        try {
            // 发送到MQ
            kafkaTemplate.send("order-topic", message.getPayload()).get(3, TimeUnit.SECONDS);
            // 发送成功,更新状态为已发送
            messageMapper.updateStatus(message.getMessageId(), MessageStatus.SENT);
        } catch (Exception e) {
            log.error("消息发送失败,messageId: {}", message.getMessageId(), e);
            // 发送失败,状态保持PENDING,等待定时任务重试
        }
    }
    
    /**
     * 定时任务 - 扫描并重发超时未确认的消息
     * 
     * 每隔30秒扫描一次,找出超过1分钟仍未发送成功的消息进行重试
     */
    @Scheduled(fixedDelay = 30000)
    public void retryPendingMessages() {
        Date deadline = new Date(System.currentTimeMillis() - 60_000);
        List<LocalMessage> pendingMessages = 
            messageMapper.selectPendingMessages(deadline, 100);
        
        for (LocalMessage message : pendingMessages) {
            // 重试次数限制,超过10次则标记为死信
            if (message.getRetryCount() >= 10) {
                messageMapper.updateStatus(message.getMessageId(), MessageStatus.DEAD);
                continue;
            }
            sendMessageWithRetry(message);
            messageMapper.incrementRetryCount(message.getMessageId());
        }
    }
}

#### 消费端处理(幂等性保证)

/**
 * 消费端 - 库存服务
 * 
 * 关键点:消费端必须实现幂等,防止消息重复投递导致的数据问题
 */
@Component
public class InventoryConsumer {
    
    @Autowired
    private InventoryMapper inventoryMapper;
    
    @Autowired
    private ConsumedMessageMapper consumedMessageMapper;
    
    /**
     * 消费订单创建消息,扣减库存
     * 
     * 幂等方案:使用消费记录表 + 唯一索引
     */
    @KafkaListener(topics = "order-topic", groupId = "inventory-group")
    public void handleOrderCreated(String payload) {
        Order order = JSON.parseObject(payload, Order.class);
        
        // 使用消费记录表做幂等(也可以用Redis SETNX)
        String consumeKey = "ORDER_CREATED:" + order.getId();
        try {
            // 尝试插入消费记录,如果主键冲突说明已处理过
            consumedMessageMapper.insert(ConsumedMessage.builder()
                .consumeKey(consumeKey)
                .consumeTime(new Date())
                .build());
        } catch (DuplicateKeyException e) {
            // 已消费过,直接返回
            log.info("消息重复消费,跳过:{}", consumeKey);
            return;
        }
        
        // 执行实际的库存扣减业务
        int affected = inventoryMapper.deductStock(order.getProductId(), 
                                                   order.getQuantity());
        if (affected == 0) {
            // 库存不足,需要告警并触发人工处理
            // 注意:这里不能直接抛异常,否则消息会重试导致死循环
            // 应该记录失败状态,发送告警,由人工介入
            log.error("库存不足,productId: {}, quantity: {}", 
                      order.getProductId(), order.getQuantity());
        }
    }
}

方案二:事务消息(RocketMQ)

RocketMQ 的事务消息解决了本地消息表方案中的两个痛点:消息表需要额外维护定时任务有延迟

#### 核心原理

graph TD A[订单服务] -->|1. 发送半消息| B[RocketMQ Broker] B -->|2. 半消息保存成功| C[执行本地事务] C -->|3a. 提交| D[Commit消息] C -->|3b. 回滚| E[Rollback消息] C -->|3c. 未知| F[Broker回查] F -->|4. 查询本地事务状态| C D -->|5. 消息可见| G[消费者] G -->|6. 消费| H[库存服务]

#### 核心源码

/**
 * 事务消息方案 - 使用RocketMQ实现
 * 
 * 核心思想:通过MQ的Half Message机制,将"发送消息"和"本地事务"解耦,
 * 由MQ Broker负责事务状态的最终确认。
 */
@Component
public class TransactionalOrderService {
    
    @Autowired
    private OrderMapper orderMapper;
    
    @Autowired
    private RocketMQTemplate rocketMQTemplate;
    
    /**
     * 创建订单 - 事务消息的入口
     */
    public Long createOrderWithTransactionalMessage(Order order) {
        // 1. 构造消息
        Message<String> message = MessageBuilder
            .withPayload(JSON.toJSONString(order))
            .setHeader("orderId", order.getId())
            .build();
        
        // 2. 发送事务消息
        // 注意:这里发送的是Half Message,对消费者不可见
        TransactionSendResult result = rocketMQTemplate.sendMessageInTransaction(
            "order-topic",
            message,
            order,  // 作为参数传给本地事务执行器
            60_000   // 3秒超时(实际生产建议30秒)
        );
        
        if (result.getLocalTransactionState() == LocalTransactionState.COMMIT_MESSAGE) {
            return order.getId();
        } else {
            throw new RuntimeException("订单创建失败,事务已回滚");
        }
    }
    
    /**
     * 本地事务执行器
     * 
     * 这个方法在发送Half Message之后执行,
     * 与Half Message的发送是"伪原子"的。
     */
    @RocketMQTransactionListener
    class OrderTransactionListener implements RocketMQLocalTransactionListener {
        
        /**
         * 执行本地事务
         */
        @Override
        @Transactional(rollbackFor = Exception.class)
        public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
            try {
                Order order = (Order) arg;
                
                // 1. 插入订单数据
                orderMapper.insert(order);
                
                // 2. 如果业务成功,提交消息
                return LocalTransactionState.COMMIT_MESSAGE;
            } catch (Exception e) {
                log.error("本地事务执行失败", e);
                // 业务失败,回滚消息
                return LocalTransactionState.ROLLBACK_MESSAGE;
            }
        }
        
        /**
         * 事务回查
         * 
         * 当Broker长时间未收到Commit/Rollback时,会主动回查。
         * 这里需要根据业务数据判断事务状态。
         */
        @Override
        public LocalTransactionState checkLocalTransaction(Message msg) {
            String orderId = msg.getHeader("orderId");
            Order order = orderMapper.selectById(orderId);
            
            if (order != null) {
                // 订单存在,说明本地事务已提交
                return LocalTransactionState.COMMIT_MESSAGE;
            }
            // 订单不存在,事务已回滚
            return LocalTransactionState.ROLLBACK_MESSAGE;
        }
    }
}

方案三:Saga模式

Saga 模式适合长事务场景,它把一个分布式事务拆分为一系列本地事务,每个本地事务都有对应的补偿操作

#### 核心流程

graph LR A[开始] --> B[T1: 创建订单] B --> C[T2: 扣减库存] C --> D[T3: 锁定优惠券] D --> E[成功] B -. 补偿C1 .-> B1[撤销订单] C -. 补偿C2 .-> C1[回补库存] D -. 补偿C3 .-> D1[解锁优惠券]

#### 核心源码(基于Spring Boot + Saga框架)

/**
 * Saga模式 - 订单创建流程
 * 
 * 使用Seata的Saga模式实现(也可以用Camel或自研)
 * 
 * 关键设计点:
 * 1. 每个步骤都是独立的本地事务
 * 2. 每个步骤都有对应的补偿操作
 * 3. 通过状态机管理整个Saga流程
 */
@Service
public class SagaOrderService {
    
    @Autowired
    private OrderClient orderClient;
    
    @Autowired
    private InventoryClient inventoryClient;
    
    @Autowired
    private CouponClient couponClient;
    
    @Autowired
    private SagaStateMachine stateMachine;
    
    /**
     * 使用Saga状态机定义业务流程
     */
    public void createOrderWithSaga(OrderRequest request) {
        // 定义Saga流程
        SagaDefinition saga = new SagaDefinition()
            // 步骤1:创建订单(正向操作)
            .addStep(
                // 正向操作
                () -> orderClient.createOrder(request),
                // 补偿操作
                (orderId) -> orderClient.cancelOrder(orderId)
            )
            // 步骤2:扣减库存
            .addStep(
                () -> inventoryClient.deductStock(request.getProductId(), 
                                                  request.getQuantity()),
                (result) -> inventoryClient.increaseStock(request.getProductId(), 
                                                          request.getQuantity())
            )
            // 步骤3:锁定优惠券
            .addStep(
                () -> couponClient.lockCoupon(request.getCouponId(), 
                                              request.getUserId()),
                (result) -> couponClient.unlockCoupon(request.getCouponId())
            );
        
        // 执行Saga流程
        try {
            stateMachine.execute(saga);
        } catch (SagaException e) {
            log.error("Saga执行失败,已自动执行补偿", e);
            // 这里可以发起告警,或者记录失败状态等待人工处理
            throw new BusinessException("订单创建失败,已回滚");
        }
    }
}

#### 基于状态机的Saga实现(更可控)

/**
 * 状态机驱动的Saga实现
 * 
 * 比顺序编排更灵活,支持条件分支、并行执行、人工介入等高级特性
 */
@Component
public class OrderSagaStateMachine {
    
    private final StateMachine<OrderState, OrderEvent> stateMachine;
    
    public OrderSagaStateMachine() {
        // 构建状态机
        StateMachineBuilder.Builder<OrderState, OrderEvent> builder = 
            StateMachineBuilder.builder();
        
        // 定义状态转换和动作
        builder.external()
            .from(OrderState.INIT)
            .to(OrderState.ORDER_CREATED)
            .event(OrderEvent.CREATE_ORDER)
            .action(createOrderAction(), cancelOrderAction());
        
        builder.external()
            .from(OrderState.ORDER_CREATED)
            .to(OrderState.STOCK_DEDUCTED)
            .event(OrderEvent.DEDUCT_STOCK)
            .action(deductStockAction(), increaseStockAction());
        
        builder.external()
            .from(OrderState.STOCK_DEDUCTED)
            .to(OrderState.COUPON_LOCKED)
            .event(OrderEvent.LOCK_COUPON)
            .action(lockCouponAction(), unlockCouponAction());
        
        builder.external()
            .from(OrderState.COUPON_LOCKED)
            .to(OrderState.COMPLETED)
            .event(OrderEvent.COMPLETE)
            .action(completeOrderAction());
        
        this.stateMachine = builder.build();
    }
    
    /**
     * 执行状态机
     */
    public void execute(OrderContext context) {
        // 开始执行,如果中途失败,状态机会自动触发补偿
        StateMachine<OrderState, OrderEvent> sm = 
            stateMachine.startWith(OrderState.INIT);
        
        sm.sendEvent(OrderEvent.CREATE_ORDER, context);
        sm.sendEvent(OrderEvent.DEDUCT_STOCK, context);
        sm.sendEvent(OrderEvent.LOCK_COUPON, context);
        sm.sendEvent(OrderEvent.COMPLETE, context);
    }
}

方案对比

| 维度 | 本地消息表 | 事务消息(RocketMQ) | Saga模式 |

|------|-----------|---------------------|----------|

| 实现复杂度 | 中等,需要额外维护消息表 | 较低,依赖MQ的能力 | 较高,需要设计状态机和补偿 |

| 一致性保证 | 最终一致,有秒级延迟 | 最终一致,延迟更低 | 最终一致,可精确控制 |

| 适用场景 | 中小规模系统,对延迟不敏感 | 已经使用RocketMQ的团队 | 长事务、复杂业务流程 |

| 侵入性 | 需要修改业务代码 | 需要引入RocketMQ | 需要引入Saga框架 |

| 运维成本 | 需要维护定时任务 | 依赖MQ的高可用 | 需要监控Saga执行状态 |

| 隔离性 | 支持 | 支持 | 支持 |

| 幂等性 | 需要消费端自行保证 | 需要消费端自行保证 | 需要每个步骤自行保证 |

各方案优缺点详解

本地消息表

  • 优点:不依赖特定中间件,任何MQ都能用;实现思路简单清晰
  • 缺点:消息表会和业务表竞争数据库资源;定时任务有延迟;消息表膨胀后需要归档

事务消息

  • 优点:RocketMQ解决了消息发送和本地事务的原子性问题;延迟更低;不需要定时任务
  • 缺点:强依赖RocketMQ(目前只有RocketMQ支持);RocketMQ的Broker需要额外配置

Saga模式

  • 优点:适合长事务和跨多服务的复杂流程;可以精细控制每个步骤的补偿逻辑
  • 缺点:实现复杂度最高;需要设计完整的补偿逻辑(有时补偿比正向逻辑还难写);没有隔离性保证

最佳实践与避坑指南

最佳实践

  1. 优先选择事务消息

如果你的团队已经在使用RocketMQ,事务消息是最优解。它同时解决了原子性和延迟问题,实现成本最低。

  1. 幂等性是必须的

无论使用哪种方案,消费端幂等都是硬性要求。推荐使用唯一业务键 + 去重表的方式,比Redis SETNX更可靠。

  1. 设置合理的重试策略
  • 消息重试间隔:指数退避(1s → 2s → 4s ...)
  • 最大重试次数:10-15次
  • 超过重试次数:进入死信队列,触发告警
  1. 监控和告警
  • 监控消息积压量、消费延迟、死信队列数量
  • 设置告警阈值,如:消息延迟超过5分钟、死信队列积压超过100条
  1. 人工补偿通道

即使自动化方案再完善,也要保留人工干预的通道。提供管理后台,可以手动重发消息、手动执行补偿。

常见坑

  1. 消息重复消费(幂等没做好)
  • 症状:库存扣成负数、优惠券被重复锁定
  • 解决方案:去重表 + 唯一索引,或者用业务ID做幂等
  1. 消息顺序问题
  • 症状:订单取消消息先于创建消息到达
  • 解决方案:使用分区顺序消息,同一个订单ID发到同一个分区
  1. 补偿操作不幂等
  • 症状:补偿操作执行两次,导致数据错乱
  • 解决方案:补偿操作也要做幂等,比如用补偿记录表
  1. 死循环重试
  • 症状:业务逻辑有问题,消费一直失败,消息无限重试
  • 解决方案:设置最大重试次数,超过后进入死信队列,人工介入
  1. 事务超时
  • 症状:本地事务执行时间过长,超过MQ回查超时时间
  • 解决方案:合理设置超时时间(建议30秒),不要在事务中调用远程服务

总结

分布式系统的事务一致性是一个"没有银弹"的问题。我在实际项目中经历过从"强一致"到"最终一致"的思维转变,这个转变对架构设计的影响是深远的。

回顾三种方案:

  • 本地消息表是最朴素的思想,用"空间换时间",适合对延迟不敏感的场景
  • 事务消息是中间件层面的创新,用"消息代理"解决原子性问题,是当前的最佳性价比选择
  • Saga模式是流程编排的思路,适合复杂的长事务场景,但需要精心的补偿设计

最后想说的是:不要为了分布式事务而分布式事务。如果你的业务可以接受将多个操作放在同一个服务中,用本地事务就够了。微服务拆分是有代价的,每个跨服务操作都需要考虑一致性方案,这会显著增加系统的复杂度。

在实际架构设计中,我通常遵循这样的决策路径:

  1. 能用本地事务解决的,不要拆服务
  2. 必须拆服务的,优先考虑事务消息
  3. 事务消息搞不定的(如长事务),才考虑Saga

希望这篇文章能帮你在分布式系统的道路上少踩一些坑。如果你有更好的方案或不同的见解,欢迎在评论区讨论。