微服务架构下,跨服务的数据操作需要分布式事务保证一致性。与单机事务不同,分布式事务面临网络分区、超时重试、部分提交等复杂问题。本文对比主流分布式事务方案,以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/