微服务架构将单体应用拆分为独立部署的服务单元,每个服务拥有自己的数据库实例。这种拆分带来了数据一致性的挑战——跨服务的业务操作无法依赖本地数据库事务保证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/