Seata AT 模式源码阅读
AT 模式是 Seata 使用最广的模式,核心思路是:代理数据源,拦截业务 SQL,自动生成前后镜像并记录 undo_log,二阶段根据 undo_log 回滚。本文从源码角度拆解整条链路,涉及 seata-rm-datasource 模块的核心类。
源码模块总览
seata 源码核心模块:
├─ seata-rm-datasource:数据源代理、镜像生成、undo_log 管理(本文主角)
├─ seata-rm:分支事务注册、二阶段指令处理
├─ seata-tm:@GlobalTransactional 全局事务管理
├─ seata-core:XID、全局事务/分支事务数据结构、RPC 协议
└─ seata-server:TC 事务协调者一、数据源代理链
AT 模式对业务完全透明的关键在于数据源代理:业务代码拿到的 DataSource 是代理对象,SQL 执行被层层拦截。
代理链
DataSourceProxy
└─ ConnectionProxy(代理 Connection)
└─ PreparedStatementProxy(代理 PreparedStatement)
└─ 真正执行 JDBC 语句DataSourceProxy
public class DataSourceProxy extends AbstractDataSourceProxy {
private final String resourceGroupId; // 资源组
private final String applicationId; // 应用 ID
private final String resourceId; // 数据源标识(jdbc url)
private final TableMetaCache tableMetaCache; // 表结构缓存
@Override
public Connection getConnection() throws SQLException {
Connection targetConnection = dataSource.getConnection();
return new ConnectionProxy(this, targetConnection); // 包一层代理
}
}每个数据源对应一个 RM,resourceId(jdbc url)作为分支事务的资源标识上报给 TC。
ConnectionProxy
ConnectionProxy 是全局事务与本地事务的桥,持有全局锁的注册逻辑:
public class ConnectionProxy extends AbstractConnectionProxy {
private boolean isGlobalLockRequire; // 是否需要全局锁(select for update)
private final List<Savepoint> savepoints = new ArrayList<>();
@Override
public PreparedStatement prepareStatement(String sql) throws SQLException {
// 根据 SQL 类型选择不同的执行器
// 查询走 SelectExecutor,DML 走 UpdateExecutor / InsertExecutor / DeleteExecutor
return new PreparedStatementProxy(this, statement, sql);
}
// 本地事务提交入口
@Override
public void commit() throws SQLException {
// 二阶段分支提交:注册分支事务,然后真正提交本地事务
...
}
}关键点:本地事务与全局事务的绑定
业务方法加了 @GlobalTransactional:
├─ 开启全局事务,生成 XID
└─ 数据源连接创建时,把 XID 绑定到当前线程上下文
ConnectionProxy 执行 commit 时:
├─ 当前线程有全局事务 XID → 走全局分支提交流程
└─ 没有 → 退化为普通本地事务二、前后镜像生成
业务 SQL 执行前,Seata 需要先拍下"前镜像",执行后再拍"后镜像",两者用于回滚与脏写检查。
SQL 执行链路
PreparedStatementProxy.executeUpdate()
└─ 根据 SQL 类型找到执行器
├─ UpdateExecutor:UPDATE 语句
├─ InsertExecutor:INSERT 语句
└─ DeleteExecutor:DELETE 语句UpdateExecutor 流程
public class UpdateExecutor extends AbstractDMLBaseExecutor {
protected Object doExecute() throws SQLException {
// 1. 解析 SQL,得到表名、条件、主键信息
// 2. 生成查询"前镜像"的 SQL:
// SELECT 主键, 被更新列 FROM 表 WHERE 原始条件 FOR UPDATE
TableRecords beforeImage = beforeImage();
// 3. 真正执行 UPDATE
int updateRows = statementProxy.getTargetStatement().executeUpdate(...);
// 4. 生成"后镜像":
// SELECT 主键, 被更新列 FROM 表 WHERE 主键 IN (前镜像中的主键)
TableRecords afterImage = afterImage(beforeImage);
// 5. 生成 undo_log,保存前/后镜像
prepareUndoLog(beforeImage, afterImage);
return updateRows;
}
}前镜像 SQL 的生成
原始 UPDATE:
UPDATE account SET balance = balance - 100 WHERE id = 1
前镜像 SQL(自动改写):
SELECT id, balance FROM account WHERE id = 1 FOR UPDATE
FOR UPDATE 的意义:
├─ 锁定该行,防止并发修改(配合全局锁形成双保险)
└─ 保证前镜像数据在事务期间不被其他本地事务篡改镜像数据结构
TableRecords(表记录集合):
├─ tableName:表名
├─ rows:每行的字段名 → 值映射
└─ 前镜像 = 变更前的行数据
后镜像 = 变更后的行数据三、undo_log 写入
前后镜像必须与业务 SQL 在同一本地事务中持久化,否则回滚无从谈起。
undo_log 表结构
CREATE TABLE `undo_log` (
`id` BIGINT(20) NOT NULL AUTO_INCREMENT,
`branch_id` BIGINT(20) NOT NULL, -- 分支事务 ID
`xid` VARCHAR(100) NOT NULL, -- 全局事务 ID
`context` VARCHAR(128) NOT NULL, -- 上下文(如序列化方式)
`rollback_info` LONGBLOB NOT NULL, -- 前后镜像的序列化内容
`log_status` INT(11) NOT NULL, -- 0 待回滚,1 已回滚
`log_created` DATETIME NOT NULL,
`log_modified` DATETIME NOT NULL,
PRIMARY KEY (`id`),
UNIQUE KEY `ux_undo_log` (`xid`, `branch_id`)
) ENGINE = InnoDB;AbstractUndoLogManager 写入逻辑
public abstract class AbstractUndoLogManager {
protected void flushUndoLogs(ConnectionProxy cp) {
// 1. 从 ConnectionProxy 取出事务期间收集的所有 undo_log
List<UndoLog> undoLogs = cp.getUndoLogs();
// 2. 序列化前后镜像(默认 jackson,可配 kryo)
// 3. 插入 undo_log 表
insertUndoLogWithNormal(undoLogs, ...);
// 4. 删除多余的历史 undo_log(按保留天数)
deleteUndoLogByLogCreated(...);
}
}为什么能同事务
业务 SQL 与 undo_log 插入用的是同一个 ConnectionProxy、
同一个本地数据库事务:
├─ 业务提交 → undo_log 一起提交
└─ 业务回滚 → undo_log 一起回滚四、分支事务注册
本地事务执行前,RM 需要向 TC 注册分支事务,拿到 branchId。
注册时机
ConnectionProxy.commit() 流程:
├─ 1. 判断当前线程是否存在全局事务(XID)
├─ 2. 存在 → 注册分支事务(向 TC 发 BranchRegisterRequest)
│ └─ TC 返回 branchId
├─ 3. 本地事务真正 commit
├─ 4. 发送 BranchReportRequest 上报分支状态
└─ 5. 提交成功 → 删除 undo_log(异步)注册消息结构
BranchRegisterRequest:
├─ xid :全局事务 ID
├─ resourceId :数据源标识(jdbc url)
├─ branchType :分支类型(AT / TCC / SAGA / XA)
├─ resourceGroupId:资源组
├─ lockKeys :本次操作涉及的行(表名:主键值集合)
└─ applicationData:附加数据lockKeys 是全局锁的核心,TC 根据它写入 lock_table。
五、全局锁管理
AT 模式用全局锁(lock_table)解决分布式写隔离,防止两个全局事务并发修改同一行。
lock_table 表结构
CREATE TABLE `lock_table` (
`row_key` VARCHAR(128) NOT NULL, -- 资源ID + 表名 + 主键值
`xid` VARCHAR(128), -- 持有锁的全局事务
`transaction_id` BIGINT,
`branch_id` BIGINT,
`resource_id` VARCHAR(256),
`table_name` VARCHAR(32),
`pk` VARCHAR(36),
PRIMARY KEY (`row_key`)
);全局锁的获取时机
UpdateExecutor 执行流程:
├─ 1. 生成前镜像(FOR UPDATE 锁住本库行)
├─ 2. 真正执行 UPDATE(本地行锁)
├─ 3. 生成后镜像
├─ 4. 注册分支事务时,把 lockKeys 传给 TC
│ └─ TC 在 lock_table 尝试插入全局锁
│ ├─ 插入成功 → 获得全局锁,继续
│ └─ 插入冲突(row_key 已存在)→ 说明其他全局事务持有
│ └─ 重试 / 抛 LockConflictException
└─ 5. 提交本地事务写隔离示例
事务A(XID=A1):UPDATE account SET balance=0 WHERE id=1 → 持有全局锁
事务B(XID=B1):UPDATE account SET balance=100 WHERE id=1
├─ 本地 FOR UPDATE 会阻塞吗?不会!行锁在 A 提交后已释放
├─ 但注册分支时发现 lock_table 中 id=1 已被 A1 锁定
└─ B 全局锁获取失败 → 重试等待 A1 提交后释放 → 再执行
结论:全局锁把"跨事务的写冲突"拦截在分支注册阶段锁等待与超时
// 全局锁获取带重试与超时(默认 45 秒)
while (remaining < maxRetry) {
try {
result = branchRegister(...); // 内部会尝试获取全局锁
break;
} catch (LockConflictException e) {
// 随机退避后重试
sleep(random delay);
}
}
// 超时仍冲突 → 抛出异常,全局事务回滚六、二阶段提交
全局事务所有分支都成功,TM 通知 TC 提交,TC 通知各 RM 做二阶段提交。
提交路径
TM.commit()
└─ GlobalCommitRequest → TC
└─ TC 状态改为 Commit
└─ 通知各分支 RM:BranchCommitRequest
└─ RM 收到后:
├─ 删除 undo_log(异步批量)
└─ 释放全局锁(删除 lock_table 记录)为什么提交只需要删 undo_log
因为业务数据已经在本地事务提交时落库,
二阶段提交只需要"清理现场":
├─ 删 undo_log:不再需要回滚信息
└─ 删全局锁:允许其他事务操作这些行异步删除失败也不影响数据一致性,靠定时任务兜底清理。
七、二阶段回滚
全局事务任一分支失败,TC 通知所有 RM 回滚。RM 用 undo_log 的前镜像恢复数据。
回滚路径
TM.rollback()
└─ GlobalRollbackRequest → TC
└─ TC 状态改为 Rollback
└─ 通知各分支 RM:BranchRollbackRequest
└─ RM 收到后执行 undo:
├─ 1. 校验 undo_log 是否有效
├─ 2. 执行脏写检查(后镜像 vs 当前数据)
├─ 3. 用前镜像生成反向 SQL 恢复数据
├─ 4. 标记 undo_log 状态为已回滚
└─ 5. 释放全局锁undo 核心流程
public void undo(ConnectionProxy cp, String xid, long branchId) {
// 1. 查 undo_log
UndoLog log = findUndoLog(xid, branchId);
// 2. 脏写检查:把当前数据与后镜像比对
TableRecords currentRecords = queryCurrentRecords(...);
boolean dirty = dataValidationManager.check(currentRecords, afterImage);
if (dirty) {
// 后镜像与当前数据不一致 → 说明被其他事务改过
// 按策略处理:抛异常人工介入 / 继续覆盖回滚
throw new BranchRollbackException("脏数据,无法回滚");
}
// 3. 用前镜像反向恢复:把当前行改回前镜像的值
// 自动生成 UPDATE 语句,逐列还原
undoExecutor.executeOn(...);
// 4. 标记 undo_log 已回滚
log.setStatus(UNDO_LOG_STATUS_ROLLBACK_DONE);
}脏写检查的意义
后镜像 = 本事务修改后的值
当前数据 = 回滚时数据库里的值
场景:事务A 回滚前,事务C 又改了同一行
├─ 当前数据 ≠ 后镜像 → 说明数据被他人修改
└─ 直接按前镜像覆盖回滚会覆盖 C 的修改 → 数据错乱
处理策略:
├─ 默认:抛异常,人工介入修复(安全优先)
└─ 可配置:允许覆盖回滚(业务兜底)回滚时如何保证行锁
回滚前会重新对相关行执行 SELECT ... FOR UPDATE:
├─ 确认行未被其他事务修改
└─ 配合脏写检查保证回滚安全八、读隔离
AT 模式默认读已提交,跨事务读可能读到未提交的中间数据,需要时用 SELECT FOR UPDATE 加全局锁。
默认读隔离(读已提交)
业务 A 修改了数据但未提交全局事务:
├─ 本地数据已改(本地事务已提交,undo_log 未删)
├─ 业务 B 普通 SELECT → 读到 A 的未提交修改(脏读)
└─ 若 A 回滚 → B 读到的数据是"幽灵数据"
为什么允许:多数业务可容忍短暂脏读,换取性能强读隔离(SELECT FOR UPDATE)
// 业务需要强一致读取时
// @GlobalTransactional 方法内使用:
SELECT * FROM account WHERE id = 1 FOR UPDATE;执行流程(SelectForUpdateExecutor):
├─ 1. 解析 SQL,得到主键
├─ 2. 向 TC 申请全局锁(lockKeys = 表名:主键)
├─ 3. 获取成功 → 执行 SELECT ... FOR UPDATE(本地行锁)
│ → 数据一定是最新已提交的(其他全局事务无法并发修改)
└─ 4. 获取失败 → 说明有未提交的全局事务持有锁 → 等待/超时九、AT 模式源码关键链路图
业务方法(@GlobalTransactional)
│ TM: GlobalTransactionalInterceptor 拦截
│ ├─ 开启全局事务 → 生成 XID
│ └─ XID 绑定到线程上下文
▼
数据源操作(DataSourceProxy 代理)
│ PreparedStatementProxy 拦截 SQL
│ ├─ UpdateExecutor:前镜像 → UPDATE → 后镜像
│ └─ prepareUndoLog:写 undo_log(同本地事务)
▼
本地事务提交(ConnectionProxy.commit)
│ ├─ 注册分支事务(BranchRegisterRequest + lockKeys)
│ ├─ TC 写 lock_table(全局锁)
│ ├─ 本地事务真正提交
│ └─ 上报分支状态(BranchReportRequest)
▼
二阶段(TC 统一指挥)
├─ 提交:RM 删 undo_log + 删全局锁
└─ 回滚:RM 脏写检查 → 前镜像恢复 → 标记已回滚 → 释放全局锁总结
AT 模式的优雅在于把分布式事务的复杂度封装在数据源代理层:业务只写普通 SQL 加一个注解,框架自动完成镜像生成、undo_log 管理、分支注册与二阶段执行。理解 DataSourceProxy → ConnectionProxy → 执行器 → 分支注册 → 全局锁 → 二阶段回滚 这条链路,就抓住了 AT 模式的源码骨架,也能解释生产中的各种异常现象(锁冲突、脏数据、undo_log 堆积等)。