Spring Boot微服务架构下的分布式事务补偿方案:Saga模式与消息中间件联动实战

微服务架构分布式事务的现实困境

微服务架构拆分带来的核心难题是分布式事务。单体应用中一个本地事务就能保证的ACID特性,在微服务间变成跨网络调用的最终一致性问题。Spring Boot生态中,两阶段提交(2PC/XA)因为同步阻塞和网络延迟,在微服务场景下几乎不可用。真正生产级采用的方案是Saga模式——将长事务拆成多个本地事务,每个本地事务提交后通过事件触发下一个步骤,任一步骤失败时执行补偿操作回滚已完成的步骤。

Saga编排模式 vs 协调模式:选型决策

Saga有两种实现模式:

编排模式(Choreography):各服务通过事件总线交互,没有中心协调器。服务A完成本地事务后发布事件,服务B监听事件后执行自己的事务并发布新事件。优点是松耦合,缺点是事务流程分散在各服务代码中,难以全局追踪。

协调模式(Orchestration):中心协调器控制整个Saga流程,向各服务发送命令并收集结果。优点是流程清晰、补偿逻辑集中管理,缺点是协调器本身成为单点。

事务步骤超过3个、涉及5个以上服务时,协调模式的可维护性优势明显。以下实战采用协调模式。

Saga协调器核心实现

Saga协调器的职责是:按序执行各步骤,记录每步状态,失败时按逆序执行补偿操作。核心数据结构:

// SagaDefinition.java
public class SagaDefinition {
    private String sagaId;
    private List<SagaStep> steps;
    
    public static class SagaStep {
        private String stepName;
        private String serviceEndpoint;
        private String compensationEndpoint;
        private StepStatus status;
        private String transactionId;
    }
    
    public enum StepStatus {
        PENDING, EXECUTING, COMPLETED, COMPENSATING, COMPENSATED, FAILED
    }
}

// SagaInstance.java
@Entity
public class SagaInstance {
    @Id
    private String sagaId;
    private String sagaType;
    private Integer currentStep;
    private SagaStatus status;
    private String payload;
    private LocalDateTime createdAt;
    private LocalDateTime updatedAt;
    
    public enum SagaStatus {
        RUNNING, COMPLETED, COMPENSATING, COMPENSATED, FAILED
    }
}

Saga协调器的主循环逻辑:

@Service
public class SagaOrchestrator {
    
    @Autowired
    private SagaInstanceRepository sagaRepo;
    
    @Autowired
    private RabbitTemplate rabbitTemplate;
    
    public void executeSaga(SagaDefinition definition, String payload) {
        SagaInstance instance = new SagaInstance();
        instance.setSagaId(UUID.randomUUID().toString());
        instance.setSagaType(definition.getSagaId());
        instance.setCurrentStep(0);
        instance.setStatus(SagaStatus.RUNNING);
        instance.setPayload(payload);
        instance.setCreatedAt(LocalDateTime.now());
        sagaRepo.save(instance);
        executeNextStep(instance, definition);
    }
    
    private void executeNextStep(SagaInstance instance, 
            SagaDefinition definition) {
        List<SagaStep> steps = definition.getSteps();
        int current = instance.getCurrentStep();
        
        if (current >= steps.size()) {
            instance.setStatus(SagaStatus.COMPLETED);
            sagaRepo.save(instance);
            return;
        }
        
        SagaStep step = steps.get(current);
        step.setStatus(StepStatus.EXECUTING);
        
        SagaCommand command = new SagaCommand();
        command.setSagaId(instance.getSagaId());
        command.setStepIndex(current);
        command.setEndpoint(step.getServiceEndpoint());
        command.setPayload(instance.getPayload());
        
        rabbitTemplate.convertAndSend(
            "saga.commands", step.getStepName(),
            JsonUtils.toJson(command));
    }
    
    @RabbitListener(queues = "saga.responses")
    public void handleStepResponse(String message) {
        SagaResponse response = JsonUtils.fromJson(message, SagaResponse.class);
        SagaInstance instance = sagaRepo
            .findById(response.getSagaId()).orElseThrow();
        SagaDefinition definition = loadDefinition(instance.getSagaType());
        
        if (response.isSuccess()) {
            instance.setCurrentStep(instance.getCurrentStep() + 1);
            definition.getSteps().get(response.getStepIndex())
                .setStatus(StepStatus.COMPLETED);
            definition.getSteps().get(response.getStepIndex())
                .setTransactionId(response.getTransactionId());
            sagaRepo.save(instance);
            executeNextStep(instance, definition);
        } else {
            startCompensation(instance, definition, response.getStepIndex());
        }
    }
}

补偿操作:逆序回滚的工程实现

当某个步骤失败时,需要对之前所有已完成的步骤按逆序执行补偿操作。补偿的关键约束是幂等性——同一个补偿操作可能被执行多次(网络超时重试场景)。

private void startCompensation(SagaInstance instance, 
        SagaDefinition definition, int failedStep) {
    instance.setStatus(SagaStatus.COMPENSATING);
    sagaRepo.save(instance);
    
    for (int i = failedStep - 1; i >= 0; i--) {
        SagaStep step = definition.getSteps().get(i);
        if (step.getStatus() == StepStatus.COMPLETED) {
            SagaCommand compensateCmd = new SagaCommand();
            compensateCmd.setSagaId(instance.getSagaId());
            compensateCmd.setStepIndex(i);
            compensateCmd.setEndpoint(step.getCompensationEndpoint());
            compensateCmd.setTransactionId(step.getTransactionId());
            compensateCmd.setPayload(instance.getPayload());
            
            rabbitTemplate.convertAndSend(
                "saga.compensations", step.getStepName(),
                JsonUtils.toJson(compensateCmd));
        }
    }
}

补偿端点的幂等实现:每个服务维护一张补偿记录表,执行补偿前先查表判断是否已执行:

@Entity
@Table(uniqueConstraints = @UniqueConstraint(
    columnNames = {"sagaId", "stepIndex"}))
public class CompensationRecord {
    @Id @GeneratedValue
    private Long id;
    private String sagaId;
    private Integer stepIndex;
    private String transactionId;
    private CompensationStatus status;
    private LocalDateTime compensatedAt;
}

@Service
public class OrderCompensationHandler {
    
    @Autowired
    private CompensationRecordRepository compRepo;
    
    @RabbitListener(queues = "saga.compensation.order")
    public void handleCompensation(String message) {
        SagaCommand cmd = JsonUtils.fromJson(message, SagaCommand.class);
        
        Optional<CompensationRecord> existing = 
            compRepo.findBySagaIdAndStepIndex(
                cmd.getSagaId(), cmd.getStepIndex());
        if (existing.isPresent() 
                && existing.get().getStatus() == CompensationStatus.DONE) {
            return; // 已补偿,直接返回
        }
        
        orderService.cancelOrder(cmd.getTransactionId());
        
        CompensationRecord record = new CompensationRecord();
        record.setSagaId(cmd.getSagaId());
        record.setStepIndex(cmd.getStepIndex());
        record.setTransactionId(cmd.getTransactionId());
        record.setStatus(CompensationStatus.DONE);
        record.setCompensatedAt(LocalDateTime.now());
        compRepo.save(record);
    }
}

消息中间件配置:RabbitMQ的可靠投递保障

Saga模式对消息中间件的核心要求是消息不丢失。RabbitMQ配置要点:

@Configuration
public class RabbitMQConfig {
    
    @Bean
    public ConnectionFactory connectionFactory() {
        CachingConnectionFactory cf = new CachingConnectionFactory();
        cf.setHost("rabbitmq-service");
        cf.setPort(5672);
        cf.setUsername("saga_user");
        cf.setPassword("saga_pass");
        cf.setPublisherConfirmType(
            CachingConnectionFactory.ConfirmType.CORRELATED);
        cf.setPublisherReturns(true);
        return cf;
    }
    
    @Bean
    public RabbitTemplate rabbitTemplate(ConnectionFactory cf) {
        RabbitTemplate template = new RabbitTemplate(cf);
        template.setMessageConverter(new Jackson2JsonMessageConverter());
        template.setMandatory(true);
        
        template.setConfirmCallback((correlationData, ack, cause) -> {
            if (!ack) {
                log.error("消息投递失败: {}", cause);
            }
        });
        
        template.setReturnsCallback(returned -> {
            log.error("消息无法路由: {} - {}", 
                returned.getMessage(), returned.getReplyText());
        });
        
        return template;
    }
    
    @Bean
    public Queue sagaResponseQueue() {
        return QueueBuilder.durable("saga.responses")
            .withArgument("x-dead-letter-exchange", "saga.dlx")
            .withArgument("x-dead-letter-routing-key", "dead")
            .build();
    }
}

死信队列(DLX)是关键配置。当消息消费失败超过重试次数后,消息进入死信队列,由人工介入处理。这避免了消费失败的消息被静默丢弃。

服务治理:Saga的超时监控与健康检查

Saga协调器需要主动监控长时间未完成的Saga实例,防止因为某个服务无响应导致整个事务悬挂:

@Scheduled(fixedRate = 30000)
public void monitorStuckSagas() {
    LocalDateTime threshold = LocalDateTime.now().minusMinutes(10);
    List<SagaInstance> stuckSagas = sagaRepo
        .findByStatusAndUpdatedAtBefore(SagaStatus.RUNNING, threshold);
    
    for (SagaInstance saga : stuckSagas) {
        log.warn("Saga超时未完成: sagaId={}, 当前步骤={}",
            saga.getSagaId(), saga.getCurrentStep());
        
        long minutes = ChronoUnit.MINUTES.between(
            saga.getUpdatedAt(), LocalDateTime.now());
        
        if (minutes > 30) {
            saga.setStatus(SagaStatus.COMPENSATING);
            sagaRepo.save(saga);
            startCompensation(saga, 
                loadDefinition(saga.getSagaType()), 
                saga.getCurrentStep());
        }
    }
}

Spring Boot微服务架构下,分布式事务没有银弹。Saga模式通过牺牲强一致性换取可用性和性能,补偿操作的设计质量直接决定数据一致性保障的可靠性。核心原则:每个本地事务必须可逆,补偿操作必须幂等,消息投递必须可靠。这三个约束不满足任何一个,Saga方案都无法上生产。

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

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

相关推荐