本文目录导读:

Java 分布式事务的解决方案通常根据业务场景的一致性要求(强一致性 vs. 最终一致性)以及技术栈(是否使用微服务、是否依赖特定中间件)来选择。
下面通过几个典型的案例场景,结合代码片段和核心思路,来讲解如何解决分布式事务问题。
案例背景:经典的下单扣库存
场景:在微服务架构中,有两个服务:
- 订单服务:创建订单。
- 库存服务:扣减库存。
问题:如果订单创建成功,但扣库存失败(或反之),会导致数据不一致。
基于可靠消息的最终一致性(最常用)
适用场景:允许短暂的数据不一致,追求高可用和高性能(大部分互联网业务)。
核心思想:利用本地消息表或 RocketMQ 事务消息,保证“本地事务”与“发消息”的原子性,然后由消费者(库存服务)保证最终执行成功。
案例流程(使用 RocketMQ 事务消息)
- 订单服务(生产者):
- 发送
half(半)消息到 RocketMQ。 - 执行本地事务(插入订单数据)。
- 根据本地事务结果,提交或回滚该
half消息。 - 如果消息提交成功,MQ确保消息至少被送达一次。
- 发送
- 库存服务(消费者):
- 监听
扣库存消息。 - 执行扣减库存(需要保持幂等性,因为消息可能重复投递)。
- 如果扣减成功,手动
ACK;如果失败,重试或转入死信队列人工处理。
- 监听
核心代码示意:
// 订单服务 - 发送事务消息
@Autowired
private RocketMQTemplate rocketMQTemplate;
public void createOrder(Order order) {
// 1. 构建消息
Message<Order> message = MessageBuilder.withPayload(order).build();
// 2. 发送事务消息
// 参数:topic、消息、路由key、参数
rocketMQTemplate.sendMessageInTransaction(
"order-topic",
message,
order.getUserId()
);
}
// 3. 实现 TransactionListener 接口
@Component
public class OrderTransactionListener implements RocketMQLocalTransactionListener {
@Autowired
private OrderMapper orderMapper;
@Override
@Transactional // 执行本地事务
public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
try {
Order order = (Order) msg.getPayload();
// 执行本地订单插入
orderMapper.insert(order);
// 如果本地事务成功,提交消息(让消费者看到)
return RocketMQLocalTransactionState.COMMIT;
} catch (Exception e) {
// 如果本地执行失败,回滚消息(消费者看不到)
return RocketMQLocalTransactionState.ROLLBACK;
}
}
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
// (重要)Broker 长时间没收到 half 消息的确认,会回调此方法进行状态回查
Order order = (Order) msg.getPayload();
// 查询订单是否存在
boolean isExists = orderMapper.selectById(order.getId()) != null;
return isExists ? RocketMQLocalTransactionState.COMMIT : RocketMQLocalTransactionState.ROLLBACK;
}
}
// 库存服务 - 消费消息
@Component
@RocketMQMessageListener(topic = "order-topic", consumerGroup = "inventory-group")
public class InventoryConsumer implements RocketMQListener<Order> {
@Autowired
private InventoryService inventoryService;
@Override
public void onMessage(Order order) {
// 1. 幂等性检查:防止重复消费导致多次扣减
if (deduplicationService.isProcessed(order.getId())) {
log.info("订单 {} 已处理,跳过", order.getId());
return;
}
// 2. 扣减库存
try {
inventoryService.deduct(order.getProductId(), order.getQuantity());
// 3. 记录处理日志(确保幂等)
deduplicationService.markAsProcessed(order.getId());
} catch (Exception e) {
// 重试机制:抛异常会让 MQ 重试(默认重试16次)
throw new RuntimeException("扣库存失败,将重试", e);
}
}
}
TCC(Try-Confirm-Cancel)模式
适用场景:对一致性要求较高,希望手动控制资源锁定粒度(如金融交易、资金账户)。
核心思想:将业务拆分为三个阶段:
- Try:预留资源(如冻结库存、冻结金额)。
- Confirm:确认执行(实际扣减)。
- Cancel:回滚释放(释放冻结资源)。
案例场景:账号间转账(A扣钱,B加钱)
实现:使用 TCC 框架(如 seata、tcc-transaction)。
// 通过 @LocalTCC 注解定义
@LocalTCC
public interface AccountService {
// Try 阶段:尝试锁定资源
@TwoPhaseBusinessAction(
name = "transfer",
commitMethod = "confirm",
rollbackMethod = "cancel"
)
boolean tryDebit(@BusinessActionContextParameter(paramName = "accountId") String accountId,
@BusinessActionContextParameter(paramName = "amount") double amount);
// Confirm 阶段:确认执行
boolean confirm(BusinessActionContext actionContext);
// Cancel 阶段:回滚
boolean cancel(BusinessActionContext actionContext);
}
服务实现:
@Service
public class AccountServiceImpl implements AccountService {
@Autowired
private AccountMapper accountMapper;
@Override
public boolean tryDebit(String accountId, double amount) {
// Try 阶段:不做扣减,而是在数据库建一张 “冻结表”
// 或者直接在用户余额字段边添加一个 “冻结金额” 字段
// UPDATE account SET frozen = frozen + ? WHERE id = ? AND balance - frozen >= ?
int result = accountMapper.freezeAmount(accountId, amount);
return result > 0;
}
@Override
public boolean confirm(BusinessActionContext actionContext) {
// Confirm:真正扣钱,释放冻结
String accountId = actionContext.getActionContext("accountId").toString();
double amount = Double.parseDouble(actionContext.getActionContext("amount").toString());
// UPDATE account SET balance = balance - ?, frozen = frozen - ? WHERE id = ?
accountMapper.realDebit(accountId, amount);
return true;
}
@Override
public boolean cancel(BusinessActionContext actionContext) {
// Cancel:释放冻结(回滚Try)
String accountId = actionContext.getActionContext("accountId").toString();
double amount = Double.parseDouble(actionContext.getActionContext("amount").toString());
// UPDATE account SET frozen = frozen - ? WHERE id = ?
accountMapper.unfreezeAmount(accountId, amount);
return true;
}
}
Seata AT(自动补偿)模式
适用场景:不想自己写TCC的代码,希望框架自动拦截SQL做反向SQL(Undo Log)。
核心思想:Seata 框架通过代理数据源,自动解析 SQL,生成“镜像”快照,事务提交时自动删除镜像;事务回滚时自动生成反向 SQL 恢复数据。
配置与使用(极其简单):
- 启动 Seata Server(TC 协调者)。
- 引入依赖
seata-spring-boot-starter。 - 在需要分布式事务的方法上添加
@GlobalTransactional注解。
@Service
public class BusinessService {
@Autowired
private OrderFeignClient orderClient; // 调用订单服务
@Autowired
private InventoryFeignClient inventoryClient; // 调用库存服务
// 核心操作:标记这个方法是分布式事务
@GlobalTransactional(name = "create-order", rollbackFor = Exception.class)
public void createOrder(Order order) {
// 1. 创建订单(本地事务)
orderClient.create(order);
// 2. 扣减库存(远程调用,也是本地事务)
inventoryClient.deduct(order.getProductId(), order.getQuantity());
// 3. 如果前面任一步失败,Seata 会自动让所有参与方回滚
// 不需要写任何回滚代码!
}
}
如何选择?
| 需求 | 推荐方案 | 理由 |
|---|---|---|
| 强一致性 | Seata AT / TCC | 要么全成功,要么全回滚。 |
| 高并发、最终一致 | RocketMQ 事务消息 | 性能好,解耦,抗压。 |
| 跨数据库/跨语言 | TCC | 不依赖特定中间件,灵活性高。 |
| 不想写回滚代码 | Seata AT | 自动通过 Undo Log 回滚。 |
| 长事务 | Saga 模式(状态机) | 通过事件驱动,适合流程长、步骤多的场景。 |
解决分布式事务的关键点
- 幂等性:无论哪种方案(尤其是消息和重试),接口必须支持幂等,防止重复执行造成数据错误。
- 资源长时间锁定:TCC 的 Try 阶段要避免长时间冻结资源(如库存、资金),否则影响并发。
- 最终一致:MQ 方案通常无法做到实时一致性,适合下单、发券等场景;不适合实时扣减唯一资源(如秒杀库存严格扣减)。
- 日志与兜底:无论哪种方案,都需要完善的失败日志和人工补偿机制,因为极端情况下(如网络分区)自动方案可能无法完全达到预期。
在实际开发中,基于可靠消息的最终一致性(方案一)是目前业务最广泛、落地成本较低的方案;而 Seata AT 则适合单体转型微服务,希望无侵入实现 ACID 的场景。