分布式事务一致性保障:Saga模式在微服务架构中的落地方案
微服务架构下,一个业务操作往往跨越多个服务,每个服务维护自己的数据库,本地事务无法保证跨服务一致性。Saga模式是目前最主流的分布式事务解决方案,但它的补偿机制设计比想象中复杂。本文从模式选型、补偿设计到代码实现,给出可落地的方案。
为什么不用两阶段提交(2PC)
2PC的核心问题是锁等待——参与者在第一阶段锁定资源后,必须等协调者通知提交或回滚。如果协调者故障或网络分区,参与者会一直持有锁,阻塞其他事务。在微服务场景下,这个锁可能跨越多个服务,影响范围巨大。
Saga的思路完全不同:不锁定资源,每个步骤都是本地事务,失败时通过补偿操作回滚已完成的步骤。牺牲了隔离性,换来了可用性和性能。
编排式Saga vs 协调式Saga
两种实现方式各有适用场景。
编排式(Choreography):各服务通过事件驱动,监听其他服务发出的事件决定下一步动作。适合流程简单、步骤少于5个的场景。
// 订单服务发布事件
@Transactional
public Order createOrder(OrderDTO dto) {
Order order = orderRepository.save(new Order(dto));
eventPublisher.publish(new OrderCreatedEvent(order.getId()));
return order;
}
// 库存服务监听事件
@EventListener
public void onOrderCreated(OrderCreatedEvent event) {
inventoryService.deduct(event.getProductId(), event.getQuantity());
eventPublisher.publish(new InventoryDeductedEvent(event.getOrderId()));
}
// 支付服务监听库存扣减事件
@EventListener
public void onInventoryDeducted(InventoryDeductedEvent event) {
paymentService.charge(event.getOrderId());
}
编排式的缺点:流程分散在各服务中,没有全局视角,难以追踪完整链路。步骤多了以后,事件的监听关系变成网状,极难维护。
协调式(Orchestration):由一个Saga协调器统一控制流程,每一步完成后由协调器决定下一步。适合步骤多、有条件分支和并行的复杂流程。
// Saga协调器定义
public class OrderSaga {
public SagaDefinition<OrderState> sagaDefinition() {
return step()
.invokeParticipant(this::createOrder)
.withCompensation(this::cancelOrder)
.step()
.invokeParticipant(this::deductInventory)
.withCompensation(this::restoreInventory)
.step()
.invokeParticipant(this::chargePayment)
.withCompensation(this::refundPayment)
.step()
.invokeParticipant(this::confirmOrder)
.build();
}
private void cancelOrder(OrderState state) {
orderService.updateStatus(state.getOrderId(), "CANCELLED");
}
private void restoreInventory(OrderState state) {
inventoryService.restore(state.getProductId(), state.getQuantity());
}
private void refundPayment(OrderState state) {
paymentService.refund(state.getPaymentId());
}
}
协调式的优点是流程集中可追踪,缺点是协调器可能成为单点。生产环境中协调器需要做持久化和重试,保证Saga不会因为协调器宕机而卡在中间态。
补偿操作的幂等性设计
Saga补偿最容易被忽略的问题是幂等性。补偿操作可能被多次执行(网络超时重试、协调器重启后重复发送),必须保证多次执行的结果一致。
设计模式:在业务表中增加补偿状态字段,补偿前先检查状态:
@Transactional
public void restoreInventory(Long orderId, String productId, int quantity) {
// 检查是否已补偿
CompensateRecord record = compensateRepo.findByOrderId(orderId);
if (record != null && record.getStatus() == COMPENSATED) {
log.info("补偿已执行,跳过: orderId={}", orderId);
return; // 幂等返回
}
// 执行补偿
inventoryRepo.addStock(productId, quantity);
// 记录补偿状态
if (record == null) {
compensateRepo.save(new CompensateRecord(orderId, COMPENSATED));
} else {
record.setStatus(COMPENSATED);
}
}
消息中间件在Saga中的角色
微服务之间的Saga通信必须使用消息中间件(如Kafka、RabbitMQ),而非HTTP同步调用。原因:HTTP调用失败时无法区分”对方没收到”和”对方处理完但响应丢了”,消息中间件通过确认机制保证至少投递一次。
// 使用Kafka发送Saga命令
@Component
public class SagaCommandPublisher {
@Autowired
private KafkaTemplate<String, SagaCommand> kafkaTemplate;
public void send(String topic, SagaCommand command) {
kafkaTemplate.send(topic, command.getSagaId(), command)
.addCallback(
result -> log.info("Saga命令发送成功: sagaId={}, step={}",
command.getSagaId(), command.getStep()),
ex -> {
log.error("Saga命令发送失败,将重试: sagaId={}", command.getSagaId(), ex);
// 重试逻辑:延迟队列或人工介入
}
);
}
}
// 消费端
@KafkaListener(topics = "saga-inventory-command")
public void handleInventoryCommand(SagaCommand command) {
try {
switch (command.getAction()) {
case DEDUCT -> inventoryService.deduct(command);
case RESTORE -> inventoryService.restore(command);
}
sagaClient.reportSuccess(command.getSagaId(), command.getStep());
} catch (Exception e) {
sagaClient.reportFailure(command.getSagaId(), command.getStep(), e.getMessage());
}
}
Saga状态持久化与故障恢复
协调器在每一步执行前后必须持久化Saga状态,否则宕机后无法知道哪些步骤已完成、需要补偿哪些步骤。
-- Saga状态表
CREATE TABLE saga_instance (
saga_id VARCHAR(64) PRIMARY KEY,
saga_type VARCHAR(128) NOT NULL,
current_step INT NOT NULL,
status ENUM('RUNNING', 'COMPENSATING', 'COMPLETED', 'FAILED'),
payload JSON,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
INDEX idx_status (status),
INDEX idx_updated (updated_at)
);
-- 步骤执行记录
CREATE TABLE saga_step (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
saga_id VARCHAR(64) NOT NULL,
step INT NOT NULL,
action VARCHAR(64) NOT NULL,
status ENUM('PENDING', 'EXECUTING', 'COMPLETED', 'COMPENSATED', 'FAILED'),
request_payload JSON,
response_payload JSON,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
UNIQUE KEY uk_saga_step (saga_id, step)
);
故障恢复机制:定时任务扫描status=RUNNING且updated_at超过阈值的Saga实例,根据最后完成的步骤决定继续执行或触发补偿。
服务治理中的Saga监控
Saga的运行状态必须可视化,否则出了问题只能查日志。核心监控指标:
1. 运行中的Saga数量和平均时长——如果持续增长说明有步骤阻塞。
2. 补偿触发率——补偿频繁说明下游服务不稳定。
3. 各步骤的成功率和耗时——定位慢步骤。
4. 卡住的Saga——长时间没有状态更新的实例,需要人工介入。
// Micrometer指标埋点
@Autowired
private MeterRegistry registry;
public void executeStep(SagaCommand command) {
Timer.Sample sample = Timer.start(registry);
try {
// 执行步骤逻辑
stepExecutor.execute(command);
registry.counter("saga.step.success", "step", command.getStep()).increment();
} catch (Exception e) {
registry.counter("saga.step.failure", "step", command.getStep(),
"error", e.getClass().getSimpleName()).increment();
throw e;
} finally {
sample.stop(registry.timer("saga.step.duration", "step", command.getStep()));
}
}
微服务架构下,分布式事务的一致性保障不是选择问题,而是必须面对的问题。Saga提供了在可用性和一致性之间取得平衡的方案,但前提是补偿逻辑设计正确、幂等性有保证、状态可追踪可恢复。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/fen-bu-shi-shi-wu-yi-zhi-xing-bao-zhang-saga-mo-shi-zai-wei/