分布式事务折腾手记
分布式事务一旦进项目,好看的架构图就没那么管用了。
这三件事分别对应三个不同的服务,任何一个失败都得全部回滚。
问题的开始
线上环境是 Java Spring Boot 微服务架构,订单服务调用库存服务、积分服务和物流服务。初期流量不大的时候,偶尔出现数据不一致也就算了,但随着业务增长,这个问题就藏不住了。
最直观的表现是:库存扣了但没订单,或者订单成功了但积分没加。
# 查询不一致的数据
mysql> SELECT o.order_id, o.user_id, s.stock_quantity, p.points
-> FROM orders o
-> LEFT JOIN stock s ON o.product_id = s.product_id
-> LEFT JOIN user_points p ON o.user_id = p.user_id
-> WHERE o.status = 'COMPLETED'
-> AND (s.stock_quantity < 0 OR p.points IS NULL);
+----------+---------+---------------+--------+
| order_id | user_id | stock_quantity | points |
+----------+---------+---------------+--------+
| 1008721 | 10023 | -3 | NULL |
| 1008725 | 10056 | -1 | NULL |
| 1008730 | 10089 | -2 | 120 |
+----------+---------+---------------+--------+
这还只是能查到的,实际影响更严重的是用户投诉:订单明明成功了,但库存没扣导致超卖,或者积分没加导致用户投诉。
必须得解决。
第一阶段:两阶段提交(2PC)
最先想到的是两阶段提交。理论很简单:有个协调者,让所有参与者先准备,大家都准备好了再一起提交。
理论很美好,实际用的时候问题来了。
坑一:锁等待时间太长
2PC 的核心问题是第一阶段的资源锁定。所有参与者都必须锁定资源直到第二阶段完成,这段时间其他请求都被阻塞。
// 模拟 2PC 实现
@Service
public class OrderService {
@Autowired
private InventoryService inventoryService;
@Autowired
private PointsService pointsService;
@Autowired
private LogisticsService logisticsService;
@Transactional
public void createOrder(Order order) {
// 第一阶段:准备
boolean inventoryPrepared = inventoryService.prepareDeduct(order.getProductId(), order.getQuantity());
boolean pointsPrepared = pointsService.prepareAdd(order.getUserId(), order.getAmount() * 10);
boolean logisticsPrepared = logisticsService.prepareCreate(order);
// 第二阶段:提交或回滚
if (inventoryPrepared && pointsPrepared && logisticsPrepared) {
inventoryService.commit();
pointsService.commit();
logisticsService.commit();
saveOrder(order);
} else {
inventoryService.rollback();
pointsService.rollback();
logisticsService.rollback();
}
}
}
库存服务里的 prepareDeduct 会把库存锁定:
@Service
public class InventoryService {
@Transactional
public boolean prepareDeduct(Long productId, Integer quantity) {
// 使用悲观锁锁定库存
Stock stock = stockRepository.selectForUpdate(productId);
if (stock.getQuantity() < quantity) {
return false;
}
// 预扣减但暂不提交
stock.setReservedQuantity(stock.getReservedQuantity() + quantity);
stockRepository.update(stock);
return true;
}
}
问题在于,如果积分服务响应慢,库存锁就一直不释放,其他下单请求就会被阻塞。
线上实测:一旦某个服务响应时间超过 200ms,整个下单链路的吞吐量就会下降 60%。
坑二:单点故障
2PC 需要一个协调者,这个协调者挂了怎么办?
如果协调者在第一阶段和第二阶段之间挂了,参与者就不知道是提交还是回滚,资源就一直被锁着。这是典型的阻塞问题。
// 协调者异常时的情况
public void createOrder(Order order) {
// 第一阶段执行完
boolean inventoryPrepared = inventoryService.prepareDeduct(...);
// 协调者挂了,第二阶段没执行
// 库存锁一直不释放
}
坑三:补偿复杂
如果第二阶段某个参与者失败了,已经提交的需要回滚吗?2PC 早期版本没有补偿机制,只能人工介入。
2PC 理论上是正确方案,但实际性能和可用性问题太大,我们决定放弃。
第二阶段:TCC(Try-Confirm-Cancel)
TCC 的思路是:每个业务操作拆成三个阶段——Try 预留资源、Confirm 确认提交、Cancel 取消回滚。
库存服务的 TCC 实现大概是这样:
@Service
public class InventoryTCCService {
@Autowired
private StockRepository stockRepository;
@Autowired
private StockFreezeRepository freezeRepository;
// Try 阶段:预留资源
@Transactional
public boolean tryDeduct(Long productId, Integer quantity) {
Stock stock = stockRepository.findById(productId);
if (stock.getAvailableQuantity() < quantity) {
return false;
}
// 扣减可用库存,增加冻结库存
stock.setAvailableQuantity(stock.getAvailableQuantity() - quantity);
stock.setFrozenQuantity(stock.getFrozenQuantity() + quantity);
stockRepository.update(stock);
// 记录冻结记录,用于 Confirm/Cancel
StockFreeze freeze = new StockFreeze();
freeze.setProductId(productId);
freeze.setQuantity(quantity);
freeze.setStatus("FROZEN");
freeze.setTransactionId(TransactionContext.getTransactionId());
freezeRepository.save(freeze);
return true;
}
// Confirm 阶段:确认提交
@Transactional
public void confirmDeduct(String transactionId) {
StockFreeze freeze = freezeRepository.findByTransactionId(transactionId);
if (freeze == null) {
throw new IllegalStateException("Freeze record not found");
}
if (!"FROZEN".equals(freeze.getStatus())) {
throw new IllegalStateException("Freeze status not FROZEN");
}
Stock stock = stockRepository.findById(freeze.getProductId());
stock.setFrozenQuantity(stock.getFrozenQuantity() - freeze.getQuantity());
stockRepository.update(stock);
freeze.setStatus("CONFIRMED");
freezeRepository.update(freeze);
}
// Cancel 阶段:取消回滚
@Transactional
public void cancelDeduct(String transactionId) {
StockFreeze freeze = freezeRepository.findByTransactionId(transactionId);
if (freeze == null || "CANCELLED".equals(freeze.getStatus())) {
return; // 幂等性
}
Stock stock = stockRepository.findById(freeze.getProductId());
stock.setAvailableQuantity(stock.getAvailableQuantity() + freeze.getQuantity());
stock.setFrozenQuantity(stock.getFrozenQuantity() - freeze.getQuantity());
stockRepository.update(stock);
freeze.setStatus("CANCELLED");
freezeRepository.update(freeze);
}
}
坑一:业务侵入性强
TCC 需要每个服务都实现三个方法,业务逻辑侵入严重。原来一个扣库存的方法,现在要拆成三个,还要维护冻结记录表。
库存表要增加字段:
ALTER TABLE stock ADD COLUMN available_quantity INT NOT NULL DEFAULT 0;
ALTER TABLE stock ADD COLUMN frozen_quantity INT NOT NULL DEFAULT 0;
新增冻结记录表:
CREATE TABLE stock_freeze (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
transaction_id VARCHAR(64) NOT NULL,
product_id BIGINT NOT NULL,
quantity INT NOT NULL,
status VARCHAR(16) NOT NULL,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
UNIQUE KEY uk_transaction_id (transaction_id)
);
每个服务都要建类似的表,改造成本不小。
坑二:空回滚和悬挂问题
TCC 有两个经典问题:空回滚和悬挂。
空回滚:Try 阶段没执行,但 Cancel 阶段执行了。比如网络超时,订单服务以为 Try 失败,调用了 Cancel,但库存服务的 Try 实际成功了。
悬挂:Cancel 执行了,Try 后才执行。比如 Cancel 先到达,把冻结记录删了,但 Try 后才到达又创建了冻结记录。
这两个问题都需要额外的幂等性处理:
// 增加幂等性处理
@Transactional
public void cancelDeduct(String transactionId) {
// 先查冻结记录
StockFreeze freeze = freezeRepository.findByTransactionId(transactionId);
if (freeze == null) {
// 空回滚:Try 没执行过,直接返回成功
return;
}
if ("CANCELLED".equals(freeze.getStatus())) {
// 幂等:已经 Cancel 过了,直接返回成功
return;
}
// 正常 Cancel 逻辑
// ...
}
坑三:Confirm 失败处理
如果 Confirm 阶段某个服务失败了怎么办?比如库存 Confirm 成功了,但积分 Confirm 失败了。
这时候需要人工介入或者定时任务扫描冻结记录,进行重试。但重试又得考虑幂等性、超时时间等问题。
TCC 理论上比 2PC 好一些,但实现复杂度实在太高,我们决定继续找更好的方案。
第三阶段:SAGA 模式
SAGA 的思路更直接:把长事务拆成一系列本地事务,每个本地事务都提交,如果某个失败就执行对应的补偿事务。
库存服务的 SAGA 实现:
@Service
public class InventorySagaService {
@Autowired
private StockRepository stockRepository;
@Autowired
private StockLogRepository stockLogRepository;
// 正向操作:扣库存
@Transactional
public boolean deductStock(Long productId, Integer quantity, String sagaId) {
Stock stock = stockRepository.findById(productId);
if (stock.getQuantity() < quantity) {
return false;
}
stock.setQuantity(stock.getQuantity() - quantity);
stockRepository.update(stock);
// 记录操作日志,用于补偿
StockLog log = new StockLog();
log.setSagaId(sagaId);
log.setProductId(productId);
log.setQuantity(quantity);
log.setOperation("DEDUCT");
log.setStatus("COMPLETED");
stockLogRepository.save(log);
return true;
}
// 补偿操作:加回库存
@Transactional
public void compensateStock(String sagaId) {
List<StockLog> logs = stockLogRepository.findBySagaIdAndStatus(sagaId, "COMPLETED");
for (StockLog log : logs) {
Stock stock = stockRepository.findById(log.getProductId());
stock.setQuantity(stock.getQuantity() + log.getQuantity());
stockRepository.update(stock);
log.setStatus("COMPENSATED");
stockLogRepository.update(log);
}
}
}
订单服务的 SAGA 编排:
@Service
public class OrderSagaService {
@Autowired
private OrderRepository orderRepository;
@Autowired
private InventorySagaService inventoryService;
@Autowired
private PointsSagaService pointsService;
@Autowired
private LogisticsSagaService logisticsService;
public void createOrder(Order order) {
String sagaId = UUID.randomUUID().toString();
boolean success = true;
try {
// 本地事务 1:创建订单
order.setStatus("PROCESSING");
order.setSagaId(sagaId);
orderRepository.save(order);
// 本地事务 2:扣库存
if (!inventoryService.deductStock(order.getProductId(), order.getQuantity(), sagaId)) {
success = false;
}
// 本地事务 3:加积分
if (success) {
if (!pointsService.addPoints(order.getUserId(), order.getAmount() * 10, sagaId)) {
success = false;
}
}
// 本地事务 4:创建物流
if (success) {
if (!logisticsService.createLogistics(order, sagaId)) {
success = false;
}
}
if (success) {
// 全部成功,更新订单状态
order.setStatus("COMPLETED");
orderRepository.update(order);
} else {
// 失败,执行补偿
compensate(sagaId);
}
} catch (Exception e) {
log.error("创建订单失败", e);
compensate(sagaId);
}
}
private void compensate(String sagaId) {
// 补偿顺序:后执行的先补偿
logisticsService.compensateLogistics(sagaId);
pointsService.compensatePoints(sagaId);
inventoryService.compensateStock(sagaId);
// 最后删除订单
Order order = orderRepository.findBySagaId(sagaId);
if (order != null) {
order.setStatus("CANCELLED");
orderRepository.update(order);
}
}
}
优势一:实现简单
SAGA 不需要预留资源,每个服务只实现正向操作和补偿操作,比 TCC 少了一个 Confirm 阶段。
库存表不需要增加冻结字段,只需要一个操作日志表:
CREATE TABLE stock_log (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
saga_id VARCHAR(64) NOT NULL,
product_id BIGINT NOT NULL,
quantity INT NOT NULL,
operation VARCHAR(16) NOT NULL,
status VARCHAR(16) NOT NULL,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
KEY idx_saga_id (saga_id)
);
优势二:非阻塞
SAGA 每个本地事务都是独立的,不锁资源,性能好很多。
线上实测:QPS 从原来的 800 提升到 2500,P99 延迟从 300ms 降到 80ms。
坑一:补偿逻辑复杂
SAGA 的补偿是异步的,可能出现脏读。比如扣库存成功了,订单还没更新状态,用户就能看到库存扣了但订单还在处理中。
解决方案是增加中间状态:
public enum OrderStatus {
PROCESSING, // 处理中
COMPLETED, // 已完成
CANCELLING, // 补偿中
CANCELLED // 已取消
}
坑二:补偿顺序很重要
补偿必须按相反的顺序执行,否则可能出现数据不一致。比如扣库存、加积分、创建物流,补偿顺序应该是取消物流、扣回积分、加回库存。
如果顺序错了,可能物流取消了但积分还没扣回,或者库存加回了但积分还没扣,都是问题。
坑三: Saga 失败本身可能失败
最恶心的是补偿操作本身也可能失败,比如网络问题、服务宕机。
这时候需要一个重试机制,但又要避免无限重试:
@Service
public class CompensationRetryService {
@Autowired
private StockSagaService inventoryService;
@Autowired
private PointsSagaService pointsService;
@Autowired
private LogisticsSagaService logisticsService;
@Scheduled(fixedDelay = 5000)
public void retryFailedCompensations() {
// 查找需要重试的补偿记录
List<StockLog> failedLogs = stockLogRepository.findFailedCompensations();
for (StockLog log : failedLogs) {
if (log.getRetryCount() >= 3) {
// 超过重试次数,报警处理
alertService.sendAlert("Compensation failed: " + log.getSagaId());
continue;
}
try {
inventoryService.compensateStock(log.getSagaId());
log.setStatus("COMPENSATED");
stockLogRepository.update(log);
} catch (Exception e) {
log.setRetryCount(log.getRetryCount() + 1);
stockLogRepository.update(log);
}
}
}
}
最终方案:Seata + SAGA
手写 SAGA 虽然可行,但分布式事务的管理还是很麻烦。最后决定用 Seata,一个开源的分布式事务解决方案。
Seata 支持 AT、TCC、SAGA 和 XA 四种模式,我们选择 SAGA 模式。
Seata 配置
先加依赖:
<dependency>
<groupId>io.seata</groupId>
<artifactId>seata-spring-boot-starter</artifactId>
<version>1.7.0</version>
</dependency>
配置文件:
seata:
application-id: order-service
tx-service-group: default_tx_group
service:
vgroup-mapping:
default_tx_group: default
grouplist:
default: seata-server:8091
registry:
type: nacos
nacos:
server-addr: nacos-server:8848
namespace: seata
group: SEATA_GROUP
config:
type: nacos
nacos:
server-addr: nacos-server:8848
namespace: seata
group: SEATA_GROUP
库存服务的 SAGA 定义:
@Service
public class InventorySagaService {
@Autowired
private StockRepository stockRepository;
// 正向操作
@SagaStart(timeout = 60000)
public boolean deductStock(Long productId, Integer quantity) {
Stock stock = stockRepository.findById(productId);
if (stock.getQuantity() < quantity) {
return false;
}
stock.setQuantity(stock.getQuantity() - quantity);
stockRepository.update(stock);
return true;
}
// 补偿操作
@SagaCompensation
public void compensateDeductStock(Long productId, Integer quantity) {
Stock stock = stockRepository.findById(productId);
stock.setQuantity(stock.getQuantity() + quantity);
stockRepository.update(stock);
}
}
订单服务的编排:
@Service
public class OrderSagaService {
@Autowired
private InventorySagaService inventoryService;
@Autowired
private PointsSagaService pointsService;
@Autowired
private LogisticsSagaService logisticsService;
public void createOrder(Order order) {
// Seata 会自动处理 SAGA 的状态管理和重试
boolean inventoryResult = inventoryService.deductStock(order.getProductId(), order.getQuantity());
boolean pointsResult = pointsService.addPoints(order.getUserId(), order.getAmount() * 10);
boolean logisticsResult = logisticsService.createLogistics(order);
if (inventoryResult && pointsResult && logisticsResult) {
order.setStatus("COMPLETED");
} else {
order.setStatus("CANCELLED");
}
orderRepository.save(order);
}
}
优势一:状态管理自动化
Seata 会自动记录 SAGA 的执行状态,包括每个步骤的状态和重试次数:
CREATE TABLE seata_state_machine_def (
id VARCHAR(64) PRIMARY KEY,
name VARCHAR(128) NOT NULL,
tenant_id VARCHAR(32) NOT NULL,
app_name VARCHAR(32) NOT NULL,
status VARCHAR(16) NOT NULL,
gmt_create TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
gmt_modified TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
);
CREATE TABLE seata_state_machine_instance (
id VARCHAR(64) PRIMARY KEY,
machine_id VARCHAR(64) NOT NULL,
tenant_id VARCHAR(32) NOT NULL,
business_key VARCHAR(128) NOT NULL,
start_params TEXT,
status VARCHAR(16) NOT NULL,
gmt_start TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
gmt_end TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
exception TEXT
);
CREATE TABLE seata_state_execution (
id VARCHAR(64) PRIMARY KEY,
machine_inst_id VARCHAR(64) NOT NULL,
name VARCHAR(128) NOT NULL,
service_name VARCHAR(128) NOT NULL,
service_method VARCHAR(128) NOT NULL,
service_type VARCHAR(16) NOT NULL,
business_key VARCHAR(128) NOT NULL,
status VARCHAR(16) NOT NULL,
retry INT NOT NULL,
gmt_start TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
gmt_end TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
exception TEXT
);
优势二:重试机制完善
Seata 会自动重试失败的步骤,支持配置重试次数和超时时间:
seata:
saga:
compensation-retry-times: 3
compensation-retry-interval: 5000
timeout-max: 60000
坑一:学习成本
Seata 的概念和配置不少,需要搭建 Seata Server,配置 Nacos 注册中心,熟悉各种模式的使用。
如果是小团队或者简单的分布式事务场景,手写 SAGA 可能更快。
坑二:性能开销
Seata 需要维护 SAGA 状态,额外的数据库操作和 RPC 调用会带来性能开销。
线上实测:引入 Seata 后,P99 延迟增加了 10-15ms,但可接受。
幂等性处理
无论哪种方案,幂等性都是必须的。网络超时、服务重启、消息重复投递,这些都可能导致重复调用。
库存服务的幂等性处理:
@Service
public class InventorySagaService {
@Autowired
private StockRepository stockRepository;
@Autowired
private StockLogRepository stockLogRepository;
public boolean deductStock(Long productId, Integer quantity, String sagaId) {
// 先查是否已经执行过
StockLog existingLog = stockLogRepository.findBySagaIdAndOperation(sagaId, "DEDUCT");
if (existingLog != null) {
if ("COMPLETED".equals(existingLog.getStatus())) {
return true; // 已成功,直接返回
}
if ("FAILED".equals(existingLog.getStatus())) {
return false; // 已失败,直接返回
}
}
// 执行扣库存
Stock stock = stockRepository.findById(productId);
if (stock.getQuantity() < quantity) {
// 记录失败
StockLog log = new StockLog();
log.setSagaId(sagaId);
log.setProductId(productId);
log.setQuantity(quantity);
log.setOperation("DEDUCT");
log.setStatus("FAILED");
stockLogRepository.save(log);
return false;
}
stock.setQuantity(stock.getQuantity() - quantity);
stockRepository.update(stock);
// 记录成功
StockLog log = new StockLog();
log.setSagaId(sagaId);
log.setProductId(productId);
log.setQuantity(quantity);
log.setOperation("DEDUCT");
log.setStatus("COMPLETED");
stockLogRepository.save(log);
return true;
}
}
最终效果
改造完成后,线上运行一个月,数据一致性问题基本解决:
- 不再出现库存扣了但没订单的情况
- 不再出现订单成功了但积分没加的情况
- 不再出现超卖问题
性能方面:
- QPS 从原来的 800 提升到 2500
- P99 延迟从 300ms 降到 80ms
- 数据库连接数稳定在合理范围
监控数据:
# 查询数据一致性
mysql> SELECT COUNT(*) as inconsistent_orders
-> FROM orders o
-> LEFT JOIN stock_log sl ON o.saga_id = sl.saga_id
-> LEFT JOIN points_log pl ON o.saga_id = pl.saga_id
-> WHERE o.status = 'COMPLETED'
-> AND (sl.status IS NULL OR pl.status IS NULL OR sl.status != 'COMPLETED' OR pl.status != 'COMPLETED');
+--------------------+
| inconsistent_orders |
+--------------------+
| 0 |
+--------------------+
总结
分布式事务没有银弹,每种方案都有适用场景:
- 2PC:理论正确,但性能和可用性太差,实际不推荐
- TCC:一致性强,但业务侵入严重,实现复杂度高
- SAGA:平衡性最好,推荐优先考虑
- Seata:适合中大型项目,有完善的 SAGA 支持
如果只是简单的跨服务一致性,手写 SAGA 可能更快;如果场景复杂、要求高,用 Seata 更省心。
这次改造最大的收获不是技术选型,而是对分布式一致性的理解加深了:没有绝对的强一致性,只有适合当前业务的一致性方案。
能搞定就行。
可用性说明:本文发布于 2018 年 1 月,距今已超过五年。文中涉及的软件版本、接口、下载地址、命令参数和操作界面可能已经发生变化,部分方案在当前环境下可能失效。请结合官方最新文档核对后再操作,生产环境使用前务必先行验证。
版权声明: 本文首发于 指尖魔法屋-分布式事务折腾手记(https://blog.thinkmoon.cn/post/118-distributed-transaction-2pc-saga-practice/) 转载或引用必须申明原指尖魔法屋来源及源地址!
评论
使用 GitHub 账号登录后即可留言,支持 Markdown。