Seata分布式事务AT模式源码分析:从全局事务到分支事务的完整链路
引言
先想象一个场景:你在电商平台下单,订单服务扣减库存,账户服务扣款,积分服务加积分。这三个操作分属三个微服务,各自使用独立的数据库。如果扣款成功但扣库存失败,用户会收到一笔“幽灵订单”——钱付了,货没发。
这就是典型的分布式事务问题。传统本地事务的ACID在跨服务、跨数据库场景下失效,我们需要一种机制来协调多个资源管理器(RM),保证要么全部成功,要么全部回滚。
业界方案众多:2PC(两阶段提交)、TCC(Try-Confirm-Cancel)、Saga、本地消息表……而Seata的AT模式(Automatic Transaction)是其中对业务侵入最小的一种——它通过代理数据源自动生成反向SQL,业务代码几乎无感知。本文将从源码层面剖析AT模式的核心机制,并配合可运行的实战代码,带你理解其设计精髓。
核心概念:生活类比 + 技术定义
生活类比:餐厅“记账本”模式
想象一家餐厅,顾客点餐(业务操作),但厨房并不直接改菜单库存,而是先在“记账本”上记录:原库存是多少,改成多少,以及改之前的快照。如果顾客取消订单(回滚),后厨就照着记账本把库存改回去;如果确认下单成功(提交),记账本撕掉即可。
AT模式就是这套逻辑:在业务SQL执行前记录数据快照(before image),执行后记录新快照(after image),如果后续某分支失败,就用反向SQL把数据恢复成before image。
技术定义
AT模式是Seata对2PC协议的优化实现,核心思想是:
- 一阶段:业务SQL在本地事务内执行,同时注册分支事务并记录undo_log。
- 二阶段提交:异步删除undo_log,释放行锁。
- 二阶段回滚:根据undo_log生成反向SQL,将数据恢复原状。
它与传统2PC的区别在于:业务SQL在本地事务中直接提交(释放数据库锁),而不是像X/A协议那样持有锁直到全局事务结束。
源码深度分析:核心链路拆解
整体架构
核心类关系
Seata的AT模式源码集中在seata-rm-datasource模块。核心类包括:
DataSourceProxy:代理用户配置的数据源,返回ConnectionProxy。
ConnectionProxy:代理JDBC Connection,拦截commit/rollback。
StatementProxy:代理Statement,在SQL执行前后记录镜像。
UndoLogManager:负责undo_log的写入与删除。
BranchTransactionalOperator:负责分支事务的注册与上报。
一阶段:SQL执行与镜像记录
当业务代码调用Connection.commit()时,ConnectionProxy会介入。我们看关键代码(版本2.x,略作简化):
// ConnectionProxy.java
public void commit() throws SQLException {
// 如果当前不在全局事务中,直接提交
if (context.isGlobalLockRequire() == false && !context.hasUndoLog()) {
targetConnection.commit();
return;
}
// 1. 生成undo_log并写入当前本地事务
processLocalCommitWithUndoLog();
// 2. 注册分支事务到TC(Transaction Coordinator)
Long branchId = branchRegister();
// 3. 上报分支事务状态
reportBranchStatus(branchId);
// 4. 本地事务提交(此时数据已可见,但全局锁未释放)
targetConnection.commit();
// 5. 移除本地线程中的事务上下文
context.reset();
}关键在于processLocalCommitWithUndoLog(),它内部会遍历当前连接执行过的所有SQL,为每条SQL生成undo_log:
private void processLocalCommitWithUndoLog() throws SQLException {
// 获取当前连接上执行过的SQL记录
List<SQLUndoLog> sqlUndoLogs = context.getUndoItems();
if (sqlUndoLogs.isEmpty()) {
return;
}
// 生成undo_log并一起写入本地事务
UndoLogManager.flushUndoLogs(sqlUndoLogs);
}镜像记录的生成:SQLVisitor
镜像记录的生成依赖SQLVisitor对SQL的解析。以UPDATE语句为例,MySQLUpdateRecognizer会解析出表名、SET字段、WHERE条件,然后执行两条查询:
- Before Image:
SELECT * FROM table WHERE [条件] FOR UPDATE(加锁)
- After Image:
SELECT * FROM table WHERE [主键](执行后)
关键代码在AbstractUndoLogGenerator中:
// UpdateExecutor.java(简化)
public ExecuteResult<UpdateExecutor> execute() throws SQLException {
// 先查before image
TableRecords beforeImage = beforeImage();
// 执行真正的更新SQL
statement.execute();
// 再查after image
TableRecords afterImage = afterImage();
// 构建undo_log
SQLUndoLog undoLog = buildUndoLog(beforeImage, afterImage);
return new ExecuteResult<>(undoLog);
}行锁与全局锁
这里有一个关键设计:Seata在before image阶段使用SELECT ... FOR UPDATE来锁定记录。但这只是本地事务锁,事务提交后锁就释放了。那怎么防止其他全局事务并发修改呢?
答案是全局锁(Global Lock)。在注册分支事务时,TC会为每行数据记录一个全局锁(在lock_table中),锁的key是resourceId + tableName + pkValue。其他事务在修改同一行数据前,会尝试获取全局锁,如果锁被占用则等待。
但这里有个性能权衡:AT模式在一阶段提交后就释放了数据库本地锁,全局锁由TC集中管理。这比X/A协议持有数据库锁到二阶段结束要高效得多,但代价是TC成为性能瓶颈。
二阶段:提交与回滚
#### 提交(异步删除undo_log)
全局事务成功,TC通知所有分支事务提交。此时RM端做的事非常简单:
// AsyncWorker.java
public void doBranchCommit() {
// 异步删除undo_log
UndoLogManager.deleteUndoLog(branchId);
}因为本地事务已经提交,数据不需要改动,只需清理日志即可。异步处理避免阻塞二阶段提交线程。
#### 回滚(反向SQL生成)
全局事务失败,TC通知分支回滚。RM端根据undo_log中的before image生成反向SQL:
// UndoLogManager.java
public void undo(Connection conn, String xid, long branchId) {
// 解析undo_log内容
UndoLogContent content = parseUndoLog();
// 根据before image生成反向SQL
String reverseSQL = buildReverseSQL(content);
// 校验当前数据是否等于after image(防止脏写)
validateCurrentData(content);
// 执行反向SQL
executeReverseSQL(conn, reverseSQL);
}脏写校验是关键:如果当前数据和after image不一致,说明有其他事务修改过数据,此时不能盲目回滚,需要人工介入。
反向SQL的生成逻辑如下:
- INSERT →
DELETE FROM table WHERE pk = ?
- DELETE →
INSERT INTO table (cols) VALUES (values)
- UPDATE →
UPDATE table SET col = before_value WHERE pk = ?
实战代码:三个完整的可运行示例
示例一:最小化AT模式集成(Spring Boot + MyBatis)
// 1. 引入依赖(pom.xml)
// <dependency>
// <groupId>io.seata</groupId>
// <artifactId>seata-spring-boot-starter</artifactId>
// <version>1.7.0</version>
// </dependency>
// 2. 配置application.yml
// seata:
// enabled: true
// tx-service-group: my_test_tx_group
// service:
// vgroup-mapping:
// my_test_tx_group: default
// grouplist:
// default: 127.0.0.1:8091
// 3. 业务代码
@Service
public class OrderServiceImpl implements OrderService {
@Autowired
private AccountFeignClient accountClient;
@Autowired
private InventoryFeignClient inventoryClient;
@Override
@GlobalTransactional(name = "create-order", rollbackFor = Exception.class)
public void createOrder(OrderRequest request) {
// 本地事务:插入订单
orderDao.insert(buildOrder(request));
// 远程调用:扣减库存(也是AT模式)
inventoryClient.deduct(request.getProductId(), request.getQuantity());
// 远程调用:扣减账户余额(AT模式)
accountClient.debit(request.getUserId(), request.getAmount());
// 如果这里抛出异常,所有已执行的本地事务都会回滚
// 包括远程服务已经提交的事务
}
}示例二:手动控制事务边界(不依赖Spring AOP)
// 适用于非Spring环境或需要更细粒度控制的场景
public class ManualTxExample {
public void manualTransaction() throws Exception {
// 1. 获取全局事务管理器
GlobalTransaction tx = GlobalTransactionContext.getCurrentOrCreate();
try {
// 2. 开启全局事务
tx.begin(60000, "manual-tx");
// 3. 执行业务操作(此时DataSourceProxy会自动拦截)
try (Connection conn = dataSourceProxy.getConnection()) {
// 业务SQL
try (PreparedStatement ps = conn.prepareStatement(
"UPDATE account SET balance = balance - ? WHERE user_id = ?")) {
ps.setBigDecimal(1, new BigDecimal("100"));
ps.setString(2, "U10001");
ps.executeUpdate();
}
}
// 4. 也可以调用其他服务的AT模式接口
// remoteService.doSomething();
// 5. 全局提交
tx.commit();
} catch (Exception e) {
// 6. 全局回滚
tx.rollback();
throw e;
}
}
}示例三:全局锁冲突处理
// 模拟两个并发事务修改同一行数据
public class LockConflictExample {
private static final int MAX_RETRY = 5;
private static final long RETRY_INTERVAL_MS = 100;
public void updateWithRetry(String userId, BigDecimal amount) {
for (int i = 0; i < MAX_RETRY; i++) {
try {
// 尝试获取全局锁并执行更新
doUpdate(userId, amount);
return; // 成功则退出
} catch (LockConflictException e) {
// 锁冲突,等待重试
if (i == MAX_RETRY - 1) {
throw new RuntimeException("获取全局锁超时", e);
}
try {
Thread.sleep(RETRY_INTERVAL_MS * (i + 1)); // 指数退避
} catch (InterruptedException ie) {
Thread.currentThread().interrupt();
throw new RuntimeException("被中断", ie);
}
}
}
}
@GlobalTransactional
public void doUpdate(String userId, BigDecimal amount) {
// 这个查询会触发SELECT FOR UPDATE(通过Seata的代理)
accountDao.updateBalance(userId, amount);
}
}方案对比:AT模式 vs TCC vs Saga vs X/A
| 方案 | 侵入性 | 一致性 | 性能 | 适用场景 |
|---|---|---|---|---|
| AT模式 | 低(需建undo_log表) | 最终一致(有中间态) | 较高(本地事务) | 通用场景,适合大多数业务 |
| TCC | 高(需实现Try/Confirm/Cancel) | 强一致 | 高 | 对一致性要求极高,如金融转账 |
| Saga | 中(需实现补偿逻辑) | 最终一致 | 高 | 长事务,如订单流程 |
| X/A | 低 | 强一致 | 低(锁持有时间长) | 单体数据库场景,已很少使用 |
AT的优势:代码侵入最小,只需一个注解。业务SQL照写,Seata自动处理回滚逻辑。
AT的劣势:
- 性能开销:每条SQL需要额外2次查询(before/after image)+ undo_log写入。
- 脏写风险:如果手动绕过Seata直接操作数据库,可能导致undo_log校验失败。
- 全局锁依赖TC,TC故障会影响整个事务。
最佳实践与避坑指南
最佳实践
- 正确配置undo_log表:每个业务库都需要建
undo_log表,SQL脚本在Seata官方仓库可以找到。表结构必须完全一致,否则会报错。
- 合理设置事务超时:
@GlobalTransactional(timeoutMills = 60000)默认60秒,如果业务逻辑耗时较长(如多个远程调用),需要适当调大,否则会提前回滚。
- 全局锁冲突重试:在高并发场景,不同全局事务可能修改同一行数据,导致
LockConflictException。建议在业务层做重试,重试间隔使用指数退避。
- 监控TC的运行状态:Seata的TC(Server端)是单点,虽然支持集群模式,但需要额外配置。建议使用
seata-server的集群模式并配合注册中心。
常见坑
坑1:MySQL8.0驱动导致undo_log写入失败
# 错误:使用旧版驱动,datetime字段类型映射问题
spring:
datasource:
driver-class-name: com.mysql.jdbc.Driver # 过时
# 正确:使用新版驱动
spring:
datasource:
driver-class-name: com.mysql.cj.jdbc.Driver同时,undo_log表的log_created和log_modified字段类型应为datetime(6),否则高并发下可能精度丢失。
坑2:批量操作无法获取正确的镜像
// 错误:批量操作可能无法生成正确的undo_log
statement.executeBatch();
// 正确:逐条执行,或使用rewriteBatchedStatements=true
// jdbc:mysql://localhost:3306/db?rewriteBatchedStatements=true坑3:多数据源场景下的事务边界模糊
如果同一个事务中操作多个数据源,事务边界会变得不清晰。建议每个数据源独立事务,通过全局事务协调。否则可能出现部分分支已提交,但部分分支未注册的情况。
坑4:本地事务与全局事务混用
// 错误:在@GlobalTransactional方法中,又开启了本地事务
@GlobalTransactional
public void badMethod() {
@Transactional // 这会创建嵌套的本地事务,可能破坏AT模式
public void innerMethod() {
// ...
}
}Seata的AT模式会代理DataSource,如果方法上同时有@Transactional和@GlobalTransactional,本地事务会先于全局事务提交,导致undo_log提前删除,回滚失败。
总结
Seata的AT模式通过“镜像记录 + 反向SQL”的巧妙设计,在保证分布式事务一致性的同时,将对业务代码的侵入降到了最低。它的核心思想是:既然业务SQL是幂等的,那我们就可以通过记录执行前后的状态,在需要时反向恢复。
但AT模式并非银弹,它牺牲了一定的性能(镜像记录开销)和一致性强度(最终一致而非强一致),换来了开发效率。在金融级场景,TCC可能是更好的选择;在长事务场景,Saga更合适;而在大多数微服务业务中,AT模式是最平衡的切入点。
延伸思考:Seata后续版本引入了@GlobalLock注解来处理“非事务内读取未提交数据”的问题,以及apm监控集成。随着云原生的发展,Seata也在探索与Kubernetes、Service Mesh的融合。理解AT模式的设计哲学,对你在其他分布式协调框架(如DTX、Saga实现)中举一反三,大有裨益。
*最后留一个问题供你思考:如果AT模式遇到“跨分库的分布式事务”,单个undo_log无法覆盖所有分库的情况,Seata是如何处理的?欢迎在评论区探讨。*