分布式事务(Seata / 2PC / TCC / 消息最终一致性)
分布式事务产生的背景
在单体应用中,所有操作都在同一个数据库内完成,依赖数据库本身的 ACID 事务即可保证数据一致性。但在微服务架构下,数据被分散在多个独立的数据库中,业务操作往往需要跨多个服务、多个数据库完成,本地事务已无法满足需求。
单体应用: 一个服务 → 一个数据库 → 本地事务
微服务架构: 服务 A → 数据库 A ? 如何保证一致性
服务 B → 数据库 B ! 分布式事务
服务 C → 数据库 C分布式事务产生的三大场景:
| 场景 | 说明 |
|---|---|
| 跨库 | 同一服务内操作多个数据库 |
| 跨服务 | 一个业务调用链涉及多个微服务 |
| 跨资源 | 同时操作数据库、消息队列、缓存等异构资源 |
CAP 理论与 BASE 理论
CAP 理论
一个分布式系统最多只能同时满足以下三个特性中的两个:
| 特性 | 说明 |
|---|---|
| C(Consistency)一致性 | 所有节点在同一时刻看到的数据相同 |
| A(Availability)可用性 | 每个请求都能获得非错误的响应 |
| P(Partition Tolerance)分区容错性 | 系统允许网络分区,仍能正常运行 |
C(一致性)
/ \
/ \
/ \
/ \
CP AP
\ /
\ /
\ /
\ /
P(分区容错性)
在实际分布式系统中,P 是必须的,因此本质是在 C 和 A 之间做权衡:
- CP 系统:Zookeeper、Etcd(强一致,牺牲部分可用性)
- AP 系统:Eureka、Cassandra(高可用,牺牲强一致性)BASE 理论
BASE 是对 CAP 中 AP 方案的延伸,核心思想是放弃强一致性,追求最终一致性。
| 要素 | 说明 |
|---|---|
| BA(Basically Available)基本可用 | 系统允许降级,保证核心功能可用 |
| S(Soft State)软状态 | 允许数据存在中间状态 |
| E(Eventually Consistent)最终一致性 | 经过一段时间后,数据最终达到一致 |
传统方案:两阶段提交(2PC)
流程
2PC 引入协调者(Coordinator) 来管理多个参与者(Participant) 的事务。
阶段一:准备阶段(Prepare / Voting)
协调者 参与者
│─── prepare ─────────→│
│ │ 执行本地事务(不提交)
│ │ 写 undo / redo log
│←────── yes / no ─────│
│ │
阶段二:提交阶段(Commit / Abort)
│─── commit / abort ──→│
│ │ 提交 or 回滚本地事务
│←────── ack ──────────│
│ │准备阶段:
- 协调者向所有参与者发送 prepare 请求
- 参与者执行本地事务,写 undo/redo 日志
- 参与者返回 Yes(就绪)或 No(失败)
提交阶段:
- 所有参与者都返回 Yes → 协调者发送 commit
- 任一参与者返回 No → 协调者发送 rollback
存在的问题
| 问题 | 说明 |
|---|---|
| 同步阻塞 | 准备阶段参与者持有资源锁,阻塞其他事务 |
| 协调者单点 | 协调者宕机,所有参与者无法决策 |
| 数据不一致 | 第二阶段网络异常,部分参与者收到 commit,部分未收到 |
| 脑裂 | 协调者发送 commit 后宕机,部分参与者提交成功,部分未知 |
三阶段提交(3PC)
3PC 在 2PC 基础上引入了超时机制和准备阶段细分,降低了阻塞范围:
阶段一:CanCommit(询问能否提交)
协调者 → 参与者:可以提交吗?
参与者 → 协调者:可以 / 不可以
阶段二:PreCommit(预提交)
协调者 → 参与者:准备提交
参与者 → 参与者:写 undo/redo 日志,返回 ACK
阶段三:DoCommit(正式提交)
协调者 → 参与者:提交 / 中止
参与者超时未收到指令 → 自动提交(相比 2PC 的改进)改进点:
- 引入超时机制,参与者超时后自动提交
- 增加 CanCommit 阶段,提前发现不能提交的节点
仍存在的问题:
- 网络分区时,参与者自动提交可能导致数据不一致
- 性能提升有限,实际应用较少
TCC 模式(Try-Confirm-Cancel)
TCC 是一种业务层的分布式事务方案,将每个分支事务分为三个操作。
流程
Try 阶段:
预留业务资源(冻结库存、冻结账户金额)
eg. 扣减库存时,将 quantity 从"可用库存"移到"冻结库存"
Confirm 阶段:
确认执行业务(将冻结资源真正扣减)
eg. 将"冻结库存"扣减,标记为已完成
Cancel 阶段:
取消执行,释放预留资源
eg. 将"冻结库存"释放回"可用库存"全局事务管理器 服务 A(库存) 服务 B(账户) 服务 C(积分)
│── Try ─────────────→│ │ │
│ │ 冻结库存 │ │
│── Try ──────────────────────────────────→│ │
│ │ │ 冻结金额 │
│── Try ────────────────────────────────────────────────────────→│
│ │ │ │ 冻结积分
│←─── 全部成功 ───────│←────────────────────│←────────────────────│
│ │ │ │
│── Confirm ─────────→│ │ │
│ │ 扣减冻结库存 │ │
│── Confirm ───────────────────────────────→│ │
│ │ │ 扣减冻结金额 │
│── Confirm ───────────────────────────────────────────────────→│
│ │ │ │ 增加积分业务侵入性分析
| 方面 | 说明 |
|---|---|
| 接口定义 | 每个业务接口需要拆分为 Try / Confirm / Cancel 三个方法 |
| 资源预留 | 业务表需设计冻结/预占状态字段 |
| 逻辑复杂 | 每个分支事务需考虑正向和逆向操作 |
| 代码量 | 相比 AT 模式,TCC 增加约 2-3 倍的业务代码 |
空回滚、幂等、悬挂问题
空回滚:Try 未执行,但 Cancel 被调用(如 Try 超时,协调者直接发 Cancel)。
// 空回滚处理:记录事务状态,Cancel 时检查 Try 是否已执行
public void cancel(BusinessActionContext ctx) {
// 检查 Try 阶段是否执行过
if (tryRecordNotExist(ctx.getXid(), ctx.getBranchId())) {
// Try 未执行,空回滚,直接返回成功
return;
}
// 正常回滚逻辑
doCancel(ctx);
}幂等:Confirm 或 Cancel 可能被重复调用。
// 幂等处理:使用事务状态表判断
public void confirm(BusinessActionContext ctx) {
String status = getTransactionStatus(ctx.getXid(), ctx.getBranchId());
if ("CONFIRMED".equals(status)) {
return; // 已确认,跳过
}
doConfirm(ctx);
updateStatus(ctx.getXid(), ctx.getBranchId(), "CONFIRMED");
}悬挂:Cancel 先于 Try 执行(网络延迟导致 Try 包晚于 Cancel 到达)。
// 悬挂处理:Cancel 中记录"已取消",Try 时检查是否已被取消
public void tryMethod(BusinessActionContext ctx) {
if (isCancelled(ctx.getXid(), ctx.getBranchId())) {
// 已收到 Cancel,不再执行 Try
return;
}
doTry(ctx);
}SAGA 模式
SAGA 将长事务拆分为一组本地事务,每个本地事务都有对应的补偿事务。
编排型 vs 协同型
编排型(Choreography):每个服务完成本地事务后,发布事件触发下一个服务。
订单服务 ──(创建订单)──→ 库存服务 ──(扣减成功)──→ 账户服务
│ │ │
│ 事件驱动,无中心协调者 │ │
└──── 失败时反向补偿 ────┘←────────────────────────┘- 优点:架构简单,无单点
- 缺点:业务流程耦合在事件中,复杂流程难以追踪
协同型(Orchestration):引入协调者(Orchestrator)统一管理流程。
Orchestrator(流程协调者)
/ | \
/ | \
订单服务 库存服务 账户服务
┌───────────────────────────────────┐
│ 1. 创建订单 │
│ 2. 扣减库存 │
│ 3. 扣减账户 │
│ 4. 如果第 3 步失败,补偿第 2 步 │
│ 5. 补偿第 1 步 │
└───────────────────────────────────┘- 优点:流程清晰,职责集中
- 缺点:协调者可能成为瓶颈和单点
补偿事务设计
补偿事务必须满足:
补偿设计原则:
1. 可补偿性:每个正向事务都有对应的逆向操作
2. 幂等性:补偿操作可重复执行
3. 交换律:正向 + 补偿 ≈ 无操作(最终状态与未执行一致)
示例:订单流程
正向事务 补偿事务
创建订单(待支付) 取消订单(关闭)
扣减库存(冻结) 释放库存(解冻)
扣减账户(冻结) 退还账户(解冻)
增加积分 扣减积分补偿示例代码:
// 正向操作
@Compensable(compensationMethod = "cancelCreateOrder")
public void createOrder(Order order) {
order.setStatus(OrderStatus.PENDING);
orderDao.insert(order);
}
// 补偿操作
public void cancelCreateOrder(Order order) {
Order db = orderDao.selectById(order.getId());
if (db != null && db.getStatus() == OrderStatus.PENDING) {
db.setStatus(OrderStatus.CANCELLED);
orderDao.updateById(db);
}
}消息最终一致性
可靠消息 + 本地消息表
核心思路:将本地事务与消息发送绑定在同一个数据库事务中,通过定时任务轮询重试确保消息可靠投递。
生产者 消费者
│ │
│ 1. 执行业务 → 插入本地消息表 │
│ (同数据库事务) │
│ │
│ 2. 定时任务轮询未发送消息 │
│ → 发送到 MQ │
│ → 标记消息为已发送 │
│ ───────────→│ 3. 消费消息
│ │ → 执行业务
│ │ → 确认消费
│ 4. 消费确认回调 │
│ → 标记消息为已完成 │
│ │-- 本地消息表设计
CREATE TABLE `t_local_message` (
`id` bigint NOT NULL AUTO_INCREMENT,
`biz_id` varchar(64) NOT NULL COMMENT '业务ID',
`biz_type` varchar(32) NOT NULL COMMENT '业务类型',
`message_body` text NOT NULL COMMENT '消息内容(JSON)',
`status` tinyint NOT NULL DEFAULT '0' COMMENT '0-待发送 1-已发送 2-已完成 3-失败',
`retry_count` int NOT NULL DEFAULT '0' COMMENT '重试次数',
`create_time` datetime NOT NULL,
`update_time` datetime NOT NULL,
PRIMARY KEY (`id`),
INDEX `idx_status` (`status`, `retry_count`)
) ENGINE=InnoDB;@Transactional
public void createOrderAndSendMessage(Order order) {
// 1. 执行业务操作
orderDao.insert(order);
// 2. 插入本地消息表(同事务)
LocalMessage msg = new LocalMessage();
msg.setBizId(order.getId().toString());
msg.setBizType("order_create");
msg.setMessageBody(JSON.toJSONString(order));
msg.setStatus(0);
localMessageDao.insert(msg);
}
// 定时任务轮询发送
@Scheduled(fixedDelay = 5000)
public void retryUnsentMessages() {
List<LocalMessage> unsent = localMessageDao.selectByStatus(0, 100);
for (LocalMessage msg : unsent) {
try {
rabbitTemplate.convertAndSend("order.exchange", "order.create", msg.getMessageBody());
localMessageDao.updateStatus(msg.getId(), 1); // 标记已发送
} catch (Exception e) {
localMessageDao.incrementRetry(msg.getId());
log.error("消息发送失败", e);
}
}
}优缺点:
| 特性 | 说明 |
|---|---|
| 优点 | 实现简单,不依赖 MQ 高级特性,与本地事务强绑定 |
| 缺点 | 耦合消息表,业务库压力大,定时扫描有延迟 |
RocketMQ 事务消息
RocketMQ 原生支持事务消息,通过半消息(Half Message) 机制实现分布式事务。
生产者 Broker
│ │
│ 1. 发送半消息(half) │
│──────────────────────────────────────→│
│ 半消息(暂不可见) │
│ │
│ 2. 执行本地事务 │
│ │
│ 3. 提交/回滚事务消息 │
│──────────────────────────────────────→│
│ │
│ (如果第 3 步未收到) │
│←──────── 4. 回查(check) ─────────────│
│ │
│ 5. 返回本地事务状态 │
│──────────────────────────────────────→│
│ │
│ (事务提交后) │
│ 半消息 → 可见 │
│ 消费者可消费 │@Component
public class OrderTransactionListener implements TransactionListener {
@Autowired
private OrderMapper orderMapper;
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
try {
// 2. 执行本地事务
Order order = (Order) arg;
orderMapper.insert(order);
return LocalTransactionState.COMMIT_MESSAGE;
} catch (Exception e) {
return LocalTransactionState.ROLLBACK_MESSAGE;
}
}
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
// 4. 回查:检查本地事务是否已提交
String orderId = msg.getKeys();
Order order = orderMapper.selectById(orderId);
if (order != null) {
return LocalTransactionState.COMMIT_MESSAGE;
}
return LocalTransactionState.UNKNOW;
}
}// 发送事务消息
@Autowired
private RocketMQTemplate rocketMQTemplate;
public void createOrderTransaction(Order order) {
Message<String> message = MessageBuilder
.withPayload(JSON.toJSONString(order))
.setHeader(MessageConst.PROPERTY_KEYS, order.getId().toString())
.build();
TransactionSendResult result = rocketMQTemplate.sendMessageInTransaction(
"order-topic",
message,
order
);
}优缺点:
| 特性 | 说明 |
|---|---|
| 优点 | 无消息表侵入业务库,可靠性高,回查机制保证最终一致 |
| 缺点 | 依赖 RocketMQ,不支持 Kafka/RabbitMQ,需部署事务消息 Broker |
Seata 框架详解
Seata(Simple Extensible Autonomous Transaction Architecture)是阿里开源的一站式分布式事务解决方案。
架构(TC / TM / RM)
┌──────────────────┐
│ TC(事务协调者) │
│ Transaction │
│ Coordinator │
└────────┬─────────┘
│
┌────────────────────┼────────────────────┐
│ │ │
┌────┴─────┐ ┌────┴─────┐ ┌────┴─────┐
│ TM │ │ TM │ │ TM │
│ 事务管理器 │ │ 事务管理器 │ │ 事务管理器 │
├──────────┤ ├──────────┤ ├──────────┤
│ RM │ │ RM │ │ RM │
│ 资源管理器 │ │ 资源管理器 │ │ 资源管理器 │
├──────────┤ ├──────────┤ ├──────────┤
│ 数据库 │ │ 数据库 │ │ 数据库 │
└──────────┘ └──────────┘ └──────────┘
服务 A 服务 B 服务 C| 角色 | 全称 | 职责 |
|---|---|---|
| TC | Transaction Coordinator | 维护全局事务和分支事务的状态,驱动提交或回滚 |
| TM | Transaction Manager | 定义全局事务范围,告知 TC 开启/提交/回滚 |
| RM | Resource Manager | 管理分支事务的资源,与 TC 通信注册分支 |
AT 模式(自动补偿)
AT 模式是 Seata 的核心模式,对业务无侵入,通过代理数据源自动生成逆向 SQL。
工作原理:
第一阶段(业务 SQL + undo log):
1. 解析业务 SQL,生成前置镜像(Before Image)
2. 执行业务 SQL
3. 生成后置镜像(After Image)
4. 将 undo log 和业务 SQL 在同一个本地事务中提交
示例:UPDATE account SET money = money - 100 WHERE id = 1
Before Image: {id: 1, money: 1000} ← 修改前数据
After Image: {id: 1, money: 900} ← 修改后数据
undo log: 插入到 seata_undo_log 表
第二阶段(提交/回滚):
提交:异步删除 undo log(快速,无需业务干预)
回滚:用 Before Image 生成逆向 SQL 恢复数据全局锁:
- AT 模式通过全局锁保证写隔离
- 事务 A 持有全局锁期间,其他事务不能修改同一数据
- 事务 A 释放全局锁后,其他事务才能继续
优缺点:
| 特性 | 说明 |
|---|---|
| 优点 | 无业务侵入,自动生成补偿 SQL,使用简单 |
| 缺点 | 性能开销较大(生成镜像、全局锁),不适用于大事务 |
TCC 模式
Seata 的 TCC 模式需要业务方实现 Try / Confirm / Cancel 接口。
@LocalTCC
public interface AccountTccAction {
@TwoPhaseBusinessAction(
name = "accountTcc",
commitMethod = "confirm",
rollbackMethod = "cancel"
)
boolean tryMethod(
@BusinessActionContextParameter(paramName = "accountId") Long accountId,
@BusinessActionContextParameter(paramName = "amount") BigDecimal amount
);
boolean confirm(BusinessActionContext ctx);
boolean cancel(BusinessActionContext ctx);
}@Service
public class AccountTccActionImpl implements AccountTccAction {
@Autowired
private AccountMapper accountMapper;
@Override
public boolean tryMethod(Long accountId, BigDecimal amount) {
// Try:冻结金额
int rows = accountMapper.freezeAmount(accountId, amount);
return rows > 0;
}
@Override
public boolean confirm(BusinessActionContext ctx) {
// Confirm:扣减冻结金额
Long accountId = Long.parseLong(ctx.getActionContext("accountId").toString());
BigDecimal amount = new BigDecimal(ctx.getActionContext("amount").toString());
accountMapper.deductFrozen(accountId, amount);
return true;
}
@Override
public boolean cancel(BusinessActionContext ctx) {
// Cancel:释放冻结金额
Long accountId = Long.parseLong(ctx.getActionContext("accountId").toString());
BigDecimal amount = new BigDecimal(ctx.getActionContext("amount").toString());
accountMapper.unfreezeAmount(accountId, amount);
return true;
}
}SAGA 模式
Seata 的 SAGA 模式通过状态机引擎实现,使用 JSON/YAML 定义流程。
# saga-definition.yml
name: "order_saga"
stateMachine:
states:
- name: "create_order"
type: "ServiceTask"
serviceName: "orderService"
serviceMethod: "create"
compensation:
serviceName: "orderService"
serviceMethod: "cancel"
- name: "deduct_stock"
type: "ServiceTask"
serviceName: "stockService"
serviceMethod: "deduct"
compensation:
serviceName: "stockService"
serviceMethod: "release"
- name: "deduct_account"
type: "ServiceTask"
serviceName: "accountService"
serviceMethod: "deduct"
compensation:
serviceName: "accountService"
serviceMethod: "refund"
# 执行顺序
transitions:
- from: "create_order"
to: "deduct_stock"
- from: "deduct_stock"
to: "deduct_account"
- from: "deduct_account"
to: "end"XA 模式
XA 模式基于数据库的 XA 协议,Seata 对其做了分布式协调封装。
XA 执行流程:
1. TM 向 TC 申请开启全局事务
2. RM(数据库)执行 XA START → 执行业务 SQL → XA END
3. RM 执行 XA PREPARE(准备阶段)
4. TC 决策:全部成功 → 各 RM 执行 XA COMMIT
有失败 → 各 RM 执行 XA ROLLBACK| 特性 | 说明 |
|---|---|
| 优点 | 强一致性,数据库原生支持,无业务侵入 |
| 缺点 | 资源锁持有时间长,性能差,依赖数据库 XA 支持 |
方案对比表
| 方案 | 一致性级别 | 性能 | 业务侵入性 | 适用场景 |
|---|---|---|---|---|
| 2PC/3PC | 强一致性 | ★ | 低 | 对一致性要求极高,低并发场景 |
| TCC | 强一致性(最终一致) | ★★★ | 高 | 高性能要求,业务可接受改造 |
| SAGA | 最终一致性 | ★★★★ | 中 | 长事务、复杂的业务流程 |
| 消息最终一致性 | 最终一致性 | ★★★★★ | 中 | 对实时性要求不高的场景 |
| Seata AT | 强一致性(全局锁) | ★★ | 低(无侵入) | 希望低成本接入的微服务项目 |
| Seata TCC | 强一致性(最终一致) | ★★★ | 高 | 需要细粒度控制事务的场景 |
| Seata SAGA | 最终一致性 | ★★★★ | 中 | 复杂长流程、需要可视化编排 |
| Seata XA | 强一致性 | ★ | 低(无侵入) | 数据一致性要求极高,并发低 |
选型建议
对一致性要求极高、并发低
└→ Seata AT / XA(自动补偿,无侵入)
高性能、业务可改造
└→ TCC / Seata TCC(资源预留,细粒度控制)
长流程、复杂业务编排
└→ SAGA / Seata SAGA(状态机编排,补偿机制)
高吞吐、最终一致性可接受
└→ 消息最终一致性(RocketMQ 事务消息)
需要强一致性且数据库支持
└→ XA(数据库原生,但性能最差)