最终一致性方案:本地消息表、事务消息、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

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