Spring Boot分布式事务实战方案:Seata AT模式与消息最终一致性对比选型

分布式事务的场景与痛点

微服务架构下单笔业务操作跨多个服务时,本地事务无法保证数据一致性。用户下单扣库存、转账涉及两个账户、订单创建同时发优惠券——每个场景都是跨库操作,部分成功部分失败就会产生脏数据。

分布式事务方案没有银弹。强一致性代价高,最终一致性有时延。选型的核心是根据业务场景的容忍度做取舍,而不是追求一个方案通吃所有场景。

Seata AT模式实战

Seata AT模式是侵入性最低的分布式事务方案,只需要加注解和配置数据源代理。适合对一致性要求高、能接受2-5ms额外延迟的场景。

服务端部署

# docker-compose-seata.yml
version: "3.8"
services:
  seata-server:
    image: seataio/seata-server:2.0.0
    ports:
      - "7091:7091"
      - "8091:8091"
    environment:
      - SEATA_IP=seata-server
      - STORE_MODE=db
    volumes:
      - ./seata-config:/seata-server/config

application.yml配置

seata:
  enabled: true
  application-id: order-service
  tx-service-group: order-tx-group
  service:
    vgroup-mapping:
      order-tx-group: default
    grouplist:
      default: seata-server:8091
  registry:
    type: nacos
    nacos:
      server-addr: nacos:8848
      namespace: seata
      group: SEATA_GROUP

数据源代理配置(必须手动配置,不能依赖自动代理):

@Configuration
public class SeataDataSourceConfig {

    @Bean
    @ConfigurationProperties(prefix = "spring.datasource")
    public DataSource dataSource() {
        return new HikariDataSource();
    }

    @Bean("seataDataSourceProxy")
    public DataSource seataDataSourceProxy(DataSource dataSource) {
        return new DataSourceProxy(dataSource);
    }

    @Bean
    @Primary
    public SqlSessionFactory sqlSessionFactory(DataSource seataDataSourceProxy) throws Exception {
        SqlSessionFactoryBean bean = new SqlSessionFactoryBean();
        bean.setDataSource(seataDataSourceProxy);
        bean.setMapperLocations(new PathMatchingResourcePatternResolver()
            .getResources("classpath:mapper/*.xml"));
        return bean.getObject();
    }
}

业务代码使用

@Service
public class OrderService {

    @Autowired
    private OrderMapper orderMapper;
    
    @Autowired
    private InventoryClient inventoryClient;
    
    @Autowired
    private AccountClient accountClient;

    @GlobalTransactional(name = "create-order", rollbackFor = Exception.class)
    public OrderDTO createOrder(OrderRequest request) {
        // 1. 创建订单(本地事务分支)
        Order order = new Order();
        order.setUserId(request.getUserId());
        order.setProductId(request.getProductId());
        order.setQuantity(request.getQuantity());
        order.setTotalAmount(request.getAmount());
        order.setStatus(OrderStatus.INIT);
        orderMapper.insert(order);

        // 2. 扣减库存(远程事务分支)
        inventoryClient.deduct(request.getProductId(), request.getQuantity());

        // 3. 扣减账户余额(远程事务分支)
        accountClient.debit(request.getUserId(), request.getAmount());

        // 4. 更新订单状态
        order.setStatus(OrderStatus.SUCCESS);
        orderMapper.updateById(order);

        return OrderConverter.toDTO(order);
    }
}

Seata AT模式的回滚机制依赖undo_log表,每个参与事务的数据库都需要建这张表。事务提交前Seata会在undo_log中记录数据的前后镜像,回滚时根据镜像恢复数据。

消息最终一致性方案

Seata AT适合强一致场景,但对于高并发业务(秒杀、抢购),两阶段提交的锁竞争会严重拖低吞吐量。消息最终一致性牺牲即时一致性,换取高吞吐和低延迟。

基于RocketMQ事务消息的实现

@Service
public class OrderServiceAsync {

    @Autowired
    private OrderMapper orderMapper;
    
    @Autowired
    private RocketMQTemplate rocketMQTemplate;

    public void createOrderAsync(OrderRequest request) {
        rocketMQTemplate.sendMessageInTransaction(
            "order-tx-topic",
            MessageBuilder.withPayload(request)
                .setHeader("userId", request.getUserId())
                .build(),
            request
        );
    }

    @RocketMQTransactionListener
    class OrderTransactionListener implements RocketMQLocalTransactionListener {

        @Autowired
        private OrderMapper orderMapper;

        @Override
        public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
            try {
                OrderRequest request = (OrderRequest) arg;
                Order order = new Order();
                order.setUserId(request.getUserId());
                order.setProductId(request.getProductId());
                order.setQuantity(request.getQuantity());
                order.setStatus(OrderStatus.PENDING);
                orderMapper.insert(order);
                return RocketMQLocalTransactionState.COMMIT;
            } catch (Exception e) {
                return RocketMQLocalTransactionState.ROLLBACK;
            }
        }

        @Override
        public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
            String orderId = (String) msg.getHeaders().get("orderId");
            Order order = orderMapper.selectById(orderId);
            if (order != null && order.getStatus() != OrderStatus.CANCELLED) {
                return RocketMQLocalTransactionState.COMMIT;
            }
            return RocketMQLocalTransactionState.ROLLBACK;
        }
    }
}

两种方案对比选型

| 维度 | Seata AT | 消息最终一致性 |
|——|———-|—————|
| 一致性模型 | 强一致(可串行化) | 最终一致(秒级延迟) |
| 性能开销 | 2-5ms/事务 + 全局锁 | 几乎无额外开销 |
| 吞吐量 | 中等(全局锁竞争) | 高(无锁) |
| 代码侵入性 | 低(注解+数据源代理) | 中(需实现事务消息) |
| 运维复杂度 | 中(需部署Seata Server) | 低(依赖MQ基础设施) |
| 典型场景 | 转账、库存+订单强一致 | 秒杀、通知、日志 |
| 故障恢复 | 自动回滚 | 需补偿机制 |

补偿机制的实现

消息最终一致性必须配套补偿机制,处理消费失败的情况:

@Service
public class CompensationService {

    @Autowired
    private CompensationLogMapper compensationMapper;
    
    @Autowired
    private RocketMQTemplate rocketMQTemplate;

    @Scheduled(fixedDelay = 60000)
    public void compensateFailedMessages() {
        List<CompensationLog> failedList = compensationMapper.selectList(
            new LambdaQueryWrapper<CompensationLog>()
                .eq(CompensationLog::getStatus, CompensationStatus.FAILED)
                .lt(CompensationLog::getRetryCount, 3)
                .lt(CompensationLog::getCreatedAt, 
                    LocalDateTime.now().minusMinutes(5))
        );

        for (CompensationLog log : failedList) {
            try {
                rocketMQTemplate.convertAndSend(
                    log.getTopic(), 
                    JSON.parseObject(log.getMessageBody(), OrderRequest.class)
                );
                log.setStatus(CompensationStatus.RETRYING);
                log.setRetryCount(log.getRetryCount() + 1);
            } catch (Exception e) {
                log.setRetryCount(log.getRetryCount() + 1);
                if (log.getRetryCount() >= 3) {
                    log.setStatus(CompensationStatus.DEAD_LETTER);
                }
            }
            compensationMapper.updateById(log);
        }
    }
}

分布式事务选型的决策路径:强一致性要求用Seata AT,高吞吐场景用消息最终一致性,两者混用时要保证同一条业务链路只用一种模型,避免嵌套事务带来的复杂性爆炸。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/springboot-fen-bu-shi-shi-wu-shi-zhan-fang-an-seataat-mo/

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

相关推荐