Appearance
分布式事务与实现
分布式事务是微服务架构中最难啃的骨头——一次业务操作跨越多个服务、多个数据库,如何保证数据一致性? 单机时代一个 @Transactional 搞定,微服务时代同样的操作变成了分布式难题。
一、为什么需要分布式事务
1.1 单机事务 vs 分布式事务
单机事务(一个数据库):
@Transactional
public void createOrder() {
orderMapper.insert(order); // 写订单表
inventoryMapper.deduct(1001); // 扣库存
accountMapper.deduct(100); // 扣余额
}
// 要么全成功,要么全回滚,一个 @Transactional 搞定
分布式事务(微服务架构):
订单服务 ──→ 订单数据库 ──→ INSERT 订单 ✅
库存服务 ──→ 库存数据库 ──→ 扣减库存 ✅
账户服务 ──→ 账户数据库 ──→ 扣减余额 ❌ 失败!
此时订单已创建,库存已扣减,但余额扣款失败——
怎么回滚?@Transactional 管不了跨服务、跨数据库!1.2 典型场景
| 场景 | 涉及服务 | 事务要求 |
|---|---|---|
| 下单 | 订单 + 库存 + 账户 + 优惠券 | 扣库存、扣余额、减优惠券、创建订单,全部成功或全部回滚 |
| 转账 | 账户 A + 账户 B | A 扣钱 + B 加钱,要么都成功,要么都不做 |
| 秒杀 | 库存 + 订单 + 消息队列 | 扣库存成功后,订单创建和消息发送也必须成功 |
| 退款 | 订单 + 账户 + 库存 | 订单状态更新、余额退回、库存恢复 |
二、CAP 理论与 BASE 理论
2.1 CAP 定理
CAP:一个分布式系统最多只能同时满足三项中的两项
C(Consistency,一致性):所有节点同一时间看到相同的数据
A(Availability,可用性):每次请求都能获得非错误的响应
P(Partition Tolerance,分区容错性):网络分区时系统仍能工作
分布式系统 P 是必须的(网络分区无法避免),
所以只能在 C 和 A 之间做取舍:
CP 系统(强一致,牺牲可用性):
→ ZooKeeper、Etcd、Consul
→ 网络分区时,宁可不可用也不返回不一致数据
AP 系统(高可用,牺牲强一致):
→ Eureka、Nacos(AP 模式)
→ 网络分区时,宁可数据不一致也要保证可用2.2 BASE 理论
BASE = 对 CAP 中 AP 方案的补充
BA(Basically Available,基本可用):
系统出现故障时,允许损失部分可用性,但保证核心功能可用
例如:下单页面正常,但优惠券查询暂时不可用
S(Soft State,软状态):
允许系统中的数据存在中间状态,该状态不影响系统整体可用性
例如:订单状态为"支付中"(中间状态),最终会变为"已支付"或"已取消"
E(Eventually Consistent,最终一致性):
不要求数据实时一致,但保证最终会达到一致
例如:支付成功后,账户余额可能 1 秒后才更新,但最终一定正确
BASE 的核心思想:放弃强一致性,追求最终一致性CAP → BASE 的演进:
CAP 告诉我们:分布式系统做不到完美的强一致 + 高可用
BASE 告诉我们:那就退一步,追求最终一致
大部分互联网业务不需要强一致——
用户支付成功后,订单状态晚 1 秒变成"已支付",完全可以接受三、分布式事务四大方案
3.1 方案全景图
| 方案 | 一致性 | 性能 | 复杂度 | 适用场景 |
|---|---|---|---|---|
| 2PC(XA) | 强一致 | 低 | 低 | 传统单体拆微服务,跨库事务 |
| TCC | 强一致 | 中 | 高 | 资金、交易等核心业务 |
| Saga | 最终一致 | 高 | 中 | 长流程、非资金类业务 |
| AT(Seata) | 最终一致 | 中 | 低 | 对业务侵入小的场景 |
| 可靠消息 | 最终一致 | 高 | 中 | 异步解耦、通知类场景 |
| 最大努力通知 | 最终一致 | 高 | 低 | 外部系统回调、通知 |
3.2 方案选择决策树
你的业务场景 → 推荐方案
─────────────────────────────────────────────────────────────
强一致性、资金交易(转账、支付) → TCC
强一致性、简单跨库事务(单体拆微服务) → 2PC(XA)
长流程、非资金类(订单创建→审核→发货) → Saga
无侵入、快速接入、非强一致 → AT(Seata)
异步解耦、通知类(短信、邮件) → 可靠消息
外部系统回调(支付回调、物流通知) → 最大努力通知四、2PC(两阶段提交)
4.1 原理
2PC = 协调者 + 参与者,分两个阶段
阶段一:投票(Prepare)
协调者 → 所有参与者:能执行吗?
参与者 → 协调者:能(YES)/ 不能(NO)
阶段二:提交(Commit)
如果所有参与者都回复 YES:
协调者 → 所有参与者:提交!
如果任意参与者回复 NO:
协调者 → 所有参与者:回滚!正常流程:
协调者 参与者 A 参与者 B
│ │ │
│──── Prepare ─────────────────→│ │
│──── Prepare ──────────────────────────────────→│
│ │ │
│←─── YES(锁定资源)──────────│ │
│←─── YES(锁定资源)──────────────────────────│
│ │ │
│──── Commit ──────────────────→│ │
│──── Commit ──────────────────────────────────→│
│ │ │
│←─── ACK ─────────────────────│ │
│←─── ACK ─────────────────────────────────────│
│ │ │
✅ 事务提交成功4.2 2PC 的致命问题
❌ 同步阻塞:Prepare 阶段,参与者锁定资源,必须等协调者通知
如果协调者挂了,参与者一直阻塞,资源无法释放
❌ 单点故障:协调者挂了,整个系统瘫痪
参与者不知道应该提交还是回滚
❌ 数据不一致:阶段二,协调者发了 Commit 后挂了
有的参与者收到 Commit 提交了,有的没收到,数据不一致
❌ 性能差:所有参与者都要锁定资源等协调者
一次事务 = 2 次网络往返 × N 个参与者4.3 Seata XA 模式(2PC 实现)
yaml
# application.yml
seata:
tx-service-group: my_tx_group
service:
vgroup-mapping:
my_tx_group: default
data-source-proxy-mode: XA # 使用 XA 模式java
@Service
public class OrderService {
@Autowired
private OrderMapper orderMapper;
@Autowired
private InventoryClient inventoryClient;
@Autowired
private AccountClient accountClient;
@GlobalTransactional // Seata 分布式事务注解
public void createOrder(Order order) {
// ① 创建订单(本地事务)
orderMapper.insert(order);
// ② 扣减库存(远程调用,Seata 自动代理)
inventoryClient.deduct(order.getProductId(), order.getQuantity());
// ③ 扣减余额(远程调用)
accountClient.deduct(order.getUserId(), order.getAmount());
// 任意一步失败,Seata 自动回滚所有操作
}
}XA 模式适用场景: 传统单体应用拆分为微服务,数据库支持 XA 协议(MySQL 5.7+、PostgreSQL、Oracle),对一致性要求高但并发量不大。
五、TCC(Try-Confirm-Cancel)
5.1 原理
TCC = 把一次业务操作拆成三个步骤
Try(尝试):预留资源,检查 + 锁定
→ 冻结库存 100 件
→ 冻结余额 100 元
Confirm(确认):确认提交,使用预留的资源
→ 扣减冻结的库存
→ 扣减冻结的余额
Cancel(取消):取消操作,释放预留的资源
→ 恢复冻结的库存
→ 恢复冻结的余额TCC 完整流程:
业务方(TM) 库存服务 账户服务
│ │ │
│──── Try ──────────────────→│ │
│ freezeStock(1001, 10) │ 冻结 10 件库存 │
│←─── OK ───────────────────│ │
│ │ │
│──── Try ────────────────────────────────────────────→│
│ freezeBalance(1001, 100) │ │ 冻结 100 元
│←─── OK ─────────────────────────────────────────────│
│ │ │
│ ⚠️ 业务校验失败,需要回滚 │ │
│ │ │
│──── Cancel ───────────────→│ │
│ unfreezeStock(1001, 10) │ 恢复 10 件库存 │
│←─── OK ───────────────────│ │
│ │ │
│──── Cancel ──────────────────────────────────────────→│
│ unfreezeBalance(1001,100) │ │ 恢复 100 元
│←─── OK ─────────────────────────────────────────────│5.2 TCC 实现
java
/**
* 库存服务——TCC 接口
*/
public interface InventoryTccService {
/**
* Try:冻结库存
* @param businessId 业务流水号,用于幂等
* @param productId 商品 ID
* @param count 冻结数量
*/
@Transactional
boolean tryFreezeStock(String businessId, Long productId, Integer count);
/**
* Confirm:确认扣减(使用冻结的库存)
*/
@Transactional
boolean confirmFreezeStock(String businessId, Long productId, Integer count);
/**
* Cancel:取消冻结(恢复库存)
*/
@Transactional
boolean cancelFreezeStock(String businessId, Long productId, Integer count);
}java
@Service
public class InventoryTccServiceImpl implements InventoryTccService {
@Autowired
private InventoryMapper inventoryMapper;
@Autowired
private FreezeLogMapper freezeLogMapper;
@Override
public boolean tryFreezeStock(String businessId, Long productId, Integer count) {
// ① 幂等校验:防止重复 Try
if (freezeLogMapper.exists(businessId)) {
return true; // 已处理过,直接返回成功
}
// ② 扣减可用库存,增加冻结库存
// UPDATE inventory SET available = available - ?, frozen = frozen + ?
// WHERE product_id = ? AND available >= ?
int rows = inventoryMapper.freeze(productId, count);
if (rows == 0) {
return false; // 库存不足
}
// ③ 记录冻结流水(用于 Confirm 和 Cancel 的幂等控制)
FreezeLog log = new FreezeLog();
log.setBusinessId(businessId);
log.setProductId(productId);
log.setCount(count);
log.setStatus("FROZEN");
freezeLogMapper.insert(log);
return true;
}
@Override
public boolean confirmFreezeStock(String businessId, Long productId, Integer count) {
// ① 幂等校验
FreezeLog log = freezeLogMapper.getByBusinessId(businessId);
if (log == null) return false;
if ("CONFIRMED".equals(log.getStatus())) return true; // 已确认
// ② 扣减冻结库存
// UPDATE inventory SET frozen = frozen - ?, total = total - ?
// WHERE product_id = ?
inventoryMapper.confirmDeduct(productId, count);
// ③ 更新流水状态
freezeLogMapper.updateStatus(businessId, "CONFIRMED");
return true;
}
@Override
public boolean cancelFreezeStock(String businessId, Long productId, Integer count) {
// ① 幂等校验
FreezeLog log = freezeLogMapper.getByBusinessId(businessId);
if (log == null) {
// 没有 Try 记录,说明 Try 根本没执行,直接返回成功
return true;
}
if ("CANCELLED".equals(log.getStatus())) return true; // 已取消
// ② 恢复库存
// UPDATE inventory SET available = available + ?, frozen = frozen - ?
// WHERE product_id = ?
inventoryMapper.unfreeze(productId, count);
// ③ 更新流水状态
freezeLogMapper.updateStatus(businessId, "CANCELLED");
return true;
}
}java
/**
* 订单服务——TCC 事务编排(TM)
*/
@Service
public class OrderTccService {
@Autowired
private InventoryTccClient inventoryTccClient;
@Autowired
private AccountTccClient accountTccClient;
@Autowired
private OrderService orderService;
public void createOrder(Order order) {
String businessId = order.getOrderId(); // 使用订单号作为业务流水号
// ① Try 阶段
boolean stockFrozen = inventoryTccClient.tryFreezeStock(
businessId, order.getProductId(), order.getQuantity());
if (!stockFrozen) {
throw new RuntimeException("库存不足");
}
boolean balanceFrozen = accountTccClient.tryFreezeBalance(
businessId, order.getUserId(), order.getAmount());
if (!balanceFrozen) {
// 余额不足,回滚库存
inventoryTccClient.cancelFreezeStock(
businessId, order.getProductId(), order.getQuantity());
throw new RuntimeException("余额不足");
}
try {
// ② 创建订单
orderService.createOrder(order);
// ③ Confirm 阶段:全部 Try 成功,确认提交
inventoryTccClient.confirmFreezeStock(
businessId, order.getProductId(), order.getQuantity());
accountTccClient.confirmFreezeBalance(
businessId, order.getUserId(), order.getAmount());
} catch (Exception e) {
// ④ Cancel 阶段:任意步骤失败,回滚所有
inventoryTccClient.cancelFreezeStock(
businessId, order.getProductId(), order.getQuantity());
accountTccClient.cancelFreezeBalance(
businessId, order.getUserId(), order.getAmount());
throw new RuntimeException("下单失败", e);
}
}
}5.3 TCC 的空回滚、幂等、悬挂
分布式事务的三大经典问题:
① 空回滚(Empty Rollback):
Try 超时,TM 发起 Cancel,但 Try 的请求还没到达 RM
→ RM 收到 Cancel,但没有对应的 Try 记录
→ 解决:Cancel 阶段发现没有 Try 记录时,直接返回成功(允许空回滚)
② 幂等(Idempotent):
Try/Confirm/Cancel 可能被多次调用(网络重试、TM 重试)
→ 解决:使用 businessId 做唯一键,已处理则直接返回成功
③ 悬挂(Suspension):
Cancel 先到达且执行完毕,之后 Try 才到达
→ Try 执行了但 Cancel 已经做过了,资源被错误冻结
→ 解决:Cancel 执行后标记 businessId 为已取消,Try 时检查标记,拒绝执行java
// 悬挂处理:Cancel 后标记,Try 时检查
public boolean tryFreezeStock(String businessId, Long productId, Integer count) {
// 检查是否已被 Cancel(防止悬挂)
if (freezeLogMapper.isCancelled(businessId)) {
return false; // 已被取消,拒绝 Try
}
// ... 正常 Try 逻辑
}
public boolean cancelFreezeStock(String businessId, Long productId, Integer count) {
// 空回滚处理
FreezeLog log = freezeLogMapper.getByBusinessId(businessId);
if (log == null) {
// 记录一条已取消的日志,防止后续 Try 悬挂
freezeLogMapper.insertCancelled(businessId);
return true; // 空回滚,允许成功
}
// ... 正常 Cancel 逻辑
}5.4 TCC 优缺点
✅ 优点:
- 强一致性:Confirm 全部成功才算成功,任一失败就全部 Cancel
- 性能优于 2PC:Try 阶段只锁定资源,不阻塞等待
- 业务可控:开发者自定义 Try/Confirm/Cancel 逻辑
❌ 缺点:
- 代码侵入大:每个业务都要写 Try/Confirm/Cancel 三个方法
- 开发成本高:空回滚、幂等、悬挂都要处理
- 不适合长流程:Confirm 阶段失败需要人工介入六、Saga 模式
6.1 原理
Saga = 把一个大事务拆成多个有序的小事务,每个小事务都有对应的补偿操作
正向:T1 → T2 → T3 → ... → Tn
补偿:C1 ← C2 ← C3 ← ... ← Cn
如果 T3 失败,依次执行 C2 → C1 回滚6.2 两种实现方式
① 编排式(Choreography,事件驱动):
每个服务做完自己的事,发布事件,下一个服务监听事件
订单服务 → 创建订单 → 发布"订单已创建"事件
库存服务 → 监听事件 → 扣库存 → 发布"库存已扣减"事件
账户服务 → 监听事件 → 扣余额 → 结束
优点:松耦合,适合简单流程
缺点:流程分散在各服务中,难以追踪
② 编排式(Orchestration,命令驱动):
一个 Saga 编排器统一调度所有步骤
Saga 编排器:
① 调用订单服务创建订单
② 调用库存服务扣库存
③ 调用账户服务扣余额
④ 失败时按相反顺序调用补偿操作
优点:流程集中管理,清晰可追踪
缺点:编排器成为新的耦合点6.3 Seata Saga 实现
java
// Saga 状态机定义(JSON 格式)
// 定义流程:创建订单 → 扣库存 → 扣余额json
{
"Name": "createOrderSaga",
"StartState": "CreateOrder",
"States": {
"CreateOrder": {
"Type": "ServiceTask",
"ServiceName": "orderService.createOrder",
"CompensateState": "CancelOrder",
"Next": "DeductStock"
},
"CancelOrder": {
"Type": "ServiceTask",
"ServiceName": "orderService.cancelOrder",
"End": true
},
"DeductStock": {
"Type": "ServiceTask",
"ServiceName": "inventoryService.deductStock",
"CompensateState": "CompensateStock",
"Next": "DeductBalance"
},
"CompensateStock": {
"Type": "ServiceTask",
"ServiceName": "inventoryService.increaseStock",
"Next": "CancelOrder"
},
"DeductBalance": {
"Type": "ServiceTask",
"ServiceName": "accountService.deductBalance",
"CompensateState": "CompensateBalance",
"End": true
},
"CompensateBalance": {
"Type": "ServiceTask",
"ServiceName": "accountService.increaseBalance",
"Next": "CompensateStock"
}
}
}java
@Service
public class OrderSagaService {
@Autowired
private SagaStateMachineEngine sagaEngine;
public void createOrder(Order order) {
// 启动 Saga 状态机,Seata 自动编排
sagaEngine.start("createOrderSaga", order);
}
}6.4 Saga 优缺点
✅ 优点:
- 高性能:每个小事务独立提交,不长时间锁定资源
- 长流程友好:适合多步骤、跨服务的长事务
- 松耦合:服务间通过事件或编排器通信
❌ 缺点:
- 最终一致性:不保证实时一致,中间状态可见
- 补偿复杂:补偿操作需要业务逻辑配合(如退款可能有手续费)
- 隔离性弱:事务 A 扣了库存但还没提交,事务 B 可能读到中间状态七、Seata AT 模式
7.1 原理
AT 模式 = 自动挡分布式事务,对业务代码零侵入
核心机制:两阶段提交 + 自动生成回滚 SQL
阶段一:
1. 执行原始 SQL(INSERT / UPDATE / DELETE)
2. 解析 SQL,获取前镜像(before image)
3. 执行业务 SQL
4. 获取后镜像(after image)
5. 将前后镜像 + 行锁插入 undo_log 表
6. 提交本地事务(释放数据库连接)
阶段二(提交):
1. 删除 undo_log 记录(异步,批量)
阶段二(回滚):
1. 从 undo_log 读取前镜像
2. 校验后镜像是否匹配(防止脏写)
3. 用前镜像生成反向 SQL 并执行
4. 删除 undo_log 记录AT 模式的全局锁机制:
事务 A:UPDATE product SET stock = stock - 10 WHERE id = 1001
→ Seata 自动加全局锁:lock:product:1001
事务 B:UPDATE product SET stock = stock - 5 WHERE id = 1001
→ 发现全局锁被占用 → 等待或超时
事务 A 提交 → 释放全局锁 → 事务 B 获取锁 → 执行7.2 AT 模式快速上手
xml
<!-- Seata 1.8.0 + Spring Boot 3.x -->
<dependency>
<groupId>io.seata</groupId>
<artifactId>seata-spring-boot-starter</artifactId>
<version>2.1.0</version>
</dependency>yaml
# application.yml
seata:
enabled: true
tx-service-group: default_tx_group
service:
vgroup-mapping:
default_tx_group: default
registry:
type: nacos
nacos:
server-addr: 127.0.0.1:8848
group: SEATA_GROUP
config:
type: nacos
nacos:
server-addr: 127.0.0.1:8848
group: SEATA_GROUP
data-source-proxy-mode: AT # AT 模式sql
-- 每个业务的数据库都需要创建 undo_log 表
CREATE TABLE `undo_log` (
`id` BIGINT(20) NOT NULL AUTO_INCREMENT,
`branch_id` BIGINT(20) NOT NULL,
`xid` VARCHAR(100) NOT NULL,
`context` VARCHAR(128) NOT NULL,
`rollback_info` LONGBLOB NOT NULL,
`log_status` INT(11) NOT NULL,
`log_created` DATETIME NOT NULL,
`log_modified` DATETIME NOT NULL,
PRIMARY KEY (`id`),
KEY `idx_unionkey` (`xid`, `branch_id`)
) ENGINE = InnoDB DEFAULT CHARSET = utf8mb4;java
@Service
public class OrderService {
@Autowired
private OrderMapper orderMapper;
@Autowired
private InventoryClient inventoryClient;
@Autowired
private AccountClient accountClient;
@GlobalTransactional(name = "createOrder", timeoutMills = 300000)
public void createOrder(Order order) {
// 业务代码完全不需要修改,和单机事务写法一模一样
// Seata 自动代理数据源,自动记录 undo_log,自动回滚
orderMapper.insert(order); // 本地事务
inventoryClient.deduct(order.getProductId(), order.getQuantity()); // 远程调用
accountClient.deduct(order.getUserId(), order.getAmount()); // 远程调用
// 任意步骤失败,Seata 自动回滚所有操作
}
}7.3 AT 模式优缺点
✅ 优点:
- 零侵入:业务代码不需要任何修改,和单机事务一样写
- 自动回滚:Seata 自动生成反向 SQL,不需要手写补偿
- 性能较好:阶段一就提交本地事务,不长期占用连接
❌ 缺点:
- 依赖数据库:undo_log 表是必须的,不支持非关系型数据库
- 全局锁:依赖 Seata Server 的全局锁,有单点风险
- 隔离性:默认读未提交,可能读到脏数据
- 仅支持 ACID 数据库:MySQL、PostgreSQL、Oracle八、可靠消息最终一致性
8.1 原理
可靠消息 = 本地事务 + 消息表 + 消息队列 + 重试机制
① 本地事务:执行业务操作 + 往"消息表"插入一条消息(同一个事务)
② 定时任务:扫描消息表,将"待发送"的消息发送到 MQ
③ 消费者:消费消息,执行业务逻辑
④ 确认删除:消费者执行成功后,回调确认,消息表标记为"已发送"8.2 实现
sql
-- 消息表(本地事务表)
CREATE TABLE `tx_message` (
`id` BIGINT AUTO_INCREMENT PRIMARY KEY,
`message_id` VARCHAR(64) NOT NULL, -- 消息 ID(幂等键)
`topic` VARCHAR(64) NOT NULL, -- MQ Topic
`tag` VARCHAR(64) DEFAULT NULL, -- MQ Tag
`body` TEXT NOT NULL, -- 消息体 JSON
`status` VARCHAR(16) NOT NULL DEFAULT 'PENDING', -- PENDING/SENT/CONFIRMED
`retry_count` INT DEFAULT 0, -- 重试次数
`max_retry` INT DEFAULT 10, -- 最大重试次数
`next_retry` DATETIME NOT NULL, -- 下次重试时间
`created_at` DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
UNIQUE KEY `uk_message_id` (`message_id`)
);java
@Service
public class OrderService {
@Autowired
private OrderMapper orderMapper;
@Autowired
private TxMessageMapper txMessageMapper;
@Autowired
private RocketMQTemplate rocketMQTemplate;
/**
* 下单:本地事务 + 消息表
*/
@Transactional
public void createOrder(Order order) {
// ① 业务操作:创建订单
orderMapper.insert(order);
// ② 同一事务中,插入消息表
TxMessage message = new TxMessage();
message.setMessageId(order.getOrderId());
message.setTopic("order-topic");
message.setTag("order-created");
message.setBody(JSON.toJSONString(order));
message.setStatus("PENDING");
message.setNextRetry(LocalDateTime.now());
txMessageMapper.insert(message);
// ③ 事务提交后,发送消息(这里简化,实际由定时任务发送)
// 如果发送失败,定时任务会重试
}
}java
/**
* 定时任务:扫描消息表,发送消息
*/
@Component
public class TxMessageSender {
@Autowired
private TxMessageMapper txMessageMapper;
@Autowired
private RocketMQTemplate rocketMQTemplate;
@Scheduled(fixedDelay = 5000) // 每 5 秒扫描一次
public void sendPendingMessages() {
// 查询待发送的消息(状态=PENDING,且到达重试时间)
List<TxMessage> messages = txMessageMapper.selectPendingMessages(
"PENDING", LocalDateTime.now(), 100);
for (TxMessage msg : messages) {
try {
// 发送到 RocketMQ
rocketMQTemplate.syncSend(
msg.getTopic() + ":" + msg.getTag(),
msg.getBody());
// 更新状态为已发送
txMessageMapper.updateStatus(msg.getMessageId(), "SENT");
} catch (Exception e) {
// 发送失败,更新重试信息
int retries = msg.getRetryCount() + 1;
if (retries > msg.getMaxRetry()) {
// 超过最大重试次数,标记为失败,人工介入
txMessageMapper.updateStatus(msg.getMessageId(), "FAILED");
} else {
// 指数退避:下次重试 = 当前时间 + 2^retries 秒
txMessageMapper.updateRetry(
msg.getMessageId(), retries,
LocalDateTime.now().plusSeconds((long) Math.pow(2, retries)));
}
}
}
}
}8.3 RocketMQ 事务消息(更优方案)
RocketMQ 事务消息 = 普通消息 + 事务状态回查
① 发送半消息(Half Message)→ MQ 保存但不可见
② 执行本地事务
③ 本地事务成功 → 提交半消息 → MQ 变为可见
本地事务失败 → 回滚半消息 → MQ 删除
④ 如果 MQ 没收到确认(超时),MQ 回查本地事务状态java
@Service
public class OrderService {
@Autowired
private RocketMQTemplate rocketMQTemplate;
public void createOrder(Order order) {
// RocketMQ 事务消息,一步到位
rocketMQTemplate.sendMessageInTransaction(
"order-topic:order-created", // topic:tag
MessageBuilder.withPayload(order).build(), // 消息体
order.getOrderId() // 业务参数
);
}
/**
* 本地事务执行回调
*/
@RocketMQTransactionListener
public class OrderTransactionListener implements RocketMQLocalTransactionListener {
@Override
public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
String orderId = (String) arg;
try {
// 执行本地事务:创建订单
orderMapper.insert(parseOrder(msg));
return RocketMQLocalTransactionState.COMMIT; // 提交消息
} catch (Exception e) {
return RocketMQLocalTransactionState.ROLLBACK; // 回滚消息
}
}
@Override
public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
String orderId = (String) msg.getHeaders().get("orderId");
// 回查:查询本地事务是否执行成功
Order order = orderMapper.selectById(orderId);
return order != null
? RocketMQLocalTransactionState.COMMIT
: RocketMQLocalTransactionState.ROLLBACK;
}
}
}九、最大努力通知
9.1 原理
最大努力通知 = 通知方尽最大努力通知,不保证一定成功
① 业务操作完成后,通知外部系统
② 如果通知失败,按一定策略重试(递增间隔)
③ 重试 N 次后仍失败,记录异常,人工介入
典型场景:支付回调通知java
@Component
public class PaymentNotifyService {
@Autowired
private RestTemplate restTemplate;
private static final int MAX_RETRIES = 5;
private static final long[] RETRY_INTERVALS = {1, 5, 15, 30, 60}; // 分钟
public void notifyMerchant(String orderId, String notifyUrl, String notifyBody) {
for (int i = 0; i < MAX_RETRIES; i++) {
try {
ResponseEntity<String> response = restTemplate.postForEntity(
notifyUrl, notifyBody, String.class);
if (response.getStatusCode().is2xxSuccessful()
&& "SUCCESS".equals(response.getBody())) {
// 通知成功
return;
}
} catch (Exception e) {
log.warn("通知商户失败,第 {} 次重试,orderId={}", i + 1, orderId);
}
// 等待后重试
try {
TimeUnit.MINUTES.sleep(RETRY_INTERVALS[i]);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
break;
}
}
// 超过最大重试次数,记录异常,人工处理
log.error("通知商户失败,已超过最大重试次数,orderId={}", orderId);
// 发送告警
alertService.send("支付通知失败", orderId);
}
}十、六大方案全面对比
| 维度 | 2PC(XA) | TCC | Saga | AT(Seata) | 可靠消息 | 最大努力通知 |
|---|---|---|---|---|---|---|
| 一致性 | 强一致 | 强一致 | 最终一致 | 最终一致 | 最终一致 | 最终一致 |
| 性能 | ⭐⭐ | ⭐⭐⭐ | ⭐⭐⭐⭐ | ⭐⭐⭐ | ⭐⭐⭐⭐ | ⭐⭐⭐⭐⭐ |
| 复杂度 | ⭐⭐ | ⭐⭐⭐⭐⭐ | ⭐⭐⭐ | ⭐⭐ | ⭐⭐⭐ | ⭐⭐ |
| 代码侵入 | 低 | 高 | 中 | 零 | 低 | 低 |
| 回滚方式 | 自动 | 手动(Cancel) | 补偿 | 自动(undo_log) | 补偿 | 重试 |
| 锁资源时间 | 长 | 短 | 短 | 短 | 无 | 无 |
| 适用场景 | 跨库事务 | 资金交易 | 长流程 | 快速接入 | 异步解耦 | 外部回调 |
| 是否依赖数据库 | 是 | 否 | 否 | 是(MySQL) | 否 | 否 |
十一、生产级最佳实践
11.1 幂等设计
分布式事务的每个参与者都必须支持幂等:
1. 使用全局唯一 ID(xid + branchId)作为幂等键
2. 数据库唯一索引 + 状态机
3. 所有 Confirm / Cancel / 补偿操作都支持重复调用11.2 监控与告警
分布式事务需要完善的监控体系:
1. 事务成功率:每分钟提交/回滚的比例
2. 事务耗时:P50 / P99 / P999
3. undo_log 堆积:未清理的 undo_log 数量
4. 重试次数:消息重试次数分布
5. 异常告警:事务失败、超时、回滚失败11.3 避免长事务
分布式事务的黄金法则:事务越短越好
1. 不要在事务中调用外部 API(RPC 除外)
2. 不要在事务中发邮件、发短信
3. 不要在事务中操作文件、生成 PDF
4. 超过 1 秒的事务考虑异步化
5. 超过 10 个参与者的考虑用 Saga 替代 TCC11.4 定时任务补偿
java
/**
* 分布式事务补偿任务:处理未完成的异常事务
*/
@Component
public class TransactionCompensationJob {
@Scheduled(fixedDelay = 60000) // 每分钟执行
public void compensate() {
// ① 处理长时间未确认的 TCC Try 记录(超过 30 分钟)
List<FreezeLog> timeoutLogs = freezeLogMapper.selectTimeoutLogs(
"FROZEN", LocalDateTime.now().minusMinutes(30));
for (FreezeLog log : timeoutLogs) {
// 自动 Cancel 释放资源
inventoryService.cancelFreezeStock(
log.getBusinessId(), log.getProductId(), log.getCount());
}
// ② 处理长时间未清理的 undo_log(超过 7 天)
undoLogMapper.deleteExpired(LocalDateTime.now().minusDays(7));
// ③ 处理失败的消息(超过 24 小时)
List<TxMessage> failedMessages = txMessageMapper.selectFailedMessages(
"FAILED", LocalDateTime.now().minusHours(24));
for (TxMessage msg : failedMessages) {
// 发送告警,人工处理
alertService.send("消息发送失败", msg.getMessageId());
}
}
}十二、面试要点
Q1:分布式事务有哪些实现方案?
六大方案:2PC(XA)强一致、TCC(Try-Confirm-Cancel)强一致、Saga 补偿模式、Seata AT 模式(自动回滚)、可靠消息最终一致性、最大努力通知。按一致性要求从强到弱选择。
Q2:TCC 和 AT 模式有什么区别?
TCC 需要手写 Try/Confirm/Cancel 三个方法,代码侵入大但灵活性高,适合资金类;AT 模式零侵入,Seata 自动生成反向 SQL,但依赖数据库,适合快速接入。
Q3:TCC 的空回滚、幂等、悬挂是什么?
空回滚:Cancel 时 Try 还没执行;幂等:同一操作重复执行;悬挂:Cancel 先到,Try 后到。都需要通过 businessId + 状态机处理。
Q4:Seata 的 AT 模式隔离级别是什么?
默认读未提交。如果业务要求读已提交,Seata 提供了 @GlobalLock 注解 + SELECT FOR UPDATE 来保证。
Q5:RocketMQ 事务消息和本地消息表怎么选?
RocketMQ 事务消息更优雅,不需要额外维护消息表,但绑定 RocketMQ;本地消息表方案通用,任何 MQ 都适用,但需要定时任务扫描。
Q6:什么场景不适合用分布式事务?
1) 非核心业务,允许少量不一致(如统计、日志);2) 超高并发场景(秒杀),分布式事务会成为瓶颈;3) 跨公司/跨系统,无法控制对方数据库。这些场景用异步 + 对账 + 人工兜底更实际。
