微服务架构下分布式事务一致性方案选型与实战

分布式事务的核心问题:一致性代价与场景分类

微服务架构下,一个业务操作跨多个服务、多个数据库,事务边界天然被打破。传统单库事务的ACID保证不再适用,分布式事务的本质是在可用性和一致性之间做权衡。没有万能方案,只有匹配场景的取舍。

分布式事务问题可以分为三类场景:跨库写一致性、跨服务写一致性、写与读的最终一致性。每类场景的最优解法完全不同,用2PC解决所有问题是最常见的过度设计。

2PC与3PC:强一致性的理论方案与现实瓶颈

两阶段提交(2PC)是最经典的分布式事务协议。Phase 1协调者通知所有参与者预提交,Phase 2根据预提交结果决定提交或回滚。

2PC的工程实现:Seata AT模式

Seata是阿里开源的分布式事务框架,AT模式是对2PC的工程化改进,对业务代码侵入性最低:

// 1. 引入Seata依赖
// Spring Boot配置
@Configuration
public class SeataConfig {
    @Bean
    public GlobalTransactionScanner globalTransactionScanner() {
        return new GlobalTransactionScanner(
            "order-service", "my-tx-group");
    }
}

// 2. 业务代码中使用@GlobalTransactional
@Service
public class OrderService {
    
    @GlobalTransactional(timeoutMills = 60000, name = "create-order")
    public OrderResult createOrder(OrderRequest request) {
        // 步骤1:创建订单(本地事务)
        Order order = orderMapper.insert(request);
        
        // 步骤2:扣减库存(远程调用)
        inventoryClient.deduct(request.getProductId(), request.getQuantity());
        
        // 步骤3:扣减账户余额(远程调用)
        accountClient.debit(request.getUserId(), order.getTotalAmount());
        
        return OrderResult.success(order);
    }
}

AT模式的工作原理:

  1. 拦截SQL,在执行前保存before-image(修改前数据快照)
  2. 执行业务SQL
  3. 保存after-image(修改后数据快照)
  4. 如果全局事务回滚,用before-image反向补偿

AT模式的问题:before/after image存储在undo_log表中,对数据库有额外写入压力;全局锁在高并发场景下成为瓶颈。

TCC模式:业务补偿的柔性事务

TCC(Try-Confirm-Cancel)将事务拆分为三个阶段,由业务代码实现补偿逻辑,不依赖数据库undo log:

@Service
public class InventoryTccService {
    
    @TwoPhaseBusinessAction(
        name = "deductInventory",
        commitMethod = "confirm",
        rollbackMethod = "cancel"
    )
    public boolean tryDeduct(
        @BusinessActionContextParameter(paramName = "productId") String productId,
        @BusinessActionContextParameter(paramName = "quantity") int quantity
    ) {
        // Try阶段:冻结库存(不是直接扣减)
        int affected = inventoryMapper.freezeStock(productId, quantity);
        if (affected == 0) {
            throw new BusinessException("库存不足");
        }
        return true;
    }
    
    public boolean confirm(BusinessActionContext context) {
        // Confirm阶段:确认扣减冻结库存
        String productId = context.getActionContext("productId", String.class);
        Integer quantity = context.getActionContext("quantity", Integer.class);
        inventoryMapper.confirmDeduct(productId, quantity);
        return true;
    }
    
    public boolean cancel(BusinessActionContext context) {
        // Cancel阶段:释放冻结库存
        String productId = context.getActionContext("productId", String.class);
        Integer quantity = context.getActionContext("quantity", Integer.class);
        inventoryMapper.releaseFrozen(productId, quantity);
        return true;
    }
}

TCC vs AT选择依据

  • AT模式:适合标准CRUD操作,无业务侵入,但全局锁可能成为瓶颈
  • TCC模式:适合非标业务逻辑(冻结/预占场景),无全局锁性能更好,但开发量是3倍

Saga模式:长事务的最终一致性方案

Saga模式将长事务拆分为多个本地事务,每个本地事务完成后触发下一个,任一步失败则逆向执行补偿。适合执行时间长的业务流程(如订单履约、跨机构转账)。

基于事件编排的Saga实现

// 订单创建Saga流程
@Service
public class OrderSagaOrchestrator {
    
    @Autowired
    private KafkaTemplate kafka;
    
    public void startCreateOrderSaga(OrderRequest request) {
        String sagaId = UUID.randomUUID().toString();
        
        // 发起第一步:创建订单
        OrderCreatedEvent event = new OrderCreatedEvent(
            sagaId, request.getOrderId(), 
            request.getProductId(), request.getQuantity()
        );
        kafka.send("order-created", event);
    }
    
    // Step 2: 监听订单创建事件 → 扣减库存
    @KafkaListener(topics = "order-created")
    public void onOrderCreated(OrderCreatedEvent event) {
        try {
            inventoryService.deduct(event.getProductId(), event.getQuantity());
            kafka.send("inventory-deducted", 
                new InventoryDeductedEvent(event.getSagaId()));
        } catch (Exception e) {
            // 库存扣减失败,发起订单取消
            kafka.send("order-cancel-requested", 
                new OrderCancelEvent(event.getSagaId(), event.getOrderId()));
        }
    }
    
    // Step 3: 监听库存扣减事件 → 扣减余额
    @KafkaListener(topics = "inventory-deducted")
    public void onInventoryDeducted(InventoryDeductedEvent event) {
        try {
            accountService.debit(event.getUserId(), event.getAmount());
            kafka.send("account-debited", 
                new AccountDebitedEvent(event.getSagaId()));
        } catch (Exception e) {
            // 余额扣减失败,逆向补偿:恢复库存 + 取消订单
            inventoryService.restore(event.getProductId(), event.getQuantity());
            kafka.send("order-cancel-requested", 
                new OrderCancelEvent(event.getSagaId(), event.getOrderId()));
        }
    }
}

Saga补偿的幂等性保障

补偿操作可能被重复执行(网络重试、消费者rebalance),每个补偿逻辑必须保证幂等:

// 幂等补偿:用唯一键去重
@Transactional
public void restoreInventory(String productId, int quantity, String sagaId) {
    // 检查是否已补偿
    if (compensationLogMapper.existsBySagaId(sagaId)) {
        log.info("补偿已执行,跳过: sagaId={}", sagaId);
        return;
    }
    
    // 执行补偿
    inventoryMapper.restoreStock(productId, quantity);
    
    // 记录补偿日志
    compensationLogMapper.insert(new CompensationLog(sagaId, "INVENTORY_RESTORE"));
}

消息中间件保障最终一致性

很多场景不需要强一致性,只需要”写操作最终生效”。消息中间件是这类场景的最优解。

本地消息表模式

核心思路:业务操作和消息发送在同一个本地事务中完成,消息先写到本地表,由后台任务异步发送到MQ:

@Service
public class PaymentService {
    
    @Transactional
    public void processPayment(PaymentRequest request) {
        // 1. 执行业务操作
        paymentMapper.insert(new Payment(request));
        accountMapper.debit(request.getUserId(), request.getAmount());
        
        // 2. 在同一事务中写入消息表
        outboxMapper.insert(new OutboxMessage(
            "payment-completed",
            JSON.toJSONString(new PaymentCompletedEvent(request)),
            LocalDateTime.now()
        ));
    }
}

// 定时任务扫描消息表并发送到MQ
@Scheduled(fixedDelay = 1000)
public void sendPendingMessages() {
    List messages = outboxMapper.findPending(100);
    for (OutboxMessage msg : messages) {
        try {
            kafka.send(msg.getTopic(), msg.getPayload()).get(5, TimeUnit.SECONDS);
            outboxMapper.markSent(msg.getId());
        } catch (Exception e) {
            outboxMapper.incrementRetry(msg.getId());
            if (msg.getRetryCount() >= 5) {
                outboxMapper.markFailed(msg.getId());
                alertService.notify("消息发送失败: " + msg.getId());
            }
        }
    }
}

API接口规范与事务边界设计

微服务间的API设计直接影响事务复杂度。基本原则:一个API调用对应一个事务边界

// 反模式:一个API内部链式调用多个服务
@PostMapping("/order")
public Result createOrder(@RequestBody OrderRequest req) {
    orderService.create(req);          // 事务1
    inventoryService.deduct(req);      // 事务2(远程调用)
    accountService.debit(req);         // 事务3(远程调用)
    // 任何一个失败,前面的不会自动回滚
}

// 正模式:明确事务边界 + 异步解耦
@PostMapping("/order")
public Result createOrder(@RequestBody OrderRequest req) {
    // 只做本地事务:创建订单 + 写消息表
    orderService.createWithEvent(req);
    // 其余步骤由Saga异步完成
    return Result.accepted("订单已受理,处理中");
}

服务治理:事务监控与故障恢复

分布式事务的运维核心是监控和告警:

// Prometheus监控Seata事务指标
// - seata_transaction_total:事务总数
// - seata_transaction_active:活跃事务数
// - seata_transaction_failed_total:失败事务数

// 关键告警规则
- alert: DistributedTxHighFailureRate
  expr: rate(seata_transaction_failed_total[5m]) / rate(seata_transaction_total[5m]) > 0.05
  labels:
    severity: warning
  annotations:
    summary: "分布式事务失败率超过5%"

// 悬挂事务处理:超时未完成的事务需要人工介入
- alert: SuspiciousLongTx
  expr: seata_transaction_active > 0 and seata_transaction_max_duration_seconds > 120
  labels:
    severity: critical

分布式事务没有银弹。场景决定方案:短事务强一致性用2PC/TCC,长流程最终一致性用Saga,跨服务通知用消息表。关键是识别业务场景的事务边界,然后用最轻量的方案满足一致性要求。过度使用强一致性事务,是微服务架构性能劣化的首要原因。

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

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

相关推荐

微服务架构下分布式事务一致性方案选型与实战

分布式事务的核心问题:一致性代价与场景分类

微服务架构下,一个业务操作跨多个服务、多个数据库,事务边界天然被打破。传统单库事务的ACID保证不再适用,分布式事务的本质是在可用性和一致性之间做权衡。没有万能方案,只有匹配场景的取舍。

分布式事务问题可以分为三类场景:跨库写一致性、跨服务写一致性、写与读的最终一致性。每类场景的最优解法完全不同,用2PC解决所有问题是最常见的过度设计。

2PC与3PC:强一致性的理论方案与现实瓶颈

两阶段提交(2PC)是最经典的分布式事务协议。Phase 1协调者通知所有参与者预提交,Phase 2根据预提交结果决定提交或回滚。

2PC的工程实现:Seata AT模式

Seata是阿里开源的分布式事务框架,AT模式是对2PC的工程化改进,对业务代码侵入性最低:

// 1. 引入Seata依赖
// Spring Boot配置
@Configuration
public class SeataConfig {
    @Bean
    public GlobalTransactionScanner globalTransactionScanner() {
        return new GlobalTransactionScanner(
            "order-service", "my-tx-group");
    }
}

// 2. 业务代码中使用@GlobalTransactional
@Service
public class OrderService {
    
    @GlobalTransactional(timeoutMills = 60000, name = "create-order")
    public OrderResult createOrder(OrderRequest request) {
        // 步骤1:创建订单(本地事务)
        Order order = orderMapper.insert(request);
        
        // 步骤2:扣减库存(远程调用)
        inventoryClient.deduct(request.getProductId(), request.getQuantity());
        
        // 步骤3:扣减账户余额(远程调用)
        accountClient.debit(request.getUserId(), order.getTotalAmount());
        
        return OrderResult.success(order);
    }
}

AT模式的工作原理:

  1. 拦截SQL,在执行前保存before-image(修改前数据快照)
  2. 执行业务SQL
  3. 保存after-image(修改后数据快照)
  4. 如果全局事务回滚,用before-image反向补偿

AT模式的问题:before/after image存储在undo_log表中,对数据库有额外写入压力;全局锁在高并发场景下成为瓶颈。

TCC模式:业务补偿的柔性事务

TCC(Try-Confirm-Cancel)将事务拆分为三个阶段,由业务代码实现补偿逻辑,不依赖数据库undo log:

@Service
public class InventoryTccService {
    
    @TwoPhaseBusinessAction(
        name = "deductInventory",
        commitMethod = "confirm",
        rollbackMethod = "cancel"
    )
    public boolean tryDeduct(
        @BusinessActionContextParameter(paramName = "productId") String productId,
        @BusinessActionContextParameter(paramName = "quantity") int quantity
    ) {
        // Try阶段:冻结库存(不是直接扣减)
        int affected = inventoryMapper.freezeStock(productId, quantity);
        if (affected == 0) {
            throw new BusinessException("库存不足");
        }
        return true;
    }
    
    public boolean confirm(BusinessActionContext context) {
        // Confirm阶段:确认扣减冻结库存
        String productId = context.getActionContext("productId", String.class);
        Integer quantity = context.getActionContext("quantity", Integer.class);
        inventoryMapper.confirmDeduct(productId, quantity);
        return true;
    }
    
    public boolean cancel(BusinessActionContext context) {
        // Cancel阶段:释放冻结库存
        String productId = context.getActionContext("productId", String.class);
        Integer quantity = context.getActionContext("quantity", Integer.class);
        inventoryMapper.releaseFrozen(productId, quantity);
        return true;
    }
}

TCC vs AT选择依据

  • AT模式:适合标准CRUD操作,无业务侵入,但全局锁可能成为瓶颈
  • TCC模式:适合非标业务逻辑(冻结/预占场景),无全局锁性能更好,但开发量是3倍

Saga模式:长事务的最终一致性方案

Saga模式将长事务拆分为多个本地事务,每个本地事务完成后触发下一个,任一步失败则逆向执行补偿。适合执行时间长的业务流程(如订单履约、跨机构转账)。

基于事件编排的Saga实现

// 订单创建Saga流程
@Service
public class OrderSagaOrchestrator {
    
    @Autowired
    private KafkaTemplate kafka;
    
    public void startCreateOrderSaga(OrderRequest request) {
        String sagaId = UUID.randomUUID().toString();
        
        // 发起第一步:创建订单
        OrderCreatedEvent event = new OrderCreatedEvent(
            sagaId, request.getOrderId(), 
            request.getProductId(), request.getQuantity()
        );
        kafka.send("order-created", event);
    }
    
    // Step 2: 监听订单创建事件 → 扣减库存
    @KafkaListener(topics = "order-created")
    public void onOrderCreated(OrderCreatedEvent event) {
        try {
            inventoryService.deduct(event.getProductId(), event.getQuantity());
            kafka.send("inventory-deducted", 
                new InventoryDeductedEvent(event.getSagaId()));
        } catch (Exception e) {
            // 库存扣减失败,发起订单取消
            kafka.send("order-cancel-requested", 
                new OrderCancelEvent(event.getSagaId(), event.getOrderId()));
        }
    }
    
    // Step 3: 监听库存扣减事件 → 扣减余额
    @KafkaListener(topics = "inventory-deducted")
    public void onInventoryDeducted(InventoryDeductedEvent event) {
        try {
            accountService.debit(event.getUserId(), event.getAmount());
            kafka.send("account-debited", 
                new AccountDebitedEvent(event.getSagaId()));
        } catch (Exception e) {
            // 余额扣减失败,逆向补偿:恢复库存 + 取消订单
            inventoryService.restore(event.getProductId(), event.getQuantity());
            kafka.send("order-cancel-requested", 
                new OrderCancelEvent(event.getSagaId(), event.getOrderId()));
        }
    }
}

Saga补偿的幂等性保障

补偿操作可能被重复执行(网络重试、消费者rebalance),每个补偿逻辑必须保证幂等:

// 幂等补偿:用唯一键去重
@Transactional
public void restoreInventory(String productId, int quantity, String sagaId) {
    // 检查是否已补偿
    if (compensationLogMapper.existsBySagaId(sagaId)) {
        log.info("补偿已执行,跳过: sagaId={}", sagaId);
        return;
    }
    
    // 执行补偿
    inventoryMapper.restoreStock(productId, quantity);
    
    // 记录补偿日志
    compensationLogMapper.insert(new CompensationLog(sagaId, "INVENTORY_RESTORE"));
}

消息中间件保障最终一致性

很多场景不需要强一致性,只需要”写操作最终生效”。消息中间件是这类场景的最优解。

本地消息表模式

核心思路:业务操作和消息发送在同一个本地事务中完成,消息先写到本地表,由后台任务异步发送到MQ:

@Service
public class PaymentService {
    
    @Transactional
    public void processPayment(PaymentRequest request) {
        // 1. 执行业务操作
        paymentMapper.insert(new Payment(request));
        accountMapper.debit(request.getUserId(), request.getAmount());
        
        // 2. 在同一事务中写入消息表
        outboxMapper.insert(new OutboxMessage(
            "payment-completed",
            JSON.toJSONString(new PaymentCompletedEvent(request)),
            LocalDateTime.now()
        ));
    }
}

// 定时任务扫描消息表并发送到MQ
@Scheduled(fixedDelay = 1000)
public void sendPendingMessages() {
    List messages = outboxMapper.findPending(100);
    for (OutboxMessage msg : messages) {
        try {
            kafka.send(msg.getTopic(), msg.getPayload()).get(5, TimeUnit.SECONDS);
            outboxMapper.markSent(msg.getId());
        } catch (Exception e) {
            outboxMapper.incrementRetry(msg.getId());
            if (msg.getRetryCount() >= 5) {
                outboxMapper.markFailed(msg.getId());
                alertService.notify("消息发送失败: " + msg.getId());
            }
        }
    }
}

API接口规范与事务边界设计

微服务间的API设计直接影响事务复杂度。基本原则:一个API调用对应一个事务边界

// 反模式:一个API内部链式调用多个服务
@PostMapping("/order")
public Result createOrder(@RequestBody OrderRequest req) {
    orderService.create(req);          // 事务1
    inventoryService.deduct(req);      // 事务2(远程调用)
    accountService.debit(req);         // 事务3(远程调用)
    // 任何一个失败,前面的不会自动回滚
}

// 正模式:明确事务边界 + 异步解耦
@PostMapping("/order")
public Result createOrder(@RequestBody OrderRequest req) {
    // 只做本地事务:创建订单 + 写消息表
    orderService.createWithEvent(req);
    // 其余步骤由Saga异步完成
    return Result.accepted("订单已受理,处理中");
}

服务治理:事务监控与故障恢复

分布式事务的运维核心是监控和告警:

// Prometheus监控Seata事务指标
// - seata_transaction_total:事务总数
// - seata_transaction_active:活跃事务数
// - seata_transaction_failed_total:失败事务数

// 关键告警规则
- alert: DistributedTxHighFailureRate
  expr: rate(seata_transaction_failed_total[5m]) / rate(seata_transaction_total[5m]) > 0.05
  labels:
    severity: warning
  annotations:
    summary: "分布式事务失败率超过5%"

// 悬挂事务处理:超时未完成的事务需要人工介入
- alert: SuspiciousLongTx
  expr: seata_transaction_active > 0 and seata_transaction_max_duration_seconds > 120
  labels:
    severity: critical

分布式事务没有银弹。场景决定方案:短事务强一致性用2PC/TCC,长流程最终一致性用Saga,跨服务通知用消息表。关键是识别业务场景的事务边界,然后用最轻量的方案满足一致性要求。过度使用强一致性事务,是微服务架构性能劣化的首要原因。

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

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

相关推荐