Spring Boot微服务分布式事务一致性实战:Saga模式与本地消息表的工程化落地方案

微服务分布式事务的困境在哪里?

单体应用拆成微服务后,一个业务操作往往跨多个服务、多个数据库。订单创建要扣库存、扣余额、写流水——每一步都独立提交,任何一步失败都会导致数据不一致。

分布式事务方案里,两阶段提交(2PC)因为同步阻塞和单点故障基本被淘汰,剩下的主流选择是Saga模式本地消息表。这两个方案不是互斥的,在不同的业务场景下各有优势,也可以组合使用。

Saga模式:集中式编排实现

Saga的核心思路是:把一个长事务拆成多个本地事务(Ti),每个本地事务有对应的补偿事务(Ci)。如果第n步失败,逆序执行前面n-1步的补偿操作,把数据回滚到一致状态。

编排方式有两种:事件驱动(Choreography)集中式(Orchestration)。事件驱动在步骤少(2-3步)时简单,步骤多了以后事件链路难以追踪,补偿顺序也容易乱。生产环境更推荐集中式编排。

集中式编排由一个Saga协调器统一控制步骤执行和补偿。以下是基于Spring Boot的实现框架:

// Saga定义:订单创建流程
@Component
public class CreateOrderSaga {

    private final SagaDefinition<OrderState> saga;

    public CreateOrderSaga(OrderService orderService,
                           InventoryService inventoryService,
                           PaymentService paymentService) {
        this.saga = SagaDefinition.<OrderState>builder()
            .step("reserve_inventory")
                .action(ctx -> inventoryService.reserve(ctx.getOrderId(), ctx.getSkuId(), ctx.getQty()))
                .compensate(ctx -> inventoryService.release(ctx.getOrderId(), ctx.getSkuId()))
            .step("deduct_balance")
                .action(ctx -> paymentService.deduct(ctx.getUserId(), ctx.getAmount()))
                .compensate(ctx -> paymentService.refund(ctx.getUserId(), ctx.getAmount()))
            .step("confirm_order")
                .action(ctx -> orderService.confirm(ctx.getOrderId()))
                .compensate(ctx -> orderService.cancel(ctx.getOrderId()))
            .build();
    }

    public SagaDefinition<OrderState> getDefinition() {
        return saga;
    }
}

Saga协调器的核心逻辑——按步骤执行,遇错逆序补偿:

// Saga执行引擎
@Service
public class SagaOrchestrator {

    private final SagaInstanceRepository sagaRepo;
    private final CompensationLogRepository compLogRepo;

    @Transactional
    public void execute(SagaDefinition<?> definition, SagaContext ctx) {
        List<SagaStep> steps = definition.getSteps();
        
        for (int i = 0; i < steps.size(); i++) {
            SagaStep step = steps.get(i);
            try {
                // 执行正向操作
                step.getAction().execute(ctx);
                // 记录已执行步骤,用于补偿
                compLogRepo.save(new CompensationLog(
                    ctx.getSagaId(), i, step.getName(), "COMPLETED"
                ));
            } catch (Exception e) {
                // 逆序补偿
                compensateReverse(steps, i - 1, ctx);
                sagaRepo.updateStatus(ctx.getSagaId(), "COMPENSATED");
                throw new SagaExecutionException("步骤[" + step.getName() + "]执行失败", e);
            }
        }
        sagaRepo.updateStatus(ctx.getSagaId(), "COMPLETED");
    }

    private void compensateReverse(List<SagaStep> steps, int from, SagaContext ctx) {
        for (int i = from; i >= 0; i--) {
            SagaStep step = steps.get(i);
            try {
                step.getCompensate().execute(ctx);
            } catch (Exception e) {
                // 补偿失败需要人工介入,记录日志并告警
                compLogRepo.save(new CompensationLog(
                    ctx.getSagaId(), i, step.getName(), "COMPENSATE_FAILED"
                ));
                alertService.sendAlert("Saga补偿失败", ctx.getSagaId(), i, e.getMessage());
            }
        }
    }
}

补偿幂等保障:最容易被忽视的环节

Saga补偿操作不是”回滚”,而是”业务撤销”——扣库存的补偿是加回库存,扣余额的补偿是退款。补偿操作必须是幂等的,因为:

– 网络超时可能导致补偿被重复执行
– 正向操作可能已经执行成功但响应丢失,协调器误判为失败后触发补偿

保障幂等的手段:

// 补偿操作幂等检查
@Component
public class InventoryService {

    @Transactional
    public void release(String orderId, String skuId) {
        // 幂等键:orderId + skuId,避免重复释放
        if (compensationLogRepository.existsByOrderIdAndAction(orderId, "RELEASE")) {
            log.info("库存释放已执行,跳过: orderId={}", orderId);
            return;
        }
        
        inventoryRepository.addStock(skuId, qty);
        compensationLogRepository.save(new CompensationRecord(
            orderId, "RELEASE", Instant.now()
        ));
    }
}

关键点:补偿日志表和业务操作在同一个本地事务里提交,确保操作和幂等记录要么一起成功要么一起失败。

本地消息表:最终一致性的朴素方案

本地消息表的思路更朴素:业务操作和消息发送放在同一个本地事务里,业务操作成功后,消息也持久化到同一数据库的消息表中。后台定时任务扫描消息表,把未发送的消息投递到MQ,消费端处理后回写状态。

这个方案没有Saga的复杂状态机,实现成本极低,适合对实时性要求不高(秒级延迟可接受)的业务场景。

// 本地消息表 — 业务端
@Service
public class OrderService {

    @Transactional
    public void createOrder(OrderDTO orderDTO) {
        // 1. 写入订单
        Order order = new Order(orderDTO);
        orderMapper.insert(order);

        // 2. 写入本地消息表(同一事务)
        OutboxMessage msg = new OutboxMessage();
        msg.setAggregateId(order.getId());
        msg.setAggregateType("ORDER");
        msg.setEventType("ORDER_CREATED");
        msg.setPayload(JsonUtils.toJson(orderDTO));
        msg.setStatus("PENDING");
        msg.setCreatedAt(Instant.now());
        outboxMapper.insert(msg);
    }
}

消息投递器——定时扫描并发送:

// 消息投递器
@Component
public class OutboxSender {

    @Scheduled(fixedDelay = 1000)  // 每秒扫描一次
    @Transactional
    public void sendPendingMessages() {
        List<OutboxMessage> messages = outboxMapper.selectPending(100);
        for (OutboxMessage msg : messages) {
            try {
                mqProducer.send(msg.getTopic(), msg.getPayload());
                outboxMapper.updateStatus(msg.getId(), "SENT");
            } catch (Exception e) {
                // 发送失败不更新状态,下次重试
                // 超过最大重试次数则标记为DEAD_LETTER
                if (msg.getRetryCount() >= MAX_RETRY) {
                    outboxMapper.updateStatus(msg.getId(), "DEAD_LETTER");
                    alertService.sendAlert("消息投递失败", msg.getId());
                } else {
                    outboxMapper.incrementRetry(msg.getId());
                }
            }
        }
    }
}

消费端的幂等保障通过唯一业务ID去重实现,逻辑和Saga补偿幂等类似,不再赘述。

Saga与本地消息表的组合方案

这两种方案并非对立。实际项目中,可以这样组合:

核心链路用Saga:涉及资金和库存的关键操作,需要精确的补偿语义和实时一致性保障
非核心通知用本地消息表:订单创建后发短信、写操作日志、更新搜索索引等,允许秒级延迟,用消息表异步投递更轻量

组合架构下,Saga正向操作的每个步骤可以同时写入本地消息表——正向操作成功后异步通知下游,补偿操作成功后异步通知回滚。这样既保证了核心链路的一致性,又不会让Saga协调器承担过多的通知职责。

选型决策参考

| 维度 | Saga集中式编排 | 本地消息表 |
|——|—————|———–|
| 一致性模型 | 准实时(毫秒级补偿) | 最终一致(秒级延迟) |
| 实现复杂度 | 中高(需要状态机和补偿框架) | 低(依赖数据库定时任务) |
| 适用步骤数 | 3步以上长流程 | 1-2步简单通知 |
| 运维成本 | 需要监控Saga状态和补偿告警 | 需要监控消息积压和DEAD_LETTER |
| 业务侵入 | 补偿逻辑需在每个服务实现 | 只需写消息表和消费端 |

没有最好的方案,只有最适合的。核心链路的一致性用Saga保障,旁路通知用消息表解耦,这套组合在微服务生产环境里被验证过多次。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/springboot-wei-fu-wu-fen-bu-shi-shi-wu-yi-zhi-xing-shi-zhan/

(0)
小编小编
上一篇 17小时前
下一篇 17小时前

相关推荐