Spring Boot微服务分布式事务方案:Seata AT模式与消息最终一致性实战

微服务架构下,跨服务的数据操作需要分布式事务保证一致性。与单机事务不同,分布式事务面临网络分区、超时重试、部分提交等复杂问题。本文对比主流分布式事务方案,以Spring Boot集成Seata和RocketMQ消息事务为例,演示最终一致性架构的完整实现。

分布式事务方案对比与选型

主流方案各有适用场景和权衡点:

2PC(两阶段提交):强一致性,但同步阻塞、性能差,单点协调者故障会导致所有参与者阻塞。适合传统数据库跨库事务,不适合微服务高并发场景。

TCC(Try-Confirm-Cancel):业务侵入性强,每个操作需实现Try、Confirm、Cancel三个接口。性能优于2PC,适合资金类强一致性场景。开发成本高,业务逻辑复杂时容易出错。

Saga:长事务拆分为多个本地事务,每个事务有对应的补偿操作。适合流程长、涉及多个服务的业务场景。不提供隔离性,中间状态可见。

消息最终一致性:通过消息中间件保证跨服务操作最终执行。开发成本最低,性能最好,但有短暂不一致窗口。适合对实时一致性要求不高的高并发场景。

选型原则:资金交易选TCC,订单流程选Saga,通用业务操作选消息最终一致性,简单跨库事务选Seata AT模式。

Seata AT模式集成Spring Boot

Seata AT模式通过SQL解析和undo_log自动生成补偿操作,业务代码侵入性最低。集成步骤:

<!-- pom.xml -->
<dependency>
    <groupId>io.seata</groupId>
    <artifactId>seata-spring-boot-starter</artifactId>
    <version>2.0.0</version>
</dependency>
<dependency>
    <groupId>com.alibaba</groupId>
    <artifactId>druid</artifactId>
    <version>1.2.20</version>
</dependency>
# application.yml
seata:
  enabled: true
  application-id: order-service
  tx-service-group: my_tx_group
  service:
    vgroup-mapping:
      my_tx_group: default
    grouplist:
      default: 127.0.0.1:8091
  data-source-proxy-mode: AT

order_service和inventory_service两个微服务参与全局事务。订单服务创建订单时调用库存服务扣减库存,任一环节失败则整体回滚:

// OrderService.java
@Service
public class OrderService {
    
    @Autowired
    private OrderMapper orderMapper;
    
    @Autowired
    private InventoryFeignClient inventoryClient;
    
    @GlobalTransactional(timeoutMills = 60000, name = "create-order-tx")
    public Order createOrder(OrderRequest request) {
        // 1. 创建订单(本地事务)
        Order order = new Order();
        order.setUserId(request.getUserId());
        order.setProductId(request.getProductId());
        order.setQuantity(request.getQuantity());
        order.setAmount(request.getAmount());
        order.setStatus("CREATED");
        orderMapper.insert(order);
        
        // 2. 远程调用扣减库存(远程事务)
        Result<Boolean> result = inventoryClient.deduct(
            request.getProductId(), 
            request.getQuantity()
        );
        
        if (result == null || !result.isSuccess()) {
            throw new BusinessException("库存扣减失败: " + 
                (result != null ? result.getMessage() : "服务调用异常"));
        }
        
        // 3. 更新订单状态
        order.setStatus("CONFIRMED");
        orderMapper.updateStatus(order);
        
        return order;
    }
}
// InventoryService.java(库存服务)
@Service
public class InventoryService {
    
    @Autowired
    private InventoryMapper inventoryMapper;
    
    public boolean deduct(Long productId, Integer quantity) {
        Inventory inventory = inventoryMapper.selectByProductId(productId);
        if (inventory == null || inventory.getStock() < quantity) {
            throw new BusinessException("库存不足");
        }
        inventory.setStock(inventory.getStock() - quantity);
        inventoryMapper.updateStock(inventory);
        return true;
    }
}

@GlobalTransactional注解标记全局事务入口,Seata自动在各参与者数据库中生成undo_log记录。任一参与者抛出异常,TC(事务协调器)通知所有参与者执行undo回滚。业务代码只需关注正常逻辑和异常抛出,补偿操作由框架自动完成。

每个参与者的数据库需创建undo_log表:

CREATE TABLE `undo_log` (
    `branch_id` BIGINT NOT NULL COMMENT '分支事务ID',
    `xid` VARCHAR(128) NOT NULL COMMENT '全局事务ID',
    `context` VARCHAR(128) NOT NULL COMMENT '上下文',
    `rollback_info` LONGBLOB NOT NULL COMMENT '回滚信息',
    `log_status` INT NOT NULL COMMENT '状态',
    `log_created` DATETIME NOT NULL COMMENT '创建时间',
    `log_modified` DATETIME NOT NULL COMMENT '修改时间',
    PRIMARY KEY (`branch_id`),
    KEY `idx_xid` (`xid`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

消息最终一致性方案实现

对于不需要强一致性的场景,基于RocketMQ事务消息实现最终一致性。以订单创建后发送通知为例:

// TransactionMQProducer配置
@Configuration
public class MQConfig {
    
    @Bean
    public TransactionMQProducer transactionProducer() throws MQClientException {
        TransactionMQProducer producer = new TransactionMQProducer("order_tx_group");
        producer.setNamesrvAddr("127.0.0.1:9876");
        producer.setTransactionListener(new OrderTransactionListener());
        producer.start();
        return producer;
    }
}

// 事务监听器
@Component
public class OrderTransactionListener implements TransactionListener {
    
    @Autowired
    private OrderMapper orderMapper;
    
    @Override
    public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
        try {
            String orderId = msg.getKeys();
            Order order = JSON.parseObject(
                new String(msg.getBody()), Order.class
            );
            orderMapper.insert(order);
            return LocalTransactionState.COMMIT_MESSAGE;
        } catch (Exception e) {
            log.error("本地事务执行失败", e);
            return LocalTransactionState.ROLLBACK_MESSAGE;
        }
    }
    
    @Override
    public LocalTransactionState checkLocalTransaction(MessageExt msg) {
        String orderId = msg.getKeys();
        Order order = orderMapper.selectById(orderId);
        if (order != null) {
            return LocalTransactionState.COMMIT_MESSAGE;
        }
        return LocalTransactionState.ROLLBACK_MESSAGE;
    }
}
// 发送事务消息
@Service
public class OrderMessageService {
    
    @Autowired
    private TransactionMQProducer producer;
    
    public void sendOrderMessage(Order order) throws MQClientException {
        Message msg = new Message(
            "ORDER_TOPIC",
            "CREATE",
            order.getId().toString(),
            JSON.toJSONString(order).getBytes()
        );
        
        TransactionSendResult result = producer.sendMessageInTransaction(msg, null);
        
        if (result.getLocalTransactionState() == LocalTransactionState.ROLLBACK_MESSAGE) {
            throw new BusinessException("订单创建失败,事务已回滚");
        }
    }
}

事务消息保证本地事务和消息发送的原子性。半消息发送成功后执行本地事务,本地事务成功则提交消息(消费者可见),失败则回滚消息。如果Producer宕机未返回确认,MQ Server定期回调checkLocalTransaction方法确定事务状态。消费端通过幂等性设计处理可能的重发消息。

消费端幂等性设计与防重消费

消息可能因网络重试被多次投递,消费端必须实现幂等性。基于Redis实现消费去重:

@Component
@RocketMQMessageListener(
    topic = "ORDER_TOPIC", 
    consumerGroup = "notification_group"
)
public class OrderNotificationConsumer 
        implements RocketMQListener<OrderMessage> {
    
    @Autowired
    private RedisTemplate<String, String> redisTemplate;
    
    @Autowired
    private NotificationService notificationService;
    
    @Override
    public void onMessage(OrderMessage message) {
        String key = "order:consumed:" + message.getId();
        
        Boolean firstConsume = redisTemplate.opsForValue()
            .setIfAbsent(key, "1", 24, TimeUnit.HOURS);
        
        if (Boolean.FALSE.equals(firstConsume)) {
            log.info("消息已消费,跳过: orderId={}", message.getId());
            return;
        }
        
        try {
            notificationService.sendOrderNotification(message);
        } catch (Exception e) {
            redisTemplate.delete(key);
            throw e;
        }
    }
}

SETNX + 过期时间的方案保证同一消息在同一时间窗内只被消费一次。业务异常时删除标记允许重试,避免因临时故障导致消息永久丢失。对于严格不允许重复的场景,可在数据库层面增加唯一索引作为兜底。

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

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

相关推荐