交互式查询场景中,一条慢 SQL 可能长时间占用数据库连接,把整个系统拖垮。
Statement.setQueryTimeout可以作为保护手段:快速失败,而不是无限等待。但”取消”到底由谁执行、如何执行,不同数据库差异很大。H2 是服务端协作式取消,MySQL 是客户端监控线程 + KILL QUERY。
最后结合 MyBatis 的
defaultStatementTimeout配置,看超时时间怎么发货作用。
Statement.setQueryTimeout(int seconds) 是 JDBC 标准的超时入口,但落地方式因数据库而异。H2 把超时实现为一条
SET命令,执行期间由算子循环周期性调用checkCanceled()检查是否到期。🎈 协作式(cooperative)的取消:依赖代码主动检查,一个长循环若没有检查点,可以”绕过”超时。
setQueryTimeout 并不是直接给 Statement 存个值,而是通过 JDBC 下发一条 SQL 命令 SET QUERY_TIMEOUT ?:
sequenceDiagram
participant 应用 as 应用
participant JDBC as JdbcStatement/Connection
participant Set as Set 命令
participant Session as Session
应用->>JDBC: 1. setQueryTimeout(1)
JDBC->>Set: prepareCommand("SET QUERY_TIMEOUT ?", 1000)
Set->>Session: session.setQueryTimeout()
Note over Session: 限幅为 maxQueryTimeout,写 this.queryTimeout
应用->>JDBC: 2. executeQuery(...)
JDBC->>Session: Command.executeQuery → setCurrentCommand
Note over Session: cancelAt = now + queryTimeout
Session->>Session: 3. 算子循环 checkCanceled()
Note over Session: now >= cancelAt → 抛 STATEMENT_WAS_CANCELED
调用链上的关键类:
/**
* JDBC 层入口:秒转毫秒,以 SET QUERY_TIMEOUT 命令下发
* @see org.h2.jdbc.JdbcStatement#setQueryTimeout
* @see org.h2.jdbc.JdbcConnection#setQueryTimeout
*/
// JdbcStatement.setQueryTimeout(seconds) → JdbcConnection.setQueryTimeout(seconds * 1000)
// → prepareCommand("SET QUERY_TIMEOUT ?", ...) → Set.update() → session.setQueryTimeout
/**
* 最终写入 session.queryTimeout(毫秒),并做全局限幅
* @see org.h2.engine.Session#setQueryTimeout
*/
public void setQueryTimeout(int queryTimeout) {
// 取数据库全局配置,session 级不能超过它
int max = database.getSettings().maxQueryTimeout;
if (max != 0 && (max < queryTimeout || queryTimeout == 0)) {
queryTimeout = max; // 超限 → 按全局上限收口
}
this.queryTimeout = queryTimeout;
}
🎈 H2 的 setQueryTimeout 作用于当前 Session,且受全局配置 maxQueryTimeout 限幅。
执行查询时先记录”取消截止时间点”,供算子循环周期检查:
/**
* 执行前:session.setCurrentCommand 计算 cancelAt = now + queryTimeout
* @see org.h2.command.Command#executeQuery
* @see org.h2.engine.Session#setCurrentCommand
*/
session.setCurrentCommand(this);
// 若 queryTimeout > 0 → cancelAt = 当前时间 + queryTimeout
/**
* 周期性检查:now >= cancelAt 则中断查询
* @see org.h2.engine.Session#checkCanceled
*/
public void checkCanceled() {
long now = getCurrentTimeNanos();
if (cancelAt != 0 && now >= cancelAt) {
throw DbException.get(ErrorCode.STATEMENT_WAS_CANCELED);
}
}
🙉 协作式取消的代价:取消依赖算子代码主动调用 checkCanceled()。若某段代码是一个不含检查点的长循环(如纯 CPU 计算、大结果集组装),超时期间不会被中断,只能等它自然结束。
MySQL 与 H2 思路完全相反:客户端驱动在发起查询前启动一个监控定时任务,超时后另开连接执行
KILL QUERY中断服务端。
/**
* 发起查询前,若启用且设置了超时,调度一个取消任务
* @see com.mysql.jdbc.StatementImpl#executeQuery
*/
if (locallyScopedConn.getEnableQueryTimeouts()
&& this.timeoutInMillis != 0
&& locallyScopedConn.versionMeetsMinimum(5, 0, 0)) {
timeoutTask = new CancelTask(this);
locallyScopedConn.getCancelTimer().schedule(timeoutTask, this.timeoutInMillis);
}
timeoutInMillis 由 setQueryTimeout(seconds) 提前换算好(seconds * 1000)。CancelTask,其内部再另起线程执行取消。CancelTask 继承 TimerTask,到点后 run() 内再启动一个取消线程,核心逻辑二选一:
/**
* 超时取消:要么直接杀连接,要么另开连接执行 KILL QUERY
* @see com.mysql.jdbc.StatementImpl.CancelTask#run
*/
if (connection.getQueryTimeoutKillsConnection()) {
// ① 配置 queryTimeoutKillsConnection=true → 直接关闭连接
toCancel.wasCancelled = true;
toCancel.wasCancelledByTimeout = true;
connection.realClose(false, false, true, new MySQLStatementCancelledException(...));
} else {
// ② 默认:复制连接属性,另开连接执行 KILL QUERY
cancelConn = connection.duplicate(); // URL 已变则回退 DriverManager.getConnection
cancelStmt = cancelConn.createStatement();
cancelStmt.execute("KILL QUERY " + connectionId);
toCancel.wasCancelled = true;
toCancel.wasCancelledByTimeout = true;
}
⭐ 关键点:
KILL QUERY <connectionId>(注意是 QUERY 而非 CONNECTION),只中断当前查询,不杀掉整个会话。execSQL 等服务端响应,收到中断后抛 MySQLTimeoutException。connection.duplicate();若原连接 URL 已变化,回退 DriverManager.getConnection(origConnURL, origConnProps) 重建连接。queryTimeoutKillsConnection=true 时不再发 KILL QUERY,直接关闭连接,手段更激进。| 维度 | H2 | MySQL 驱动 |
|---|---|---|
| 取消发起方 | 服务端(执行算子循环) | 客户端(监控定时任务) |
| 取消方式 | 协作式检查 checkCanceled() |
新连接执行 KILL QUERY |
| 取消粒度 | 查询内部检查点 | 服务端线程级中断 |
| 盲区 | 无检查点的长循环可”绕过” | 服务端无响应时监控任务徒劳 |
| 额外开销 | 无 | 每次超时需新建连接执行 KILL |
超时配置最常见的落点之一是 MyBatis,全局配置:
<configuration>
<settings>
<!-- 数据库查询/更新执行超过 5 秒仍未响应则超时 -->
<setting name="defaultStatementTimeout" value="5"/>
</settings>
</configuration>
创建 Statement 后,MyBatis 统一把超时值透传给 JDBC:
/**
* 设置 Statement 超时:优先 MappedStatement 级别,其次全局默认
* @see org.apache.ibatis.executor.statement.BaseStatementHandler#setStatementTimeout
*/
protected void setStatementTimeout(Statement stmt) throws SQLException {
Integer timeout = mappedStatement.getTimeout();
if (timeout != null) {
stmt.setQueryTimeout(timeout); // ① MappedStatement 级别
return;
}
Integer defaultTimeout = configuration.getDefaultStatementTimeout();
if (defaultTimeout != null) {
stmt.setQueryTimeout(defaultTimeout); // ② 全局默认
}
}
🎈 优先级:MappedStatement.getTimeout() > configuration.getDefaultStatementTimeout() > 不设置。
SET QUERY_TIMEOUT + checkCanceled()),MySQL 走客户端监控线程 + KILL QUERY。defaultStatementTimeout 将超时统一透传到 JDBC,实现”一处配置、全局生效”。上一篇《Insight h2database MVStore Table 新存储引擎 MVCC 实现原理》梳理了 MVTable 的整体架构与可见性机制。
本文聚焦 UPDATE 操作在 MVTable 中的执行细节。
与直觉上的”原地修改某个字段”不同,MVTable 将 UPDATE 拆解为 remove 旧行 + add 新行两步,每一步都遵循
trySet的”先写 undo、再写数据”。巧妙的是:undoLog 中每条
VersionedValue#value是前一步 data 的值,由此 undoLog 串成一条可从最新值回溯到最初已提交版本的单向链。
先用一张表回顾 MVTable 与 RegularTable 的核心差异:
| 维度 | RegularTable(旧) | MVTable(新/默认) |
|---|---|---|
| 多版本载体 | row.sessionId + row.deleted |
VersionedValue.operationId |
| 未提交数据存放 | 内存 delta 集合 | 持久化 undoLog(write-ahead) |
| 冲突检测 | 行属性对比 sessionId |
trySet 比对 operationId |
| 可见性回溯 | delta 过滤 | undoLog 单向链表回溯 |
👉 完整对比见《H2 数据库 MVStore 与 Regular 存储引擎对比》
MVTable 把版本信息编码进每个值自身,而非行对象上:
/**
* 数据 MVMap 中存储的值,不是裸值,而是带版本的包装
* @see org.h2.mvstore.db.TransactionStore.VersionedValue
*/
static class VersionedValue {
public long operationId; // transactionId(高位) + logId(低位),0 表示已提交
public Object value; // 真正的行数据,null 表示删除标记
}
operationId == 0 → 已提交,全体事务可见operationId != 0 → 属于某个未提交事务,对其他事务隔离operationId 字段,同时承担了隔离(可见性判定)与锁(冲突检测)两个职责。所有写入(Insert / Update / Delete)最终都汇聚到 trySet,它是冲突检测的唯一入口:
/**
* 尝试写入/删除值。若该 key 被其他未提交事务占用,返回 false
* @see org.h2.mvstore.db.TransactionStore.TransactionMap#trySet
*/
public boolean trySet(K key, V value, boolean onlyIfUnchanged) {
VersionedValue current = map.get(key);
VersionedValue newValue = new VersionedValue();
newValue.operationId = getOperationId(transaction.transactionId, transaction.logId);
newValue.value = value;
if (current == null) {
// ① 全新 key:先写 undo → 再 putIfAbsent
transaction.log(mapId, key, current); // undo 存 null
if (map.putIfAbsent(key, newValue) != null) {
transaction.logUndo(); return false;
}
return true;
}
int tx = getTransactionId(current.operationId);
if (tx == 0) {
// ② 已有已提交值:先写 undo(旧值) → 再 CAS 覆盖
transaction.log(mapId, key, current); // undo 存 旧值
if (!map.replace(key, current, newValue)) {
transaction.logUndo(); return false;
}
return true;
}
if (tx == transaction.transactionId) {
// ③ 本事务自己改过:允许重入覆盖
transaction.log(mapId, key, current);
if (!map.replace(key, current, newValue)) {
transaction.logUndo(); return false;
}
return true;
}
// ★ 属于其他未提交事务 → 并发冲突
return false;
}
设计模式(write-ahead 原则):
⭐ transaction.log(mapId, key, current) 总是在数据变更前执行 —— 先把”修改前的值”记入 undoLog,再修改 data MVMap。
🎈 MVTable 处理 UPDATE 并不是”原地改一个值”,而是拆成 两次 trySet 调用,每次对应一个 logId:
sequenceDiagram
participant 调用方 as MVTable.updateRow
participant 主键索引 as MVPrimaryIndex
participant txMap as TransactionMap
调用方->>主键索引: 1. 删除旧行 removeRow(row)
主键索引->>txMap: trySet(key, null) · logId=0
Note over txMap: undo 记旧值 → data 写 null(删除标记)
调用方->>主键索引: 2. 插入新行 addRow(newRow)
主键索引->>txMap: trySet(key, newValue) · logId=1
Note over txMap: undo 记 null(删除标记) → data 写新行
🎈 此为全文最核心的设计:
undoLog 单向回溯链:
bjz (op=operationId_1)
→ 查 undoLog[operationId_1] → oldValue = null (op=operationId_0)
→ 查 undoLog[operationId_0] → oldValue = bjn (op=0, 链尾,已提交)
每条 undo 记录的 VersionedValue#value(即 d[2]),恰好是上一步 trySet 写入 data 的值。由此 undoLog 自然串成一条可回溯的单向链。
/**
* 可见性回溯:顺 operationId 在 undoLog 中查找"被覆盖前的值"
* @see org.h2.mvstore.db.TransactionStore.TransactionMap#getValue
*/
VersionedValue getValue(K key, long maxLog, VersionedValue data) {
while (true) {
if (data == null) return null;
long id = data.operationId;
if (id == 0) return data; // 已提交 → 可见
int tx = getTransactionId(id);
if (tx == transaction.transactionId) {
if (getLogId(id) < maxLog) return data; // 本事务自己的修改 → 可见
}
// 其他未提交事务 → 从 undoLog 回溯到上一个版本
Object[] d = transaction.store.undoLog.get(id); // d[2] 即 oldValue
data = (d == null) ? map.get(key) : (VersionedValue) d[2];
// 循环直到拿到一个可见版本
}
}
🙉 注意:getValue 不回走”原始 MVMap”(因 MVStore 是 Copy-on-Write),而是通过 undoLog 中记录的 oldValue(VersionedValue) 拿到历史版本。越靠近链尾的值越接近最初状态。
以下用两个浏览器 Session 模拟并发 UPDATE,观察 MVCC 行为:
-- 建表并初始化数据
CREATE TABLE city (
id INT(10) NOT NULL AUTO_INCREMENT PRIMARY KEY,
code VARCHAR(40) NOT NULL,
name VARCHAR(40) NOT NULL
);
INSERT INTO city VALUES(1, 'bjx', '北京西');
INSERT INTO city VALUES(2, 'bjn', '北京南');
-- Session A:开启事务,更新 id=2
SET AUTOCOMMIT OFF;
UPDATE city SET code = 'bjz' WHERE id = 2;
-- Session A 事务内查询 → 看到 bjz(自己的修改)
SELECT * FROM city WHERE id = 2;
-- 结果: id=2, code='bjz', name='北京南'
此时 data MVMap 中 id=2 的值为 VersionedValue{operationId=txA_log1, value=(2,'bjz','北京南')},undoLog 有 2 条记录。
-- Session B:查询同一条数据 → 通过 undoLog 回溯,仍看到 bjn
SELECT * FROM city WHERE id = 2;
-- 结果: id=2, code='bjn', name='北京南' → 读到旧版本 ✔
-- Session B:尝试并发更新同一条数据
UPDATE city SET code = 'sjx' WHERE id = 2;
-- ❌ 报错: CONCURRENT_UPDATE_1
⭐ 并发流程解读:
| Session B | 底层机制 | 结果 |
|---|---|---|
SELECT |
getValue 检测到 operationId 属于 txA → 顺 undoLog 回溯 → 读到链尾 bjn(op=0) |
读到旧值 ✔ |
UPDATE |
trySet 检测到 operationId 属于 txA ≠ txB → 返回 false → 上层抛 CONCURRENT_UPDATE_1 |
并发冲突 ❌ |
demo 演示了同一事务内 UPDATE 的完整状态变化,可逐步查看 undoLog 与 data MVMap 的实时联动:
undoLog 与数据 MVMap 中 VersionedValue 的实时变化 · 事务 tx5👉 点击”下一步”观察 logId=0(删除旧行)→ logId=1(插入新行)两个阶段中 undoLog 如何逐步增长,以及回溯链
bjz → null → bjn如何形成。👉 点击”模拟 COMMIT”观察:operationId 全部置零,undoLog 清空,新值成为已提交版本。
提交的本质:把本事务所有改动值的 operationId 清零,并清空 undoLog:
/**
* @see org.h2.mvstore.db.TransactionStore#commit
*/
void commit(Transaction t, long maxLogId) {
for (long logId = 0; logId < maxLogId; logId++) {
Long undoKey = getOperationId(t.getId(), logId);
Object[] op = undoLog.get(undoKey);
MVMap<Object, VersionedValue> map = openMap((Integer) op[0]);
Object key = op[1];
VersionedValue value = map.get(key);
if (value.value == null) {
map.remove(key); // 删除标记 → 物理删除
} else {
VersionedValue v2 = new VersionedValue();
v2.value = value.value; // operationId 归零 → 全局可见
map.put(key, v2);
}
undoLog.remove(undoKey); // 清理 undo
}
endTransaction(t);
}
回滚从大 logId 向小 logId 倒序遍历 undoLog,用 oldValue 逐级还原:
/**
* @see org.h2.mvstore.db.TransactionStore#rollbackTo
*/
// 遍历本事务的 undoLog,把每个 key 恢复成 oldValue
// 新增的删掉,修改/删除的还原
对于上例的 UPDATE(bjn → bjz):
undoLog[operationId_1] → 还原为 null(op=operationId_0)undoLog[operationId_0] → 还原为 bjn(op=0)逐级逆向走完这条链,data 回到最初状态,如同 UPDATE 从未发生。
trySet,每次对应一个 logId。trySet 遵循 write-ahead 原则:每条数据写入前,先 transaction.log() 把旧值记入 undoLog,保证崩溃后可恢复。VersionedValue#value 恰好是前一步 data 的值,自然形成 bjz → null → bjn 的单向回溯链,getValue 通过遍历此链实现多版本读。trySet 中通过比对 operationId 实现:若目标值属于其他未提交事务即返回 false,上层转为 CONCURRENT_UPDATE_1。在日常业务开发中,经常需要在事务中完成数据库操作后,对外发送 MQ 消息通知下游消费。但在事务提交前就发送消息,消费者可能立即处理并查询数据库,此时事务还未提交,导致数据不一致。
Spring 框架提供了
TransactionSynchronizationAdapter作为最小代价的同步方案,在afterCommit()回调中执行外部交互,确保消息只在事务成功提交后才发送。本文将围绕该机制展开,同时扩展到
TransactionAwareCacheManagerProxy的应用场景,并简要介绍其他同步方案。
假设一个典型的业务场景:创建异常单并通知下游系统。
sequenceDiagram
participant 业务服务
participant 数据库
participant MQ消息队列
participant 消费者
业务服务->>数据库: 1. 开启事务
业务服务->>数据库: 2. 插入业务记录
业务服务->>MQ消息队列: 3. 发送 MQ 消息
业务服务->>数据库: 4. 提交事务
MQ消息队列->>消费者: 5. 消息被立即消费
消费者->>数据库: 6. 查询数据
❌ 消息在事务提交前发出,消费者查不到刚写入的数据,导致处理异常❌ 如果事务最终回滚,消息已经发出且无法撤回,造成数据不一致这是分布式系统中的事务与外部交互一致性问题,核心矛盾在于:
| 时序 | 消息发送时机 | 结果 |
|---|---|---|
| 事务提交前发消息 | 消息已送达,数据未落库 | ❌ 消费者查不到数据 |
| 事务提交后发消息 | 数据已落库,消息后发送 | ✅ 消费者可正常处理 |
🎈 不只 MQ 消息,任何在事务中需要与外部系统交互的场景(缓存更新、文件写入、RPC 调用等)都会面临同样的问题。
Spring 通过 TransactionSynchronizationManager 在事务生命周期中提供多个扩展点,允许注册自定义的同步回调:
| 回调方法 | 触发时机 | 适用场景 |
|---|---|---|
afterCommit() |
事务已提交 | ✅ 发送 MQ、清除缓存 |
afterCompletion(int status) |
事务结束 | 资源清理(区分提交/回滚) |
beforeCommit(boolean readOnly) |
事务提交前 | 最终校验 |
beforeCompletion() |
事务结束前 | 连接释放前的收尾工作 |
⭐ 核心思路:在 afterCommit() 中执行消息发送,利用事务提交后再执行的特性,天然解决消息与数据的时序问题。
sequenceDiagram
participant 业务服务
participant TransactionSynchronizationManager
participant 数据库
participant MQ消息队列
业务服务->>数据库: 1. 开启事务
业务服务->>数据库: 2. 执行业务操作
业务服务->>TransactionSynchronizationManager: 3. registerSynchronization<br/>(注册 afterCommit)
业务服务->>数据库: 4. 提交事务
数据库-->>TransactionSynchronizationManager: 5. 事务提交成功
TransactionSynchronizationManager->>MQ消息队列: 6. 触发 afterCommit()<br/>发送 MQ 消息
✔ 消息发送时数据库操作已全部提交✔ 事务回滚时 afterCommit() 不会被调用,消息不会错误发送核心代码如下,通过在方法中判断当前是否存在活动事务,将消息发送注册到 afterCommit() 回调:
/**
* 发送 MQ 任务消息,确保在事务提交后执行
*
* @see org.springframework.transaction.support.TransactionSynchronizationAdapter#afterCommit
*/
public <T extends MQTaskMessageBaseDto> void sendMQTask(T messageBody) {
if (TransactionSynchronizationManager.isSynchronizationActive()) {
// 存在活动事务 → 注册 afterCommit 回调,延迟到事务提交后发送
TransactionSynchronizationManager.registerSynchronization(
new TransactionSynchronizationAdapter() {
@Override
public void afterCommit() {
log.info("事务已提交,发送MQ消息:{}", messageBody);
doSendMQTask(messageBody);
}
});
} else {
// 不存在事务 → 直接发送
log.info("无事务上下文,直接发送MQ消息:{}", messageBody);
doSendMQTask(messageBody);
}
}
/**
* 实际执行消息发送
*/
private void doSendMQTask(MQTaskMessageBaseDto messageBody) {
// 具体的 MQ 发送逻辑
}
🎈 注意:当 TransactionSynchronizationManager.isSynchronizationActive() 为 false 时(非事务场景),需要直接发送消息,避免回调永远不被触发。
优点:
✔ Spring 原生支持,无需引入额外依赖✔ 实现简单,代码侵入性低✔ 无额外延迟,消息在事务提交后立即发送不足:
❌ 消息发送失败不会影响已提交的事务,需要配合重试或补偿机制❌ 仅适用于单库事务,对分布式事务场景无能为力❌ 如果存在主从延迟,消费者可能从从库读取到旧数据,仍会出现短暂不一致❌ 如果存在网络延迟(消息到达延迟),需要额外的最终一致性保障🙉 对于主从延迟和网络延迟问题,这不是 TransactionSynchronizationAdapter 本身的缺陷,而是分布式系统固有的挑战。简单的 afterCommit() 无法感知这些延迟,需要引入更复杂的方案(如本地消息表 + 定时补偿)。
TransactionAwareCacheManagerProxy 是 Spring 基于同样的同步机制,解决缓存与事务不一致问题的典型案例。
当使用非事务感知的缓存管理器(如 SimpleCacheManager)时,如果先更新缓存再回滚事务:
sequenceDiagram
participant 业务服务
participant 数据库
participant 缓存
业务服务->>数据库: 1. 开启事务
业务服务->>缓存: 2. put(key, value)
业务服务->>数据库: 3. 事务回滚!
Note over 缓存: ❌ 缓存中残留脏数据!<br/>数据库已是旧值
❌ 数据库回滚后数据恢复原值,但缓存中已写入新值,造成脏数据TransactionAwareCacheManagerProxy 基于组合模式,通过 TransactionAwareCacheDecorator 装饰目标 Cache 实例,将缓存写操作与 Spring 事务绑定:
/**
* TransactionAwareCacheDecorator 核心逻辑(伪代码示意)
*
* @see org.springframework.cache.transaction.TransactionAwareCacheDecorator
*/
public class TransactionAwareCacheDecorator implements Cache {
private final Cache targetCache; // 被装饰的真实 Cache
@Override
public void put(Object key, Object value) {
if (TransactionSynchronizationManager.isSynchronizationActive()) {
// 存在活动事务 → 注册 afterCommit 回调
TransactionSynchronizationManager.registerSynchronization(
new TransactionSynchronizationAdapter() {
@Override
public void afterCommit() {
targetCache.put(key, value); // 事务提交后真正写入缓存
}
});
} else {
targetCache.put(key, value); // 非事务场景直接写入
}
}
// ❌ evict、clear 等写操作同样延迟到 afterCommit
}
核心行为:
| 操作 | 存在活动事务时 | 无事务时 |
|---|---|---|
put(key, value) |
注册到 afterCommit,事务提交后才写入 |
立即写入 |
evict(key) |
注册到 afterCommit,事务提交后才清除 |
立即清除 |
clear() |
注册到 afterCommit,事务提交后才清空 |
立即清空 |
get(key) |
直接读取(读操作不受事务影响) | 直接读取 |
✔ 事务回滚 → afterCommit() 不会被调用 → 缓存保持旧值,与数据库一致✔ 事务提交 → afterCommit() 被调用 → 缓存写入新值,与数据库一致通过 TransactionAwareCacheManagerProxy 包装已有的 CacheManager:
@Bean
public CacheManager cacheManager(CacheManager targetCacheManager) {
// 包装目标 CacheManager,使其具备事务感知能力
return new TransactionAwareCacheManagerProxy(targetCacheManager);
}
🎈 如果使用了 @EnableCaching 且配置了 JCacheCacheManager、ConcurrentMapCacheManager 等非事务感知的实现,建议通过这种方式增强事务一致性。
除了 TransactionSynchronizationAdapter 之外,针对事务与外部交互的同步问题,还有以下可选方案:
| 方案 | 核心思路 | 适用场景 | 复杂度 |
|---|---|---|---|
| 本地消息表 | 将消息写入本地数据库表,定时扫描发送 + 确认消费 | 跨服务、需要最终一致性 | ⭐⭐⭐ |
| RocketMQ 事务消息 | 利用 RocketMQ 事务消息机制,回查本地事务状态 | 已使用 RocketMQ 的场景 | ⭐⭐⭐ |
| 定时扫描补偿 | 定时扫描待发送/失败记录,进行补发 | 简单场景、对时效性要求不高 | ⭐⭐ |
| 事件溯源 (Event Sourcing) | 以事件为事实源,状态从事件派生 | 复杂的领域模型 | ⭐⭐⭐⭐ |
| TCC 分布式事务 | Try-Confirm-Cancel 三阶段协议 | 严格的跨服务数据一致性 | ⭐⭐⭐⭐⭐ |
👉 方案选型建议:优先从最简单的方式开始。如果只是单库事务中的消息发送或缓存更新,TransactionSynchronizationAdapter 是最低成本的选择;当面临跨数据库、跨服务场景时,再考虑本地消息表或 RocketMQ 事务消息等方案。
TransactionSynchronizationAdapter 是 Spring 提供的最小代价方案,通过 afterCommit() 将外部交互延迟到事务成功提交后,解决消息发送与数据不一致问题TransactionAwareCacheManagerProxy 也是基于同样的同步机制,通过装饰模式将缓存写操作绑定到事务生命周期,避免事务回滚后的脏数据在深入阅读 H2 数据库源码时,会发现内部存在两套完全独立的存储实现——
MVStore和PageStore。初次接触时容易混淆:同名类
FileStore为何分布在不同包下?MVMap和Page又是什么关系?本文通过层级对照的方式,将两条并行的存储栈拉平对比,帮助建立起清晰的结构认知。
MVStore为默认引擎,采用MVCC机制;Regular对应传统的PageStore栈。
Database 在启动时根据配置决定走哪条存储栈,MVStore 为当前默认选项:
// @see org.h2.engine.Database#open
// 依据配置 MV_STORE 决定使用哪套引擎
if (dbSettings.mvStore) {
// 走 MVStore 栈(新)
} else {
// 走 PageStore 栈(旧)
}
🎈 两套栈互不混用,表引擎、存储层、文件 I/O 全是独立实现。理解这一点,就不会再把 org.h2.mvstore.FileStore 和 org.h2.store.FileStore 搞混了。
graph TB
DB["Database 数据库引擎"]
DB -->|MVCC/默认| MVStack["MVStore 栈"]
DB -->|传统| PSStack["PageStore 栈"]
subgraph MVStack["MVStore 存储栈 (新)"]
MVTable["MVTable extends TableBase"]
MVTable --> MVStore["MVStore 存储引擎"]
MVTable --> TxStore["TransactionStore 事务层"]
MVStore --> MVMap["MVMap 有序 K-V 树"]
MVStore --> MVFileStore["mvstore.FileStore"]
MVMap --> MVPage["mvstore.Page"]
end
subgraph PSStack["PageStore 存储栈 (旧)"]
RegularTable["RegularTable extends TableBase"]
RegularTable --> PageStore["PageStore 存储引擎"]
PageStore --> PSFileStore["store.FileStore"]
PageStore --> PSPage["store.Page 固定页"]
end
⚠️ 关键点:
org.h2.mvstore.FileStore和org.h2.store.FileStore是两个不同的类,分属两套栈,互不相关。
// MVMap 属于 MVStore
public class MVMap<K, V> {
MVStore store;
}
// MVStore 管理多个 MVMap
public class MVStore {
ConcurrentHashMap<Integer, MVMap<?, ?>> maps;
FileStore fileStore; // 底层 I/O
MVMap<String, String> meta; // 元信息 map
}
// MVTable 持有 TransactionStore
public class MVTable extends TableBase {
TransactionStore store;
// 主键/索引数据落在若干 MVMap 上
}
核心关系:
MVMap.store → 每个 MVMap 属于一个 MVStoreMVStore.maps → 一个 MVStore 管理多个 MVMapMVStore.fileStore → MVStore 持有一个 mvstore.FileStore 做底层 I/OMVStore.meta → 特殊的 MVMap<String,String> 记录所有 map 的元信息MVTable.store → MVTable 持有 TransactionStore,后者包装 MVStore// RegularTable 的索引持有 PageStore
public class PageDataIndex extends Index {
PageStore store;
}
public class PageStore {
FileStore file; // store.FileStore,与 mvstore.FileStore 不同
}
// Database 整库共享一个 PageStore 实例
public class Database {
PageStore pageStore;
}
核心关系:
RegularTable 的索引(PageDataIndex、PageBtreeIndex)持有 PageStorePageStore.file → PageStore 持有一个 store.FileStoreDatabase.pageStore → 整库共享一个 PageStore 实例graph LR
subgraph L1["表层"]
A1["RegularTable"]
B1["MVTable"]
end
subgraph L2["事务层"]
B2["TransactionStore"]
end
subgraph L3["存储引擎层"]
A3["PageStore"]
B3["MVStore"]
end
subgraph L4["逻辑数据结构层"]
A4["store.Page 固定页 B-Tree"]
B4["MVMap + mvstore.Page"]
end
subgraph L5["文件 I/O 层"]
A5["store.FileStore"]
B5["mvstore.FileStore"]
end
A1 --> A3 --> A4 --> A5
B1 --> B2 --> B3 --> B4 --> B5
RegularTable vs MVTable:均继承 TableBase,对应传统表与 MVCC 表。MVTable 通过 TransactionStore 支持 MVCC;RegularTable 直接耦合在 PageStore 中,没有独立事务层。PageStore 和 MVStore 是两套核心引擎,分别服务各自的上层表。PageStore 侧使用 store.Page 固定页 B-Tree;MVStore 侧使用 MVMap(有序 K-V 树),内部节点为 mvstore.Page。store.FileStore 和 mvstore.FileStore 同名不同类,各自只服务本栈。| 维度 | MVStore 栈 | PageStore 栈 |
|---|---|---|
| 默认启用 | ✔ 默认 | ❌ 需显式配置 |
| 事务支持 | MVCC 多版本并发 |
默认传统锁机制(也可开启 MVCC) |
| 中间层 | MVMap(多 K-V 树容器) |
无,直接用固定页 B-Tree |
| 文件 I/O | mvstore.FileStore |
store.FileStore |
| 表实现 | MVTable |
RegularTable |
🎈 最关键的区别:MVMap 是 MVStore 独有的中间层。MVStore 是“多个 MVMap 的容器”,每张 MVTable 的数据/索引就是若干 MVMap;PageStore 侧没有这层,直接用固定页 B-Tree 索引。
🎈 RegularTable 有完整的 MVCC 事务语义,包括读己之写、隐藏其他事务未提交数据、写写冲突检测、提交与回滚;只是底层仍受 PageStore 全局同步与内存 undoLog 的限制,不如 MVTable 的 MVStore 原生 MVCC。
RegularTable→PageStore→store.FileStoreMVTable→(TransactionStore)→MVStore→mvstore.FileStoreMVMap 是 MVStore 独有的中间层,理解这点就能区分两套引擎的核心设计差异FileStore 同名不同类,各自只服务本栈,是最底层的文件读写封装业务操作复杂后,发送 MQ 消息和数据库事务往往存在时序错位。事务尚未提交,消息已被消费者处理,导致数据查询失败。
本文的核心对象是 Spring 的
TransactionSynchronizationAdapter,通过注册事务提交回调,以最小代价实现”事务成功后才发送消息”的同步保障。关键结论:
TransactionSynchronization是单体应用内解决事务与外部交互一致性的轻量级方案,但若涉及主从延迟或网络延迟,则需要更复杂的分布式改造。
当业务操作(如创建订单、异常单)和发送 MQ 任务消息在同一事务中执行时,如果在事务提交前就发送消息,消息消费者可能立即处理并查询数据库,此时事务还未提交,导致数据不一致。
/**
* 问题代码示意:事务与消息发送时序错误
*/
@Transactional
public void createOrder(Order order) {
orderMapper.insert(order); // 1. 数据库操作
mqProducer.send(orderMessage); // 2. 事务未提交,消息已发送 ❌
// 3. 其他耗时操作(业务逻辑、外部调用等)。消息消费者查询不到新增的数据 ❌
// 4. 事务提交在此之后
}
🙉 消费者收到消息后查询数据库,可能查不到刚插入的数据——因为事务还未提交。
核心矛盾在于事务边界与外部交互边界的错位:
| 阶段 | 消息发送时机 | 消费者查询结果 |
|---|---|---|
| 事务提交前发消息 | 消息已发送,但数据未提交 | ❌ 查不到数据 |
| 事务提交后发消息 | 数据已提交,消息后发送 | ✅ 正常查询 |
🎈 在单体应用、同库事务场景下,理想的方案是:让消息发送这个”外部交互”延迟到事务真正提交之后。
Spring 提供了 TransactionSynchronization 接口,允许将自定义逻辑注册到事务生命周期中。其中 afterCommit() 回调会在事务成功提交后触发,afterCompletion() 则在事务完成(无论成功或失败)后触发。
🎈 选择 afterCommit() 而非 afterCompletion() 的原因是:只有事务成功提交后才需要发送消息;事务回滚时,数据未落库,不应发送消息。
/**
* 事务提交后才发送 MQ 消息
* @see org.springframework.transaction.support.TransactionSynchronizationAdapter
*/
public <T extends MQTaskMessageBaseDto> void sendMQTask(T messageBody) {
if (TransactionSynchronizationManager.isSynchronizationActive()) {
log.info("sendMQTask--事务提交后再发送消息:{}", messageBody);
// 注册事务同步回调,延迟到事务提交后执行
TransactionSynchronizationManager.registerSynchronization(
new TransactionSynchronizationAdapter() {
@Override
public void afterCommit() {
doSendMQTask(messageBody); // 事务已提交,安全发送
}
});
} else {
log.info("sendMQTask--直接发送消息:{}", messageBody);
doSendMQTask(messageBody); // 无事务上下文,直接发送
}
}
sequenceDiagram
participant 业务服务
participant 数据库
participant MQ消息队列
业务服务->>数据库: 1. 开启事务
业务服务->>数据库: 2. 执行数据库操作
业务服务->>业务服务: 3. registerSynchronization<br/>(afterCommit 回调)
业务服务->>业务服务: 4. 其他耗时操作<br/>(业务逻辑、外部调用等)
业务服务->>数据库: 5. 提交事务
数据库-->>业务服务: 6. 事务提交成功
业务服务->>MQ消息队列: 7. afterCommit() 触发<br/>发送 MQ 任务消息
Note over MQ消息队列: ✅ 消息发送时数据已提交
✔ 这种方案的优势在于零侵入——业务代码不需要关心事务状态,统一封装在消息发送层。
⭐ 优点:复杂度低,Spring 原生支持; 实现难度低,代码量少
🎈 不足:如果涉及主从延迟(消息发送后消费者读从库,数据尚未同步)或网络延迟(消息先于事务提交传播到消费者),则需要更复杂的分布式改造方案(如本地消息表、事务消息等)。
TransactionSynchronizationAdapter 的应用不仅限于 MQ 消息发送。Spring Cache 模块中,TransactionAwareCacheManagerProxy 正是基于同样的思想,为缓存操作添加事务感知能力。
某些缓存管理器(如 SimpleCacheManager)不支持直接的事务感知。如果在事务中先更新缓存、后回滚事务,缓存中就会残留脏数据:
@Transactional
public void updateData(Data data) {
cache.put(data.getId(), data); // 1. 缓存已更新
dataMapper.update(data); // 2. 数据库操作
// 3. 抛出异常,事务回滚
// 结果:缓存是新数据,数据库是旧数据 ❌
}
TransactionAwareCacheManagerProxy 基于组合模式,为不支持事务感知的 CacheManager 添加事务同步能力。它本质上是 TransactionAwareCacheDecorator 装饰目标 Cache 实例,将缓存写操作与 Spring 管理的事务绑定。
/**
* 事务感知缓存装饰器
* @see org.springframework.cache.transaction.TransactionAwareCacheDecorator
*/
public class TransactionAwareCacheDecorator implements Cache {
private final Cache targetCache;
@Override
public void put(final Object key, final Object value) {
if (TransactionSynchronizationManager.isSynchronizationActive()) {
// 存在活动事务时,延迟缓存写操作到事务提交后
TransactionSynchronizationManager.registerSynchronization(
new TransactionSynchronizationAdapter() {
@Override
public void afterCommit() {
targetCache.put(key, value); // 事务提交后真正写入缓存
}
});
} else {
targetCache.put(key, value); // 无事务,直接写入
}
}
// evict、clear 等操作同理延迟执行
}
🎈 核心逻辑与 MQ 发送方案完全一致:存在活动事务时,put、evict、clear 等操作被注册到 TransactionSynchronizationManager,延迟到事务成功提交后才真正执行。若事务回滚,这些缓存操作自然失效,不会污染缓存。
✔ 这种装饰器模式的设计非常优雅:不修改原有缓存管理器的实现,仅通过代理层增强事务感知能力,符合开闭原则。
在分布式场景或更高可靠性要求下,还有以下方案可供选择:
| 方案 | 复杂度 | 可靠性 | 实现难度 | 适用场景 |
|---|---|---|---|---|
| 本地消息表 | ⭐⭐⭐ | 高 | 中 | 跨服务 / 跨库场景,需持久化消息 |
| RocketMQ 事务消息 | ⭐⭐⭐ | 高 | 中 | 已使用 RocketMQ,两阶段提交 |
| 定时扫描补偿 | ⭐⭐ | 中 | 低 | 简单场景,允许最终一致性延迟 |
| Seata 分布式事务 | ⭐⭐⭐⭐ | 高 | 高 | 强一致性要求,改造代价大 |
🎈 选型建议:单体应用优先使用 TransactionSynchronization(本文方案);跨服务、跨库场景考虑本地消息表或 RocketMQ 事务消息;对一致性要求极高且可接受高改造成本时,再考虑分布式事务框架。
TransactionSynchronizationAdapter 是 Spring 原生的事务生命周期钩子,以最小代码量实现”事务提交后执行外部操作”的同步保障。afterCommit 回调,绑定到事务成功提交后才真正执行。TransactionAwareCacheManagerProxy 是同一思想的经典应用,通过装饰器模式为缓存添加事务感知,避免事务回滚后的脏数据问题。