单体架构使用数据库本地事务即可保证 ACID,但微服务架构下业务操作跨越多个服务和数据库,分布式事务成为必须解决的问题。本文系统讲解从强一致性到最终一致性的各种分布式事务方案。
1. 分布式事务分类
| 方案 | 一致性 | 性能 | 复杂度 | 适用场景 |
|---|---|---|---|---|
| 2PC | 强一致 | 低 | 中 | 短事务、低并发 |
| 3PC | 强一致 | 低 | 高 | 2PC 改进,较少使用 |
| TCC | 最终一致 | 高 | 高 | 高并发、核心资产业务 |
| Saga | 最终一致 | 高 | 中 | 长事务、业务流程 |
| 本地消息表 | 最终一致 | 高 | 中 | 异步场景 |
| Seata AT | 最终一致 | 较高 | 低 | 侵入小、快速接入 |
| 最大努力通知 | 最终一致 | 高 | 低 | 对账场景 |
2. 两阶段提交 (2PC)
2PC (Two-Phase Commit) 是最经典的分布式事务协议,由协调者 (Coordinator) 和参与者 (Participants) 组成。
2.1 协议流程
Phase 1 (投票阶段):
Coordinator → Prepare → Participant A
Coordinator → Prepare → Participant B
Coordinator → Prepare → Participant C
A: Undo/Redo log written, return Yes
B: Undo/Redo log written, return Yes
C: Failed to prepare, return No
Phase 2 (提交/回滚):
Coordinator (收到 No) → Rollback → A
Coordinator (收到 No) → Rollback → B
2.2 2PC 的状态机
Coordinator 状态:
INIT → PREPARING → PREPARED → COMMITTING → COMMITTED
→ ABORTING → ABORTED
Participant 状态:
INIT → READY → COMMITTED/ABORTED
2.3 2PC 的问题
| 问题 | 原因 | 后果 |
|---|---|---|
| 同步阻塞 | 参与者需锁定资源等待协调者指令 | 性能差 |
| 单点故障 | 协调者宕机,参与者一直阻塞 | 可用性低 |
| 数据不一致 | 协调者发送 Commit 后宕机,部分参与者未收到 | 脑裂 |
| _TIMEOUT | 网络超时导致不确定性 | 需人工干预 |
2.4 Java 实现示例
@Component
public class XAOrderService {
@Autowired
private JdbcTemplate orderJdbc;
@Autowired
private JdbcTemplate inventoryJdbc;
@Transactional(rollbackFor = Exception.class)
public void createOrder(OrderRequest request) {
// XA 两阶段提交,通过 JTA 管理
// atomikos 或 narayana 实现
UserTransaction ut = userTransactionManager.getUserTransaction();
try {
ut.begin();
// 操作订单库
orderJdbc.update("INSERT INTO orders (user_id, amount) VALUES (?, ?)",
request.getUserId(), request.getAmount());
// 操作库存库
inventoryJdbc.update("UPDATE inventory SET count = count - ? WHERE sku = ?",
request.getQuantity(), request.getSku());
ut.commit();
} catch (Exception e) {
try {
ut.rollback();
} catch (Exception ex) {
log.error("Rollback failed", ex);
}
throw new BusinessException("Order creation failed", e);
}
}
}
2.5 3PC 改进
3PC 增加了一个 CanCommit 阶段,减少协调者宕机导致的阻塞时间,但实现复杂且网络开销更大,实际很少使用。
CanCommit → PreCommit → DoCommit
3. TCC (Try-Confirm-Cancel)
TCC 由阿里提出,将业务操作拆分为三个阶段,通过业务逻辑保证最终一致性。
3.1 TCC 三阶段
| 阶段 | 操作 | 资源状态 |
|---|---|---|
| Try | 预留资源,执行业务检查 | 资源被冻结/预扣 |
| Confirm | 真正执行业务 | 资源正式扣除 |
| Cancel | 释放预留资源 | 回滚到初始状态 |
用户下单:
Try: 冻结库存 1,冻结余额 100
库存: 10 → 可用 9, 冻结 1
余额: 500 → 可用 400, 冻结 100
Confirm: 确认订单
库存: 可用 9, 冻结 0, 已售 +1
余额: 可用 400, 冻结 0, 已扣 100
Cancel: 取消订单
库存: 可用 10, 冻结 0
余额: 可用 500, 冻结 0
3.2 TCC 实现
public interface InventoryTccAction {
@TwoPhaseBusinessAction(name = "inventoryTccAction",
tryMethod = "tryDeduct",
confirmMethod = "commit",
rollbackMethod = "rollback")
boolean tryDeduct(@BusinessActionContextParameter(paramName = "sku") String sku,
@BusinessActionContextParameter(paramName = "count") int count);
boolean commit(BusinessActionContext context);
boolean rollback(BusinessActionContext context);
}
@Service
public class InventoryTccActionImpl implements InventoryTccAction {
@Autowired
private InventoryMapper inventoryMapper;
@Override
public boolean tryDeduct(String sku, int count) {
// 检查库存
Inventory inventory = inventoryMapper.selectBySku(sku);
if (inventory.getAvailable() < count) {
throw new BusinessException("Insufficient inventory");
}
// 冻结库存
return inventoryMapper.freeze(sku, count) > 0;
}
@Override
public boolean commit(BusinessActionContext context) {
String sku = context.getActionContext("sku");
int count = Integer.parseInt(context.getActionContext("count"));
// 将冻结库存转为实际扣减
return inventoryMapper.confirmDeduct(sku, count) > 0;
}
@Override
public boolean rollback(BusinessActionContext context) {
String sku = context.getActionContext("sku");
int count = Integer.parseInt(context.getActionContext("count"));
// 释放冻结库存
return inventoryMapper.unfreeze(sku, count) > 0;
}
}
3.3 TCC 注意事项
- 幂等性:Confirm 和 Cancel 必须幂等,可能因网络重试而多次执行
- 空回滚:Try 尚未执行就触发 Cancel,需要记录 Try 是否执行过
- 悬挂:Cancel 先执行,Try 后执行(超时导致的乱序),需防悬挂
// 幂等性控制
@Override
public boolean commit(BusinessActionContext context) {
String xid = context.getXid();
// 检查是否已提交
if (tccLogMapper.isCommitted(xid)) {
return true; // 已处理,直接返回
}
// ... 执行业务
tccLogMapper.markCommitted(xid);
return true;
}
4. Saga 模式
Saga 将长事务拆分为多个本地事务,每个本地事务提交后立即释放资源,通过补偿操作回滚。
4.1 两种 Saga 实现
| 类型 | 机制 | 代表框架 |
|---|---|---|
| 编排式 (Choreography) | 每个服务完成本地事务后发送事件触发下一个服务 | 事件驱动 |
| 编排式 (Orchestration) | 中央协调器统一调度各服务的执行和补偿 | Camunda, Apache Camel |
4.2 编排式 Saga (Orchestration)
@Service
public class OrderSagaOrchestrator {
@Autowired
private StateMachineFactory<OrderStatus, OrderEvent> stateMachineFactory;
public void startOrderSaga(Order order) {
StateMachine<OrderStatus, OrderEvent> sm = stateMachineFactory.getStateMachine();
sm.start();
// 状态流转:
// CREATED → [create inventory reservation] → INVENTORY_RESERVED
// → [create payment] → PAYMENT_COMPLETED
// → [ship order] → SHIPPED
// → [complete] → COMPLETED
// 补偿链(反向执行):
// PAYMENT_COMPLETED → [refund payment] → INVENTORY_RESERVED
// → [release inventory] → CREATED
}
}
// Spring State Machine 配置
@Configuration
public class OrderSagaConfig extends StateMachineConfigurerAdapter<OrderStatus, OrderEvent> {
@Override
public void configure(StateMachineTransitionConfigurer<OrderStatus, OrderEvent> transitions)
throws Exception {
transitions
.withExternal()
.source(OrderStatus.CREATED)
.target(OrderStatus.INVENTORY_RESERVED)
.event(OrderEvent.RESERVE_INVENTORY)
.action(reserveInventoryAction())
.and()
.withExternal()
.source(OrderStatus.INVENTORY_RESERVED)
.target(OrderStatus.PAYMENT_COMPLETED)
.event(OrderEvent.PROCESS_PAYMENT)
.action(processPaymentAction())
.and()
// 补偿路径
.withExternal()
.source(OrderStatus.PAYMENT_COMPLETED)
.target(OrderStatus.INVENTORY_RESERVED)
.event(OrderEvent.REFUND_PAYMENT)
.action(refundPaymentAction());
}
}
4.3 事件编排式 Saga
// 订单服务
@Transactional
public void createOrder(OrderRequest request) {
Order order = orderRepository.save(new Order(request));
eventPublisher.publish(new OrderCreatedEvent(order.getId(), request));
}
// 库存服务监听
@EventListener
@Transactional
public void onOrderCreated(OrderCreatedEvent event) {
inventoryService.reserve(event.getSku(), event.getQuantity());
eventPublisher.publish(new InventoryReservedEvent(event.getOrderId()));
}
// 支付服务监听
@EventListener
@Transactional
public void onInventoryReserved(InventoryReservedEvent event) {
try {
paymentService.charge(event.getOrderId());
eventPublisher.publish(new PaymentCompletedEvent(event.getOrderId()));
} catch (Exception e) {
eventPublisher.publish(new PaymentFailedEvent(event.getOrderId()));
}
}
// 补偿:支付失败释放库存
@EventListener
@Transactional
public void onPaymentFailed(PaymentFailedEvent event) {
inventoryService.release(event.getOrderId());
orderService.cancel(event.getOrderId());
}
5. 本地消息表
基于可靠消息实现最终一致性,适用于异步场景。
5.1 核心思想
业务操作和消息记录在同一个本地事务中:
BEGIN
INSERT INTO orders (...) -- 业务表
INSERT INTO message_queue (topic, payload, status) -- 消息表
COMMIT
后台任务轮询消息表,发送到消息队列
消息消费方处理完毕后 ACK,消息表状态更新为 DONE
5.2 实现
@Service
public class ReliableMessageService {
@Transactional
public void createOrderWithMessage(OrderRequest request) {
// 1. 保存订单
Order order = orderRepository.save(new Order(request));
// 2. 记录消息(同库同事务)
OutboxMessage message = new OutboxMessage();
message.setTopic("order_created");
message.setPayload(JsonUtils.toJson(new OrderCreatedEvent(order)));
message.setStatus(MessageStatus.PENDING);
message.setRetryCount(0);
outboxRepository.save(message);
}
// 定时任务:轮询消息表
@Scheduled(fixedRate = 5000)
public void pollOutboxMessages() {
List<OutboxMessage> pending = outboxRepository
.findByStatusAndRetryCountLessThan(MessageStatus.PENDING, 3);
for (OutboxMessage msg : pending) {
try {
kafkaTemplate.send(msg.getTopic(), msg.getPayload()).get(5, TimeUnit.SECONDS);
msg.setStatus(MessageStatus.SENT);
} catch (Exception e) {
msg.setRetryCount(msg.getRetryCount() + 1);
}
outboxRepository.save(msg);
}
}
// 消息消费确认
@KafkaListener(topics = "order_created")
public void handleOrderCreated(String payload, Acknowledgment ack) {
OrderCreatedEvent event = JsonUtils.fromJson(payload, OrderCreatedEvent.class);
// 幂等处理
if (processedMessageRepository.existsByMessageId(event.getMessageId())) {
ack.acknowledge();
return;
}
inventoryService.deduct(event.getSku(), event.getQuantity());
processedMessageRepository.save(new ProcessedMessage(event.getMessageId()));
ack.acknowledge();
}
}
6. Seata AT 模式
Seata (Simple Extensible Autonomous Transaction Architecture) 是阿里开源的分布式事务解决方案,AT 模式对业务零侵入。
6.1 AT 模式原理
1. 一阶段:业务 SQL → 解析 SQL → 查询前镜像 → 执行业务 SQL → 查询后镜像 → 记录 UNDO_LOG
2. 二阶段成功:异步删除 UNDO_LOG
3. 二阶段回滚:用 UNDO_LOG 的前镜像生成反向 SQL 回滚
6.2 部署与配置
# application.yml
seata:
enabled: true
application-id: ${spring.application.name}
tx-service-group: my_tx_group
service:
vgroup-mapping:
my_tx_group: default
client:
rm:
async-commit-buffer-limit: 10000
tm:
commit-retry-count: 5
rollback-retry-count: 5
datasource:
proxy-datasource: true
6.3 业务代码(零侵入)
@Service
public class BusinessService {
@Autowired
private OrderService orderService;
@Autowired
private StorageService storageService;
@Autowired
private AccountService accountService;
@GlobalTransactional(name = "create-order", rollbackFor = Exception.class)
public void createOrder(Order order) {
// 1. 创建订单
orderService.create(order);
// 2. 扣减库存
storageService.deduct(order.getCommodityCode(), order.getCount());
// 3. 扣减账户余额
accountService.debit(order.getUserId(), order.getMoney());
// 任意步骤异常,全局回滚
}
}
Seata 代理数据源自动拦截 SQL,生成 UNDO_LOG。
6.4 Seata 四种模式对比
| 模式 | 侵入性 | 适用场景 | 性能 |
|---|---|---|---|
| AT | 零侵入 | CRUD 为主 | 较高 |
| TCC | 高(需实现三接口) | 复杂业务 | 高 |
| Saga | 中(状态机配置) | 长事务、业务流程 | 高 |
| XA | 低 | 兼容传统 2PC | 低 |
7. 方案选择决策树
是否需要实时强一致性?
├── 是 → 2PC/XA(短事务、低并发)
│ 或:业务层面通过单服务聚合避免分布式事务
└── 否 → 最终一致性
├── 是否需要高并发?
│ ├── 是 → TCC(金融核心、库存扣减)
│ └── 否 → Saga(业务流程、长事务)
└── 是否需异步解耦?
├── 是 → 本地消息表 / Outbox
└── 否 → Seata AT(快速接入、零侵入)
总结
| 维度 | 2PC/XA | TCC | Saga | 本地消息表 | Seata AT |
|---|---|---|---|---|---|
| 一致性 | 强 | 最终 | 最终 | 最终 | 最终 |
| 性能 | 低 | 高 | 高 | 高 | 较高 |
| 复杂度 | 中 | 高 | 中 | 中 | 低 |
| 侵入性 | 低 | 高 | 中 | 中 | 零 |
| 回滚能力 | 自动 | 业务补偿 | 业务补偿 | 无法自动回滚 | 自动 |
| 适用 | 传统系统 | 核心资产 | 业务流程 | 异步通知 | 通用 |
分布式事务选型建议:
- 能不用就不用:通过业务设计避免分布式事务(如将操作收敛到单一服务)
- 最终一致性优先:绝大多数场景最终一致性足够
- TCC 用于核心资产:资金、库存等高并发且需精确控制
- Seata AT 快速接入:已有系统改造,不想改动业务代码
- Saga 处理长事务:订单流程、审批流等跨多服务的业务流程
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。