微服务架构中的分布式事务方案:Saga模式与TCC实现详解

微服务架构将单体应用拆分为独立部署的服务单元,每个服务拥有自己的数据库实例。这种拆分带来了数据一致性的挑战——跨服务的业务操作无法依赖本地数据库事务保证ACID特性。分布式事务方案成为微服务架构绕不开的核心问题。本文对比Saga和TCC两种主流方案,结合Spring Boot框架说明在Java/Go实战中的具体实现路径。

分布式事务的核心挑战与CAP约束

分布式事务的本质问题在于网络分区、节点故障、超时重试等不确定性因素。CAP定理指出,在网络分区发生时,一致性和可用性不可兼得。实际工程中通常选择BASE模型(Basically Available, Soft state, Eventually consistent)而非强一致性,通过补偿机制和幂等设计实现最终一致性。

主流分布式事务方案对比:2PC(两阶段提交)保证强一致性但阻塞资源、可用性差;Saga通过正反向操作序列实现最终一致性,适合长流程业务;TCC通过Try-Confirm-Cancel三阶段实现准实时一致性,适合对一致性要求较高的资金类业务;本地消息表+消息中间件实现异步最终一致性,适合对延迟容忍度高的场景。

Saga模式:编排式与编舞式的实现选择

Saga模式将一个跨服务事务拆分为一系列本地事务,每个本地事务有对应的补偿操作。如果某一步失败,按逆序执行已完成步骤的补偿操作。Saga有两种协调方式:编排式(Orchestration)由中央协调器统一调度;编舞式(Choreography)通过事件驱动各服务自行响应。

编排式Saga使用SEATA框架实现的代码示例:

// Saga协调器(Spring Boot)
@Service
public class OrderSagaOrchestrator {

    @Autowired
    private OrderService orderService;
    @Autowired
    private InventoryService inventoryService;
    @Autowired
    private PaymentService paymentService;

    public OrderResult executeOrder(OrderRequest request) {
        SagaOrchestrator saga = SagaOrchestrator.create();
        
        // 定义Saga流程:创建订单 -> 扣库存 -> 支付
        saga
            .process(
                // 正向操作
                () -> orderService.create(request),
                // 补偿操作
                order -> orderService.cancel(order.getId())
            )
            .process(
                () -> inventoryService.deduct(request.getProductId(), request.getQuantity()),
                result -> inventoryService compensable(result.getdeductId())
            )
            .process(
                () -> paymentService.charge(request.getPaymentInfo()),
                payment -> paymentService.refund(payment.getPaymentId())
            );
        
        try {
            saga.execute();
            return OrderResult.success();
        } catch (SagaExecutionException e) {
            // Saga框架自动执行补偿,此处处理最终失败
            log.error("Saga执行失败,已执行补偿", e);
            return OrderResult.failure("订单处理失败");
        }
    }
}

// 各服务实现TCC接口
@LocalTCC
public interface InventoryTccAction {
    
    @TwoPhaseBusinessAction(name = "inventoryTccAction", 
        commitMethod = "commit", rollbackMethod = "rollback")
    boolean deduct(BusinessActionContext ctx,
                   @BusinessActionContextParameter(paramName = "productId") String productId,
                   @BusinessActionContextParameter(paramName = "quantity") int quantity);
    
    boolean commit(BusinessActionContext ctx);
    boolean rollback(BusinessActionContext ctx);
}

编舞式Saga通过消息中间件解耦,各服务监听事件并自行决策:

// 订单服务发布事件
@Service
public class OrderService {
    
    @Autowired
    private KafkaTemplate<String, String> kafka;
    
    @Transactional
    public Order create(OrderRequest request) {
        Order order = orderRepo.save(request.toEntity());
        // 发布订单创建事件
        kafka.send("order-events", new OrderCreatedEvent(order.getId()).toJson());
        return order;
    }
    
    // 监听库存扣减结果
    @KafkaListener(topics = "inventory-events")
    public void onInventoryEvent(InventoryEvent event) {
        if (event.getType() == InventoryEvent.Type.DEDUCTED) {
            // 库存扣减成功,继续支付流程
            kafka.send("order-events", new OrderConfirmedEvent(event.getOrderId()).toJson());
        } else if (event.getType() == InventoryEvent.Type.DEDUCT_FAILED) {
            // 库存扣减失败,取消订单
            cancelOrder(event.getOrderId());
        }
    }
}

编舞式的优势在于服务间松耦合,新增步骤只需新增事件监听器。劣势是流程可视化困难,调试复杂。编排式的流程清晰,但协调器成为单点。实践中长流程业务(如订单全生命周期)倾向使用编排式,短链路实时性要求高的场景使用编舞式。

TCC模式:Try-Confirm-Cancel三阶段实现

TCC模式将每个操作拆分为三个阶段:Try阶段预留资源(如冻结金额、预占库存),Confirm阶段确认提交,Cancel阶段释放预留资源。TCC的关键约束是三个阶段都必须幂等,因为网络超时可能导致重试。

// 账户服务TCC实现
@Service
public class AccountTccService {

    // Try:冻结金额
    @Transactional
    public TccResult tryFreeze(String accountId, BigDecimal amount, String xid) {
        // 幂等检查:同一个xid只执行一次
        if (tccLogRepository.existsByXidAndPhase(xid, "TRY")) {
            return TccResult.alreadyProcessed();
        }
        
        Account account = accountRepository.findById(accountId)
            .orElseThrow(() -> new BusinessException("账户不存在"));
        
        if (account.getAvailableBalance().compareTo(amount) < 0) {
            return TccResult.failure("余额不足");
        }
        
        // 冻结金额:从可用余额扣除,加入冻结余额
        account.setAvailableBalance(account.getAvailableBalance().subtract(amount));
        account.setFrozenBalance(account.getFrozenBalance().add(amount));
        accountRepository.save(account);
        
        // 记录TCC日志
        tccLogRepository.save(new TccLog(xid, accountId, "TRY", amount));
        return TccResult.success();
    }
    
    // Confirm:扣减冻结金额
    @Transactional
    public TccResult confirmDeduct(String xid) {
        if (tccLogRepository.existsByXidAndPhase(xid, "CONFIRM")) {
            return TccResult.alreadyProcessed();
        }
        
        TccLog tryLog = tccLogRepository.findByXidAndPhase(xid, "TRY");
        Account account = accountRepository.findById(tryLog.getAccountId()).get();
        
        // 从冻结余额中扣除
        account.setFrozenBalance(account.getFrozenBalance().subtract(tryLog.getAmount()));
        accountRepository.save(account);
        
        tccLogRepository.save(new TccLog(xid, tryLog.getAccountId(), "CONFIRM", tryLog.getAmount()));
        return TccResult.success();
    }
    
    // Cancel:解冻金额
    @Transactional
    public TccResult cancelFreeze(String xid) {
        if (tccLogRepository.existsByXidAndPhase(xid, "CANCEL")) {
            return TccResult.alreadyProcessed();
        }
        
        TccLog tryLog = tccLogRepository.findByXidAndPhase(xid, "TRY");
        if (tryLog == null) {
            // Try未执行,空回滚
            tccLogRepository.save(new TccLog(xid, null, "CANCEL", BigDecimal.ZERO));
            return TccResult.success();
        }
        
        Account account = accountRepository.findById(tryLog.getAccountId()).get();
        // 将冻结金额退回可用余额
        account.setFrozenBalance(account.getFrozenBalance().subtract(tryLog.getAmount()));
        account.setAvailableBalance(account.getAvailableBalance().add(tryLog.getAmount()));
        accountRepository.save(account);
        
        tccLogRepository.save(new TccLog(xid, tryLog.getAccountId(), "CANCEL", tryLog.getAmount()));
        return TccResult.success();
    }
}

TCC实现中最容易出错的三个边界情况:空回滚(Try未执行但收到Cancel请求,需检测Try日志是否存在)、悬挂(Cancel先于Try到达,需在Try前检查Cancel是否已执行)、幂等控制(同一阶段重复调用需返回相同结果)。上述代码通过tcc_log表记录每个阶段的执行状态来处理这三种情况。

Go语言中的Saga实现:DTM框架实战

Go语言生态中,DTM(Distributed Transaction Manager)是轻量级分布式事务管理器,支持Saga、TCC、XA、二阶段消息等多种模式。以下是在Go微服务中使用DTM实现Saga的示例:

package main

import (
    "github.com/dtm-labs/dtmcli"
    "github.com/gin-gonic/gin"
)

const dtmServer = "http://dtm:36789/api/dtmsvr"

func main() {
    r := gin.Default()
    
    // 订单创建接口
    r.POST("/api/order/create", func(c *gin.Context) {
        req := &OrderRequest{}
        c.Bind(req)
        
        // 构建Saga
        saga := dtmcli.NewSaga(dtmServer, "order-"+req.OrderID).
            Add(
                // 正向操作:创建订单
                "http://order-svc/api/order/create",
                // 补偿操作:取消订单
                "http://order-svc/api/order/cancel",
                req,
            ).
            Add(
                // 正向操作:扣减库存
                "http://inventory-svc/api/inventory/deduct",
                // 补偿操作:恢复库存
                "http://inventory-svc/api/inventory/restore",
                &DeductReq{ProductID: req.ProductID, Quantity: req.Qty},
            ).
            Add(
                // 正向操作:支付
                "http://payment-svc/api/payment/charge",
                // 补偿操作:退款
                "http://payment-svc/api/payment/refund",
                &PaymentReq{OrderID: req.OrderID, Amount: req.Amount},
            )
        
        // 设置超时和重试策略
        saga.TimeoutToFail = 1800  // 30分钟超时
        saga.RetryInterval = 10    // 重试间隔10秒
        saga.PanicOptional = true
        
        // 提交Saga到DTM服务器
        err := saga.Submit()
        if err != nil {
            c.JSON(500, gin.H{"error": err.Error()})
            return
        }
        c.JSON(200, gin.H{"msg": "order processing"})
    })
    
    r.Run(":8080")
}

DTM将事务协调逻辑从业务服务中剥离到独立的事务管理器,业务服务只需实现正向操作和补偿操作的HTTP接口。DTM自动处理超时重试、补偿调度和状态管理。对于需要服务治理和限流熔断保护的高并发设计场景,DTM与Sentinel集成可在事务重试时自动触发熔断,避免故障扩散。

消息中间件在最终一致性中的角色

本地消息表方案是另一种常用的最终一致性实现。核心思路是业务操作和消息写入在同一个本地事务中完成,后台任务轮询消息表并将消息投递到消息中间件。消费端通过幂等消费保证不重复处理:

// 本地消息表方案
@Service
public class OrderService {
    
    @Transactional
    public void createOrder(OrderRequest request) {
        // 1. 写订单表
        Order order = orderRepo.save(request.toEntity());
        
        // 2. 写本地消息表(同一事务)
        Message message = new Message();
        message.setTopic("order-created");
        message.setBody(toJson(new OrderCreatedEvent(order)));
        message.setStatus("PENDING");
        messageRepo.save(message);
    }
}

// 后台任务轮询投递
@Scheduled(fixedRate = 1000)
public void publishMessages() {
    List<Message> pending = messageRepo.findByStatus("PENDING", PageRequest.of(0, 100));
    for (Message msg : pending) {
        try {
            kafkaTemplate.send(msg.getTopic(), msg.getBody()).get();
            msg.setStatus("SENT");
            messageRepo.save(msg);
        } catch (Exception e) {
            // 投递失败,下次重试
            msg.incrementRetryCount();
            if (msg.getRetryCount() > MAX_RETRY) {
                msg.setStatus("FAILED");
                alertService.notify("消息投递失败: " + msg.getId());
            }
            messageRepo.save(msg);
        }
    }
}

消息中间件选型取决于业务场景的吞吐和可靠性要求。RocketMQ原生支持事务消息,无需本地消息表即可实现类似效果。Kafka需要配合本地消息表或Kafka事务API使用。RabbitMQ通过 publisher confirm 机制保障消息不丢失,但在高吞吐场景下性能不如Kafka和RocketMQ。API接口规范方面,消息体建议遵循CloudEvents规范,便于跨语言解析和链路追踪。

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

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

相关推荐