为什么需要分布式事务?
在微服务架构中,一个业务操作往往涉及多个服务,例如下单需要扣库存、生成订单、扣减账户余额。如果某个步骤失败,如何保证数据最终一致?传统数据库ACID事务无法跨服务,因此需要分布式事务方案。本文将手把手带你实现三种主流方案:Saga、TCC和最终一致性(基于消息队列),并分享踩坑经验。
场景定义
假设一个电商下单流程:
- 订单服务创建订单(状态为待支付)
- 库存服务扣减库存
- 账户服务扣减余额
要求:要么全部成功,要么全部回滚(或补偿)。
方案一:Saga 模式(基于事件/编排)
Saga 将长事务拆分为多个本地事务,每个本地事务有对应的补偿操作。如果某个步骤失败,则反向执行补偿。
实现思路
使用事件编排:每个服务完成后发布事件,触发下一步。失败时发布补偿事件。
代码示例(Node.js + RabbitMQ)
// orchestrator.js - 编排器
const amqp = require('amqplib');
async function startSaga(orderId, productId, userId, amount) {
const conn = await amqp.connect('amqp://localhost');
const channel = await conn.createChannel();
const exchange = 'saga_exchange';
await channel.assertExchange(exchange, 'topic', { durable: true });
// 步骤1: 创建订单
const orderEvent = { orderId, userId, amount, status: 'PENDING' };
channel.publish(exchange, 'order.create', Buffer.from(JSON.stringify(orderEvent)));
console.log('Order created event published');
// 监听后续事件...
// 实际应使用消费者处理回调,此处简化
}
// order-service.js - 订单服务消费者
channel.consume('order_queue', async (msg) => {
const event = JSON.parse(msg.content.toString());
if (event.status === 'PENDING') {
// 创建订单(本地事务)
await db.insertOrder(event);
// 发布订单已创建事件
channel.publish(exchange, 'order.created', Buffer.from(JSON.stringify(event)));
channel.ack(msg);
} else if (event.status === 'COMPENSATE') {
// 补偿:取消订单
await db.updateOrder(event.orderId, { status: 'CANCELLED' });
channel.ack(msg);
}
});
// 类似实现库存和账户服务
优缺点
- 优点:简单,易于理解,适合长事务
- 缺点:需要实现补偿逻辑,可能导致数据中间状态
方案二:TCC 模式(Try-Confirm-Cancel)
TCC 将每个服务操作分为两个阶段:Try(预留资源)、Confirm(确认执行)、Cancel(取消释放)。
实现思路
每个服务提供 Try、Confirm、Cancel 三个接口。事务管理器协调。
代码示例(Java + Spring Boot)
// 账户服务 TCC 接口
public interface AccountTccService {
@TwoPhaseBusinessAction(name = "deduct", commitMethod = "confirm", rollbackMethod = "cancel")
void tryDeduct(BusinessActionContext context,
@BusinessActionContextParameter("userId") String userId,
@BusinessActionContextParameter("amount") double amount);
boolean confirm(BusinessActionContext context);
boolean cancel(BusinessActionContext context);
}
// 实现
@Service
public class AccountTccServiceImpl implements AccountTccService {
@Override
public void tryDeduct(BusinessActionContext context, String userId, double amount) {
// Try:冻结金额
accountDao.freezeBalance(userId, amount);
}
@Override
public boolean confirm(BusinessActionContext context) {
// Confirm:扣除冻结金额
String userId = context.getActionContext("userId").toString();
double amount = Double.parseDouble(context.getActionContext("amount").toString());
accountDao.deductFrozen(userId, amount);
return true;
}
@Override
public boolean cancel(BusinessActionContext context) {
// Cancel:解冻金额
String userId = context.getActionContext("userId").toString();
double amount = Double.parseDouble(context.getActionContext("amount").toString());
accountDao.unfreezeBalance(userId, amount);
return true;
}
}
注意事项
- Try 阶段必须保证资源预留成功,否则不执行 Confirm
- Confirm 和 Cancel 必须幂等
- 需要事务管理器(如 Seata、ByteTCC)
方案三:最终一致性(基于消息队列)
利用消息队列异步处理,保证最终一致。本地事务和消息发送绑定,通过重试和幂等实现。
实现思路
- 订单服务本地事务插入订单,同时发送一条“扣库存”消息到队列
- 库存服务消费消息,执行扣减,成功后发送“扣余额”消息
- 每个步骤失败则重试,直到成功
代码示例(RocketMQ 事务消息)
// 订单服务发送半消息
TransactionMQProducer producer = new TransactionMQProducer();
producer.setNamesrvAddr("localhost:9876");
producer.setTransactionListener(new TransactionListener() {
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
// 执行本地事务(创建订单)
Order order = (Order) arg;
try {
orderDao.insert(order);
return LocalTransactionState.COMMIT_MESSAGE;
} catch (Exception e) {
return LocalTransactionState.ROLLBACK_MESSAGE;
}
}
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
// 回查订单状态
String orderId = msg.getKeys();
Order order = orderDao.selectById(orderId);
if (order != null) {
return LocalTransactionState.COMMIT_MESSAGE;
}
return LocalTransactionState.ROLLBACK_MESSAGE;
}
});
producer.start();
// 发送半消息
Message msg = new Message("inventory_topic", "", orderId, "扣减库存".getBytes());
SendResult result = producer.sendMessageInTransaction(msg, order);
幂等性保证
消费者必须实现幂等(例如通过唯一键去重):
// 库存服务消费
@RocketMQMessageListener(topic = "inventory_topic", consumerGroup = "inventory_group")
public class InventoryConsumer implements RocketMQListener<String> {
@Override
public void onMessage(String message) {
String orderId = message;
// 使用订单ID作为幂等键
if (deduplicationService.isProcessed(orderId)) {
return;
}
// 扣减库存
inventoryDao.deduct(orderId);
deduplicationService.markProcessed(orderId);
}
}
对比与选型建议
| 方案 | 一致性 | 复杂度 | 性能 | 适用场景 | |——|——–|——–|——|———-| | Saga | 最终一致 | 中 | 高 | 长事务,允许中间状态 | | TCC | 强一致(两阶段) | 高 | 中 | 短事务,对一致性要求高 | | 最终一致性 | 最终一致 | 低 | 高 | 异步场景,容忍延迟 |
最佳实践:
- 优先考虑最终一致性,避免分布式事务
- 如果必须,Saga 适合业务流程长、步骤多的场景
- TCC 适合资金类业务,但实现成本高
常见坑与踩坑经验
- Saga 补偿逻辑必须幂等:补偿可能重复执行
- TCC Cancel 可能失败:需要重试机制,且 Confirm 和 Cancel 必须幂等
- 消息队列顺序:使用分区保证同一订单的消息顺序
- 事务消息回查超时:合理设置回查间隔,避免堆积
总结
本文通过代码示例详细介绍了 Saga、TCC 和最终一致性三种分布式事务方案。实际项目中建议优先使用消息队列实现最终一致性,只有在强一致需求下才考虑 TCC。无论哪种方案,幂等性和补偿机制都是关键。
延伸阅读:
- Seata 框架:AT 模式自动补偿
- RocketMQ 事务消息官方文档
- 分布式事务模式:Saga vs TCC vs AT
希望本文能帮助你在项目中正确选择并实现分布式事务。