最终一致性方案:本地消息表、事务消息、Saga模式
引言
先讲一个我亲身踩过的坑。
2023年,我当时所在的团队负责一个电商中台的订单服务重构。业务逻辑并不复杂:用户下单后,需要扣减库存、锁定优惠券、增加积分。最初是单机应用,一个本地事务包住所有操作,干干净净。
微服务拆分后,订单、库存、优惠券、积分各成独立服务,各自拥有独立数据库。问题来了:用户下单成功后,如果库存扣减失败怎么办?如果优惠券锁定成功但积分增加失败怎么办?
当时我们天真地采用了先更新订单状态,再依次调用其他服务的方案。上线第一个月,就出现了大量数据不一致:有订单显示已支付但库存没扣、有优惠券被锁定但订单已取消。客服工单爆满,凌晨三点被叫起来排查数据。
后来我们调研并落地了最终一致性方案,才真正解决了这个问题。
本文将以订单场景为主线,深入剖析三种主流方案:本地消息表、事务消息、Saga模式。我会给出源码级分析和完整可运行的代码示例,并对比它们的优缺点和适用场景。
核心概念
生活类比:支付宝转账的"异步记账"
想象一个场景:你在餐厅吃完饭,用支付宝扫码付款。收银台系统需要做两件事:扣你的钱、记餐厅的账。
如果这两个操作必须同时成功或同时失败(强一致),那么支付宝和餐厅系统之间就需要一个分布式事务协调者,整个过程耗时可能从毫秒级变成秒级——你肯定不愿意在收银台等3秒。
现实中的做法是:支付宝先扣你的钱,生成一条转账记录(本地事务),然后异步通知餐厅的账户系统入账。如果通知失败,后台会有补偿任务不断重试,直到餐厅入账成功。你扫码后立刻能走人,餐厅晚几秒收到钱也无所谓。
这就是最终一致性:系统不保证任意时刻数据完全一致,但保证在一段时间后,所有副本最终达到一致状态。
技术定义
在分布式系统中,跨服务的业务操作无法用单一数据库的ACID事务来保证。最终一致性方案的核心思想是:
- 将一个大事务拆分为多个本地事务,每个服务在自己的数据库内完成一个本地事务
- 通过消息或事件驱动,将本地事务的结果异步传播给其他服务
- 引入补偿机制,当后续步骤失败时,回滚前面已成功的操作
三种主流方案的区别在于:如何保证"本地事务"与"消息发送"的原子性,以及如何设计补偿逻辑。
源码/原理深度分析
方案一:本地消息表
这是最经典、也是实现成本最低的方案,由 eBay 提出。核心思想是:将业务操作和消息写入放在同一个本地事务中。
#### 核心流程
#### 关键源码
/**
* 本地消息表方案 - 订单创建服务
*
* 核心思想:订单数据和消息数据在同一个本地事务中写入,
* 保证"业务操作"和"消息发送"的原子性。
*/
@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 的事务消息解决了本地消息表方案中的两个痛点:消息表需要额外维护、定时任务有延迟。
#### 核心原理
#### 核心源码
/**
* 事务消息方案 - 使用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 模式适合长事务场景,它把一个分布式事务拆分为一系列本地事务,每个本地事务都有对应的补偿操作。
#### 核心流程
#### 核心源码(基于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模式
- 优点:适合长事务和跨多服务的复杂流程;可以精细控制每个步骤的补偿逻辑
- 缺点:实现复杂度最高;需要设计完整的补偿逻辑(有时补偿比正向逻辑还难写);没有隔离性保证
最佳实践与避坑指南
最佳实践
- 优先选择事务消息
如果你的团队已经在使用RocketMQ,事务消息是最优解。它同时解决了原子性和延迟问题,实现成本最低。
- 幂等性是必须的
无论使用哪种方案,消费端幂等都是硬性要求。推荐使用唯一业务键 + 去重表的方式,比Redis SETNX更可靠。
- 设置合理的重试策略
- 消息重试间隔:指数退避(1s → 2s → 4s ...)
- 最大重试次数:10-15次
- 超过重试次数:进入死信队列,触发告警
- 监控和告警
- 监控消息积压量、消费延迟、死信队列数量
- 设置告警阈值,如:消息延迟超过5分钟、死信队列积压超过100条
- 人工补偿通道
即使自动化方案再完善,也要保留人工干预的通道。提供管理后台,可以手动重发消息、手动执行补偿。
常见坑
- 消息重复消费(幂等没做好)
- 症状:库存扣成负数、优惠券被重复锁定
- 解决方案:去重表 + 唯一索引,或者用业务ID做幂等
- 消息顺序问题
- 症状:订单取消消息先于创建消息到达
- 解决方案:使用分区顺序消息,同一个订单ID发到同一个分区
- 补偿操作不幂等
- 症状:补偿操作执行两次,导致数据错乱
- 解决方案:补偿操作也要做幂等,比如用补偿记录表
- 死循环重试
- 症状:业务逻辑有问题,消费一直失败,消息无限重试
- 解决方案:设置最大重试次数,超过后进入死信队列,人工介入
- 事务超时
- 症状:本地事务执行时间过长,超过MQ回查超时时间
- 解决方案:合理设置超时时间(建议30秒),不要在事务中调用远程服务
总结
分布式系统的事务一致性是一个"没有银弹"的问题。我在实际项目中经历过从"强一致"到"最终一致"的思维转变,这个转变对架构设计的影响是深远的。
回顾三种方案:
- 本地消息表是最朴素的思想,用"空间换时间",适合对延迟不敏感的场景
- 事务消息是中间件层面的创新,用"消息代理"解决原子性问题,是当前的最佳性价比选择
- Saga模式是流程编排的思路,适合复杂的长事务场景,但需要精心的补偿设计
最后想说的是:不要为了分布式事务而分布式事务。如果你的业务可以接受将多个操作放在同一个服务中,用本地事务就够了。微服务拆分是有代价的,每个跨服务操作都需要考虑一致性方案,这会显著增加系统的复杂度。
在实际架构设计中,我通常遵循这样的决策路径:
- 能用本地事务解决的,不要拆服务
- 必须拆服务的,优先考虑事务消息
- 事务消息搞不定的(如长事务),才考虑Saga
希望这篇文章能帮你在分布式系统的道路上少踩一些坑。如果你有更好的方案或不同的见解,欢迎在评论区讨论。