一、核心架构与设计思想
1. 事务同步机制全景图
图表
代码
复制
下载
全屏
graph TB
subgraph "Spring事务管理器"
TM[PlatformTransactionManager]
TS[TransactionSynchronizationManager]
TSR[TransactionSynchronizationRegistry]
end
subgraph "事务资源"
DS1[DataSource 1]
DS2[DataSource 2]
DS3[DataSource 3]
end
subgraph "同步处理器"
SYNC1[TransactionSynchronization 1]
SYNC2[TransactionSynchronization 2]
SYNC3[TransactionSynchronization 3]
end
TM –>|管理| TS
TS –>|注册| TSR
TS –>|绑定资源| DS1
TS –>|绑定资源| DS2
TS –>|绑定资源| DS3
TS –>|注册回调| SYNC1
TS –>|注册回调| SYNC2
TS –>|注册回调| SYNC3
style TS fill:#e1f5e1
2. 线程绑定的事务上下文
java
复制
下载
/**
* TransactionSynchronizationManager 核心数据结构
* 基于ThreadLocal实现线程隔离的事务上下文
*/
public abstract class TransactionSynchronizationManager {
// 核心ThreadLocal变量
private static final ThreadLocal<Map<Object, Object>> resources =
new NamedThreadLocal<>("Transactional resources");
private static final ThreadLocal<Set<TransactionSynchronization>> synchronizations =
new NamedThreadLocal<>("Transaction synchronizations");
private static final ThreadLocal<String> currentTransactionName =
new NamedThreadLocal<>("Current transaction name");
private static final ThreadLocal<Boolean> currentTransactionReadOnly =
new NamedThreadLocal<>("Current transaction read-only status");
private static final ThreadLocal<Integer> currentTransactionIsolationLevel =
new NamedThreadLocal<>("Current transaction isolation level");
private static final ThreadLocal<Boolean> actualTransactionActive =
new NamedThreadLocal<>("Actual transaction active");
/**
* 事务同步注册表 – 管理所有同步回调
*/
private static class SynchronizationRegistry {
// 按顺序执行的同步器
private List<TransactionSynchronization> orderedSynchronizations;
// 按优先级分组的同步器
private Map<Integer, List<TransactionSynchronization>> prioritizedSyncs;
// 同步状态
private boolean synchronizationActive = false;
}
}
二、核心实现机制
1. 完整的资源绑定管理
java
复制
下载
/**
* 事务资源管理器 – 管理所有绑定到当前线程的资源
*/
public class TransactionResourceManager {
/**
* 绑定资源到当前事务
*/
public static void bindResource(Object key, Object value) {
Assert.notNull(key, "Key must not be null");
Assert.notNull(value, "Value must not be null");
Map<Object, Object> map = resources.get();
// 如果当前线程没有资源映射,创建一个
if (map == null) {
map = new HashMap<>();
resources.set(map);
}
Object oldValue = map.put(key, value);
// 检查资源是否已绑定
if (oldValue != null) {
throw new IllegalStateException(
"Already value [" + oldValue + "] for key [" +
key + "] bound to thread [" + Thread.currentThread().getName() + "]");
}
if (logger.isTraceEnabled()) {
logger.trace("Bound value [" + value + "] for key [" + key + "] to thread");
}
}
/**
* 解除资源绑定
*/
public static Object unbindResource(Object key) throws IllegalStateException {
Object value = doUnbindResource(key);
if (value == null) {
throw new IllegalStateException(
"No value for key [" + key + "] bound to thread");
}
return value;
}
private static Object doUnbindResource(Object key) {
Map<Object, Object> map = resources.get();
if (map == null) {
return null;
}
Object value = map.remove(key);
// 如果map为空,清理ThreadLocal
if (map.isEmpty()) {
resources.remove();
}
if (value != null && logger.isTraceEnabled()) {
logger.trace("Removed value [" + value + "] for key [" + key + "] from thread");
}
return value;
}
/**
* 获取绑定的资源
*/
public static Object getResource(Object key) {
Map<Object, Object> map = resources.get();
if (map == null) {
return null;
}
Object value = map.get(key);
if (value instanceof ResourceHolder) {
ResourceHolder holder = (ResourceHolder) value;
if (holder.isVoid()) {
map.remove(key);
if (map.isEmpty()) {
resources.remove();
}
value = null;
}
}
return value;
}
/**
* 资源持有器 – 封装真实资源
*/
public interface ResourceHolder {
void released();
boolean isVoid();
boolean shouldReleaseBeforeCompletion();
void setSynchronizedWithTransaction(boolean synched);
boolean isSynchronizedWithTransaction();
}
/**
* 数据源资源持有器
*/
public static class DataSourceResourceHolder implements ResourceHolder {
private final DataSource dataSource;
private Connection connection;
private boolean synchronizedWithTransaction = false;
private boolean released = false;
public DataSourceResourceHolder(DataSource dataSource) {
this.dataSource = dataSource;
}
public Connection getConnection() throws SQLException {
if (connection == null || connection.isClosed()) {
connection = dataSource.getConnection();
}
return connection;
}
@Override
public void released() {
this.released = true;
if (connection != null) {
try {
if (!connection.getAutoCommit()) {
connection.rollback();
}
connection.close();
} catch (SQLException e) {
logger.debug("Could not close JDBC connection", e);
}
}
}
@Override
public boolean isVoid() {
return released;
}
}
}
2. 事务同步回调机制
java
复制
下载
/**
* 事务同步回调管理器
* 管理事务生命周期中的各种回调
*/
public class TransactionSynchronizationManager {
/**
* 事务同步回调接口
*/
public interface TransactionSynchronization extends Ordered {
/** 事务挂起时调用 */
void suspend();
/** 事务恢复时调用 */
void resume();
/** 刷新会话状态到数据库 */
void flush();
/** 提交前调用(在commit之前) */
void beforeCommit(boolean readOnly);
/** 提交前调用(在commit之后,beforeCompletion之前) */
void beforeCompletion();
/** 提交后调用(在commit之后) */
void afterCommit();
/** 事务完成后调用(无论提交还是回滚) */
void afterCompletion(int status);
/**
* 执行顺序
* @return 优先级,值越小优先级越高
*/
@Override
default int getOrder() {
return Ordered.LOWEST_PRECEDENCE;
}
}
/**
* 注册事务同步器
*/
public static void registerSynchronization(TransactionSynchronization synchronization)
throws IllegalStateException {
Assert.notNull(synchronization, "TransactionSynchronization must not be null");
Set<TransactionSynchronization> syncs = synchronizations.get();
if (syncs == null) {
throw new IllegalStateException(
"Transaction synchronization is not active");
}
syncs.add(synchronization);
if (logger.isDebugEnabled()) {
logger.debug("Registered synchronization: " + synchronization);
}
}
/**
* 获取所有同步器(按优先级排序)
*/
public static List<TransactionSynchronization> getSynchronizations()
throws IllegalStateException {
Set<TransactionSynchronization> syncs = synchronizations.get();
if (syncs == null) {
throw new IllegalStateException(
"Transaction synchronization is not active");
}
// 如果只有一个同步器,直接返回
if (syncs.size() == 1) {
return new ArrayList<>(syncs);
}
// 按优先级排序
List<TransactionSynchronization> sorted = new ArrayList<>(syncs);
AnnotationAwareOrderComparator.sort(sorted);
return sorted;
}
/**
* 触发提交前回调
*/
public static void triggerBeforeCommit(boolean readOnly) {
for (TransactionSynchronization synchronization : getSynchronizations()) {
try {
synchronization.beforeCommit(readOnly);
} catch (Throwable t) {
logger.error("TransactionSynchronization.beforeCommit threw exception", t);
}
}
}
/**
* 触发完成前回调
*/
public static void triggerBeforeCompletion() {
for (TransactionSynchronization synchronization : getSynchronizations()) {
try {
synchronization.beforeCompletion();
} catch (Throwable t) {
logger.error("TransactionSynchronization.beforeCompletion threw exception", t);
}
}
}
/**
* 触发提交后回调
*/
public static void triggerAfterCommit() {
invokeAfterCommit(getSynchronizations());
}
private static void invokeAfterCommit(List<TransactionSynchronization> synchronizations) {
if (synchronizations != null) {
for (TransactionSynchronization synchronization : synchronizations) {
try {
synchronization.afterCommit();
} catch (Throwable t) {
logger.error("TransactionSynchronization.afterCommit threw exception", t);
}
}
}
}
/**
* 触发完成后回调
*/
public static void triggerAfterCompletion(int status) {
List<TransactionSynchronization> synchronizations = getSynchronizations();
invokeAfterCompletion(synchronizations, status);
}
private static void invokeAfterCompletion(
List<TransactionSynchronization> synchronizations, int status) {
if (synchronizations != null) {
for (TransactionSynchronization synchronization : synchronizations) {
try {
synchronization.afterCompletion(status);
} catch (Throwable t) {
logger.error("TransactionSynchronization.afterCompletion threw exception", t);
}
}
}
}
}
篇幅限制下面就只能给大家展示小册部分内容了。整理了一份核心面试笔记包括了:Java面试、Spring、JVM、MyBatis、Redis、MySQL、并发编程、微服务、Linux、Springboot、SpringCloud、MQ、Kafc
需要全套面试笔记及答案 【点击此处即可/免费获取】
三、完整的事务同步工作流
1. 事务同步生命周期
java
复制
下载
/**
* 事务同步生命周期管理器
*/
public class TransactionSynchronizationLifecycle {
/**
* 开启事务同步
*/
public static void initSynchronization() throws IllegalStateException {
if (isSynchronizationActive()) {
throw new IllegalStateException(
"Cannot activate transaction synchronization – already active");
}
logger.trace("Initializing transaction synchronization");
synchronizations.set(new LinkedHashSet<>());
}
/**
* 清理事务同步
*/
public static void clearSynchronization() throws IllegalStateException {
if (!isSynchronizationActive()) {
throw new IllegalStateException(
"Cannot deactivate transaction synchronization – not active");
}
logger.trace("Clearing transaction synchronization");
synchronizations.remove();
}
/**
* 挂起事务同步
*/
public static TransactionSynchronizationState suspend() {
if (!isSynchronizationActive()) {
return null;
}
Set<TransactionSynchronization> suspendedSyncs = synchronizations.get();
// 触发挂起回调
for (TransactionSynchronization sync : suspendedSyncs) {
sync.suspend();
}
// 保存当前状态
TransactionSynchronizationState state = new TransactionSynchronizationState(
currentTransactionName.get(),
currentTransactionReadOnly.get(),
currentTransactionIsolationLevel.get(),
actualTransactionActive.get()
);
// 清理ThreadLocal
clear();
return state;
}
/**
* 恢复事务同步
*/
public static void resume(TransactionSynchronizationState state) {
if (state == null) {
return;
}
// 恢复事务状态
if (state.getTransactionName() != null) {
currentTransactionName.set(state.getTransactionName());
}
if (state.getReadOnly() != null) {
currentTransactionReadOnly.set(state.getReadOnly());
}
if (state.getIsolationLevel() != null) {
currentTransactionIsolationLevel.set(state.getIsolationLevel());
}
if (state.getActive() != null) {
actualTransactionActive.set(state.getActive());
}
// 触发恢复回调
Set<TransactionSynchronization> syncs = synchronizations.get();
if (syncs != null) {
for (TransactionSynchronization sync : syncs) {
sync.resume();
}
}
}
/**
* 完整的事务提交流程
*/
public static void processCommit(DefaultTransactionStatus status) throws TransactionException {
try {
// 1. 触发提交前回调
triggerBeforeCommit(status.isReadOnly());
// 2. 如果有保存点,释放保存点
if (status.hasSavepoint()) {
if (status.isDebug()) {
logger.debug("Releasing transaction savepoint");
}
status.releaseHeldSavepoint();
}
// 3. 如果是新事务,提交
else if (status.isNewTransaction()) {
if (status.isDebug()) {
logger.debug("Initiating transaction commit");
}
doCommit(status);
}
// 4. 如果有全局回滚标记,设置回滚
else if (status.hasTransaction()) {
if (isGlobalRollbackOnly()) {
if (status.isDebug()) {
logger.debug("Participating transaction failed – marking existing transaction as rollback-only");
}
doSetRollbackOnly(status);
} else {
if (status.isDebug()) {
logger.debug("Participating in existing transaction");
}
}
}
// 5. 触发完成前回调
triggerBeforeCompletion();
// 6. 触发提交后回调(仅在提交成功时)
triggerAfterCommit();
} catch (Exception ex) {
// 7. 处理提交失败
doRollbackOnCommitException(status, ex);
throw ex;
} finally {
// 8. 清理资源
cleanupAfterCompletion(status);
}
}
/**
* 完整的事务回滚流程
*/
public static void processRollback(DefaultTransactionStatus status) throws TransactionException {
try {
// 1. 如果有保存点,回滚到保存点
if (status.hasSavepoint()) {
if (status.isDebug()) {
logger.debug("Rolling back transaction to savepoint");
}
status.rollbackToHeldSavepoint();
}
// 2. 如果是新事务,回滚
else if (status.isNewTransaction()) {
if (status.isDebug()) {
logger.debug("Initiating transaction rollback");
}
doRollback(status);
}
// 3. 参与现有事务,标记为回滚
else if (status.hasTransaction()) {
if (status.isDebug()) {
logger.debug("Participating in existing transaction – marking as rollback-only");
}
doSetRollbackOnly(status);
}
// 4. 触发完成前回调
triggerBeforeCompletion();
} catch (Exception ex) {
triggerAfterCompletion(TransactionSynchronization.STATUS_UNKNOWN);
throw ex;
} finally {
// 5. 清理资源
cleanupAfterCompletion(status);
}
}
/**
* 提交后清理资源
*/
private static void cleanupAfterCompletion(DefaultTransactionStatus status) {
// 设置完成状态
status.setCompleted();
// 清理同步
if (status.isNewSynchronization()) {
clearSynchronization();
}
// 清理资源
if (status.isNewTransaction()) {
doCleanupAfterCompletion(status.getTransaction());
}
// 如果有挂起的资源,恢复
if (status.getSuspendedResources() != null) {
Object transaction = status.getTransaction();
if (transaction != null) {
resumeAfterCompletion(transaction, (SuspendedResourcesHolder) status.getSuspendedResources());
}
}
// 触发完成后回调
int completionStatus = status.isRollbackOnly() ?
TransactionSynchronization.STATUS_ROLLED_BACK :
TransactionSynchronization.STATUS_COMMITTED;
triggerAfterCompletion(completionStatus);
}
}
2. 嵌套事务与挂起恢复
java
复制
下载
/**
* 嵌套事务与事务挂起管理器
*/
public class NestedTransactionManager {
/**
* 挂起当前事务
*/
public static SuspendedResourcesHolder suspend(@Nullable Object transaction)
throws TransactionException {
if (TransactionSynchronizationManager.isSynchronizationActive()) {
// 触发同步器的挂起回调
List<TransactionSynchronization> suspendedSynchronizations =
doSuspendSynchronization();
try {
Object suspendedResources = null;
if (transaction != null) {
// 挂起事务资源
suspendedResources = doSuspend(transaction);
}
// 获取当前事务状态
String name = TransactionSynchronizationManager.getCurrentTransactionName();
boolean readOnly = TransactionSynchronizationManager.isCurrentTransactionReadOnly();
Integer isolationLevel = TransactionSynchronizationManager.getCurrentTransactionIsolationLevel();
boolean wasActive = TransactionSynchronizationManager.isActualTransactionActive();
// 清理当前事务状态
TransactionSynchronizationManager.clear();
// 创建挂起资源持有器
return new SuspendedResourcesHolder(
suspendedResources,
suspendedSynchronizations,
name,
readOnly,
isolationLevel,
wasActive
);
} catch (RuntimeException | Error ex) {
// 挂起失败,恢复同步
doResumeSynchronization(suspendedSynchronizations);
throw ex;
}
} else if (transaction != null) {
// 没有同步激活,只挂起事务
Object suspendedResources = doSuspend(transaction);
return new SuspendedResourcesHolder(suspendedResources);
} else {
// 既没有同步也没有事务
return null;
}
}
/**
* 恢复挂起的事务
*/
public static void resume(@Nullable Object transaction,
@Nullable SuspendedResourcesHolder resourcesHolder)
throws TransactionException {
if (resourcesHolder != null) {
Object suspendedResources = resourcesHolder.getSuspendedResources();
if (suspendedResources != null) {
doResume(transaction, suspendedResources);
}
List<TransactionSynchronization> suspendedSynchronizations =
resourcesHolder.getSuspendedSynchronizations();
if (suspendedSynchronizations != null) {
// 恢复事务状态
TransactionSynchronizationManager.setActualTransactionActive(
resourcesHolder.wasActive());
TransactionSynchronizationManager.setCurrentTransactionName(
resourcesHolder.getName());
TransactionSynchronizationManager.setCurrentTransactionReadOnly(
resourcesHolder.isReadOnly());
TransactionSynchronizationManager.setCurrentTransactionIsolationLevel(
resourcesHolder.getIsolationLevel());
// 恢复同步
doResumeSynchronization(suspendedSynchronizations);
}
}
}
/**
* 挂起同步器
*/
private static List<TransactionSynchronization> doSuspendSynchronization() {
List<TransactionSynchronization> suspendedSynchronizations =
TransactionSynchronizationManager.getSynchronizations();
for (TransactionSynchronization synchronization : suspendedSynchronizations) {
synchronization.suspend();
}
TransactionSynchronizationManager.clearSynchronization();
return suspendedSynchronizations;
}
/**
* 恢复同步器
*/
private static void doResumeSynchronization(List<TransactionSynchronization> synchronizations) {
TransactionSynchronizationManager.initSynchronization();
for (TransactionSynchronization synchronization : synchronizations) {
TransactionSynchronizationManager.registerSynchronization(synchronization);
synchronization.resume();
}
}
/**
* 挂起资源持有器
*/
public static class SuspendedResourcesHolder {
private final Object suspendedResources;
private final List<TransactionSynchronization> suspendedSynchronizations;
private final String name;
private final boolean readOnly;
private final Integer isolationLevel;
private final boolean wasActive;
public SuspendedResourcesHolder(Object suspendedResources) {
this(suspendedResources, null, null, false, null, false);
}
public SuspendedResourcesHolder(
Object suspendedResources,
List<TransactionSynchronization> suspendedSynchronizations,
String name,
boolean readOnly,
Integer isolationLevel,
boolean wasActive) {
this.suspendedResources = suspendedResources;
this.suspendedSynchronizations = suspendedSynchronizations;
this.name = name;
this.readOnly = readOnly;
this.isolationLevel = isolationLevel;
this.wasActive = wasActive;
}
}
}
四、高级应用场景
1. 自定义事务同步器实现
java
复制
下载
/**
* 自定义事务同步器 – 完整的生产级实现
*/
@Component
public class CustomTransactionSynchronization
implements TransactionSynchronization, Ordered {
private final ThreadLocal<TransactionContext> contextHolder =
new NamedThreadLocal<>("Transaction Context");
private final AuditLogger auditLogger;
private final CacheManager cacheManager;
private final MessageQueuePublisher messagePublisher;
// 事务上下文
@Data
public static class TransactionContext {
private String transactionId;
private long startTime;
private List<OperationLog> operationLogs = new ArrayList<>();
private Map<String, Object> metadata = new HashMap<>();
public void logOperation(String operation, Object… args) {
OperationLog log = new OperationLog();
log.setTimestamp(System.currentTimeMillis());
log.setOperation(operation);
log.setArgs(args);
operationLogs.add(log);
}
}
@Override
public int getOrder() {
// 设置为高优先级,确保在其他同步器之前执行
return Ordered.HIGHEST_PRECEDENCE;
}
@Override
public void suspend() {
TransactionContext context = contextHolder.get();
if (context != null) {
context.logOperation("TRANSACTION_SUSPENDED");
auditLogger.logAudit("transaction_suspended",
Map.of("transactionId", context.getTransactionId()));
}
}
@Override
public void resume() {
TransactionContext context = contextHolder.get();
if (context != null) {
context.logOperation("TRANSACTION_RESUMED");
auditLogger.logAudit("transaction_resumed",
Map.of("transactionId", context.getTransactionId()));
}
}
@Override
public void flush() {
// 在Hibernate等ORM框架刷新会话时调用
TransactionContext context = contextHolder.get();
if (context != null) {
context.logOperation("FLUSH_OPERATION");
logger.debug("Flushing transaction context: {}", context.getTransactionId());
}
}
@Override
public void beforeCommit(boolean readOnly) {
TransactionContext context = contextHolder.get();
if (readOnly) {
logger.debug("Read-only transaction, skipping beforeCommit processing");
return;
}
if (context != null) {
context.logOperation("BEFORE_COMMIT", readOnly);
try {
// 1. 验证业务规则
validateBusinessRules();
// 2. 准备缓存更新
prepareCacheUpdates();
// 3. 生成审计日志
generateAuditLogs(context);
// 4. 准备消息队列消息
prepareEventMessages(context);
} catch (Exception e) {
logger.error("beforeCommit failed", e);
throw new TransactionException("Pre-commit validation failed", e);
}
}
}
@Override
public void beforeCompletion() {
TransactionContext context = contextHolder.get();
if (context != null) {
context.logOperation("BEFORE_COMPLETION");
try {
// 执行最后检查
performFinalChecks();
// 更新事务统计
updateTransactionStats(context);
} catch (Exception e) {
logger.warn("beforeCompletion failed, but transaction will continue", e);
}
}
}
@Override
public void afterCommit() {
TransactionContext context = contextHolder.get();
if (context != null) {
context.logOperation("AFTER_COMMIT");
try {
// 1. 异步更新缓存
updateCachesAsync(context);
// 2. 发送消息队列事件
publishEventsAsync(context);
// 3. 记录成功审计日志
auditLogger.logSuccess("transaction_committed",
Map.of(
"transactionId", context.getTransactionId(),
"duration", System.currentTimeMillis() – context.getStartTime(),
"operations", context.getOperationLogs().size()
));
// 4. 触发后续处理
triggerPostCommitActions(context);
} catch (Exception e) {
// afterCommit中的异常不应该回滚事务
logger.error("afterCommit processing failed", e);
handlePostCommitFailure(context, e);
}
}
}
@Override
public void afterCompletion(int status) {
TransactionContext context = contextHolder.get();
try {
if (context != null) {
context.logOperation("AFTER_COMPLETION", status);
// 根据状态处理
switch (status) {
case STATUS_COMMITTED:
handleCommittedTransaction(context);
break;
case STATUS_ROLLED_BACK:
handleRolledBackTransaction(context);
break;
case STATUS_UNKNOWN:
handleUnknownTransaction(context);
break;
}
// 清理线程本地变量
cleanupThreadLocal(context);
// 记录完成日志
logTransactionCompletion(context, status);
}
} finally {
// 确保清理context
contextHolder.remove();
}
}
/**
* 事务开始时的初始化
*/
@EventListener
public void onTransactionStart(TransactionStartedEvent event) {
TransactionContext context = new TransactionContext();
context.setTransactionId(generateTransactionId());
context.setStartTime(System.currentTimeMillis());
context.getMetadata().put("event", event);
contextHolder.set(context);
// 注册同步器
TransactionSynchronizationManager.registerSynchronization(this);
logger.debug("Transaction context initialized: {}", context.getTransactionId());
}
/**
* 更新缓存(异步)
*/
private void updateCachesAsync(TransactionContext context) {
CompletableFuture.runAsync(() -> {
try {
// 获取所有需要更新的缓存键
Set<String> cacheKeys = extractCacheKeysFromContext(context);
// 批量更新缓存
cacheManager.refreshCaches(cacheKeys);
logger.debug("Updated {} caches for transaction: {}",
cacheKeys.size(), context.getTransactionId());
} catch (Exception e) {
logger.error("Cache update failed for transaction: {}",
context.getTransactionId(), e);
}
}, cacheRefreshExecutor);
}
/**
* 发布事件(异步)
*/
private void publishEventsAsync(TransactionContext context) {
List<DomainEvent> events = extractEventsFromContext(context);
for (DomainEvent event : events) {
messagePublisher.publishAsync(event)
.thenAccept(result -> {
logger.debug("Event published: {}", event.getEventId());
})
.exceptionally(ex -> {
logger.error("Failed to publish event: {}", event.getEventId(), ex);
// 记录失败,可以重试
recordFailedEvent(event, ex);
return null;
});
}
}
/**
* 处理已提交的事务
*/
private void handleCommittedTransaction(TransactionContext context) {
// 1. 更新监控指标
metricsCollector.recordTransactionSuccess(
context.getTransactionId(),
System.currentTimeMillis() – context.getStartTime(),
context.getOperationLogs().size()
);
// 2. 清理临时数据
cleanupTemporaryData(context);
// 3. 通知相关方
notifyStakeholders(context);
}
/**
* 处理已回滚的事务
*/
private void handleRolledBackTransaction(TransactionContext context) {
// 1. 记录回滚原因
auditLogger.logFailure("transaction_rolled_back",
Map.of(
"transactionId", context.getTransactionId(),
"reason", extractRollbackReason(context),
"duration", System.currentTimeMillis() – context.getStartTime()
));
// 2. 清理部分提交的数据
rollbackPartialChanges(context);
// 3. 补偿操作
executeCompensationActions(context);
}
}
篇幅限制下面就只能给大家展示小册部分内容了。整理了一份核心面试笔记包括了:Java面试、Spring、JVM、MyBatis、Redis、MySQL、并发编程、微服务、Linux、Springboot、SpringCloud、MQ、Kafc
需要全套面试笔记及答案 【点击此处即可/免费获取】
2. 多数据源事务同步
java
复制
下载
/**
* 多数据源事务同步管理器
*/
@Component
public class MultiDataSourceSynchronizationManager {
private final Map<String, DataSource> dataSources;
private final Map<String, PlatformTransactionManager> transactionManagers;
/**
* 注册多数据源同步器
*/
public void registerMultiDataSourceSynchronizations() {
for (String dataSourceName : dataSources.keySet()) {
TransactionSynchronization synchronization =
new DataSourceSpecificSynchronization(dataSourceName);
TransactionSynchronizationManager.registerSynchronization(synchronization);
}
}
/**
* 数据源特定的同步器
*/
private class DataSourceSpecificSynchronization
implements TransactionSynchronization, Ordered {
private final String dataSourceName;
private Connection connection;
public DataSourceSpecificSynchronization(String dataSourceName) {
this.dataSourceName = dataSourceName;
}
@Override
public int getOrder() {
// 数据源同步器优先级较低
return Ordered.LOWEST_PRECEDENCE + 100;
}
@Override
public void beforeCommit(boolean readOnly) {
if (!readOnly) {
try {
// 准备数据源特定的提交前操作
prepareDataSourceForCommit();
// 设置保存点(如果需要)
if (supportsSavepoints()) {
setSavepoint();
}
} catch (SQLException e) {
throw new DataAccessException(
"Failed to prepare data source for commit: " + dataSourceName, e);
}
}
}
@Override
public void beforeCompletion() {
try {
// 确保所有语句都执行完毕
if (connection != null && !connection.getAutoCommit()) {
// 刷新批处理
flushBatchStatements();
// 检查约束
checkConstraints();
}
} catch (SQLException e) {
logger.warn("Failed to complete operations for data source: " +
dataSourceName, e);
}
}
@Override
public void afterCommit() {
// 数据源特定的提交后处理
CompletableFuture.runAsync(() -> {
try {
// 异步更新从库
updateReadReplicas();
// 清理临时表
cleanupTempTables();
// 更新统计信息
updateStatistics();
} catch (Exception e) {
logger.error("Post-commit processing failed for data source: " +
dataSourceName, e);
}
}, asyncExecutor);
}
@Override
public void afterCompletion(int status) {
try {
if (status == STATUS_COMMITTED) {
// 提交成功后的清理
cleanupAfterCommit();
} else if (status == STATUS_ROLLED_BACK) {
// 回滚后的恢复
recoverAfterRollback();
// 记录回滚日志
logRollback(dataSourceName);
}
// 关闭连接
closeConnectionQuietly();
} finally {
// 从ThreadLocal清理连接
clearConnectionHolder();
}
}
private void prepareDataSourceForCommit() throws SQLException {
DataSource dataSource = dataSources.get(dataSourceName);
ConnectionHolder holder = (ConnectionHolder)
TransactionSynchronizationManager.getResource(dataSource);
if (holder != null) {
connection = holder.getConnection();
// 设置连接属性
if (!connection.getAutoCommit()) {
// 确保隔离级别
applyIsolationLevel();
// 设置超时
applyStatementTimeout();
// 启用/禁用约束检查
configureConstraints();
}
}
}
}
/**
* 连接持有器管理器
*/
public static class ConnectionHolderManager {
/**
* 获取或创建连接持有器
*/
public static ConnectionHolder getConnectionHolder(
DataSource dataSource, boolean allowCreate) {
ConnectionHolder holder = (ConnectionHolder)
TransactionSynchronizationManager.getResource(dataSource);
if (holder == null && allowCreate) {
holder = new ConnectionHolder(dataSource);
TransactionSynchronizationManager.bindResource(dataSource, holder);
}
return holder;
}
/**
* 连接持有器
*/
public static class ConnectionHolder implements ResourceHolder {
private final DataSource dataSource;
private Connection connection;
private boolean transactionActive = false;
private boolean readOnly = false;
private Integer isolationLevel;
private boolean synchronizedWithTransaction = false;
public ConnectionHolder(DataSource dataSource) {
this.dataSource = dataSource;
}
public Connection getConnection() throws SQLException {
if (this.connection == null) {
this.connection = dataSource.getConnection();
// 应用事务设置
applyTransactionSettings();
}
return this.connection;
}
private void applyTransactionSettings() throws SQLException {
if (transactionActive && !connection.getAutoCommit()) {
if (readOnly) {
connection.setReadOnly(true);
}
if (isolationLevel != null) {
connection.setTransactionIsolation(isolationLevel);
}
}
}
@Override
public void released() {
closeConnection();
clear();
}
@Override
public boolean isVoid() {
return connection == null;
}
private void closeConnection() {
if (this.connection != null) {
try {
if (!this.connection.getAutoCommit()) {
this.connection.rollback();
}
this.connection.close();
} catch (SQLException ex) {
logger.debug("Could not close JDBC connection", ex);
} finally {
this.connection = null;
}
}
}
}
}
}
五、Spring框架集成
1. 声明式事务与同步器集成
java
复制
下载
/**
* Spring声明式事务管理器扩展
*/
@Configuration
@EnableTransactionManagement
public class TransactionSyncConfiguration {
@Bean
public PlatformTransactionManager transactionManager(DataSource dataSource) {
DataSourceTransactionManager transactionManager =
new DataSourceTransactionManager(dataSource);
// 启用同步管理器
transactionManager.setTransactionSynchronization(
DataSourceTransactionManager.SYNCHRONIZATION_ALWAYS);
// 设置超时
transactionManager.setDefaultTimeout(30); // 30秒
// 设置只读优化
transactionManager.setEnforceReadOnly(true);
return transactionManager;
}
/**
* 事务拦截器增强
*/
@Bean
public TransactionInterceptor transactionInterceptor(
PlatformTransactionManager transactionManager) {
TransactionInterceptor interceptor = new TransactionInterceptor();
interceptor.setTransactionManager(transactionManager);
// 配置事务属性源
NameMatchTransactionAttributeSource source =
new NameMatchTransactionAttributeSource();
// 方法级别的事务配置
Properties properties = new Properties();
properties.setProperty("get*",
"PROPAGATION_REQUIRED,readOnly,-Exception");
properties.setProperty("find*",
"PROPAGATION_REQUIRED,readOnly,-Exception");
properties.setProperty("save*",
"PROPAGATION_REQUIRED,-Exception");
properties.setProperty("update*",
"PROPAGATION_REQUIRED,-Exception");
properties.setProperty("delete*",
"PROPAGATION_REQUIRED,-Exception");
properties.setProperty("*",
"PROPAGATION_REQUIRED,-Exception");
source.setProperties(properties);
interceptor.setTransactionAttributeSource(source);
// 设置事务同步管理器
interceptor.setTransactionSynchronizationManager(
new EnhancedTransactionSynchronizationManager());
return interceptor;
}
/**
* 增强的事务同步管理器
*/
public static class EnhancedTransactionSynchronizationManager
extends TransactionSynchronizationManager {
private final List<TransactionSynchronizationInitializer> initializers;
public EnhancedTransactionSynchronizationManager() {
this.initializers = Arrays.asList(
new AuditSynchronizationInitializer(),
new CacheSynchronizationInitializer(),
new EventSynchronizationInitializer(),
new MetricsSynchronizationInitializer()
);
}
@Override
public void initSynchronization() throws IllegalStateException {
super.initSynchronization();
// 初始化所有同步器
for (TransactionSynchronizationInitializer initializer : initializers) {
TransactionSynchronization synchronization =
initializer.createSynchronization();
if (synchronization != null) {
registerSynchronization(synchronization);
}
}
// 触发初始化完成事件
publishEvent(new SynchronizationInitializedEvent(this));
}
@Override
public void clearSynchronization() throws IllegalStateException {
// 触发清理前事件
publishEvent(new SynchronizationClearingEvent(this));
super.clearSynchronization();
}
}
/**
* 事务同步器初始化接口
*/
public interface TransactionSynchronizationInitializer {
TransactionSynchronization createSynchronization();
}
}
2. 事务事件发布机制
java
复制
下载
/**
* 事务事件发布管理器
*/
@Component
public class TransactionEventPublisher {
private final ApplicationEventPublisher eventPublisher;
private final TransactionSynchronizationManager syncManager;
/**
* 注册事务事件同步器
*/
@PostConstruct
public void registerEventSynchronization() {
syncManager.registerSynchronization(new EventPublishingSynchronization());
}
/**
* 事件发布同步器
*/
private class EventPublishingSynchronization implements TransactionSynchronization {
private final ThreadLocal<List<DomainEvent>> pendingEvents =
new ThreadLocal<>();
private final ThreadLocal<List<DomainEvent>> committedEvents =
new ThreadLocal<>();
@Override
public void beforeCommit(boolean readOnly) {
if (!readOnly) {
// 收集待发布的事件
List<DomainEvent> events = collectPendingEvents();
pendingEvents.set(events);
// 验证事件
validateEvents(events);
}
}
@Override
public void afterCommit() {
List<DomainEvent> events = pendingEvents.get();
if (events != null && !events.isEmpty()) {
// 保存已提交事件
committedEvents.set(new ArrayList<>(events));
// 异步发布事件
publishEventsAsync(events);
// 清理待处理事件
pendingEvents.remove();
}
}
@Override
public void afterCompletion(int status) {
if (status == STATUS_COMMITTED) {
// 标记事件为已发布
List<DomainEvent> events = committedEvents.get();
if (events != null) {
markEventsAsPublished(events);
committedEvents.remove();
}
} else {
// 事务回滚,清理事件
pendingEvents.remove();
committedEvents.remove();
clearPendingEvents();
}
}
private List<DomainEvent> collectPendingEvents() {
// 从ThreadLocal收集事件
List<DomainEvent> events = new ArrayList<>();
// 从各种上下文收集事件
events.addAll(TransactionEventContextHolder.getEvents());
events.addAll(DomainEventCollector.getPendingEvents());
events.addAll(EntityChangeTracker.getChangeEvents());
return events;
}
private void publishEventsAsync(List<DomainEvent> events) {
CompletableFuture.runAsync(() -> {
for (DomainEvent event : events) {
try {
eventPublisher.publishEvent(event);
logEventPublished(event);
} catch (Exception e) {
logEventPublishFailure(event, e);
handlePublishFailure(event, e);
}
}
}, eventPublishingExecutor);
}
}
/**
* 事务事件上下文持有器
*/
public static class TransactionEventContextHolder {
private static final ThreadLocal<EventContext> contextHolder =
new ThreadLocal<>();
@Data
public static class EventContext {
private String transactionId;
private List<DomainEvent> events = new ArrayList<>();
private Map<String, Object> metadata = new HashMap<>();
public void addEvent(DomainEvent event) {
events.add(event);
}
public void addMetadata(String key, Object value) {
metadata.put(key, value);
}
}
public static EventContext getContext() {
EventContext context = contextHolder.get();
if (context == null) {
context = new EventContext();
context.setTransactionId(UUID.randomUUID().toString());
contextHolder.set(context);
}
return context;
}
public static List<DomainEvent> getEvents() {
EventContext context = getContext();
return context.getEvents();
}
public static void clear() {
contextHolder.remove();
}
}
}
六、性能优化与监控
1. 事务同步性能监控
java
复制
下载
/**
* 事务同步性能监控器
*/
@Component
public class TransactionSyncPerformanceMonitor {
private final MeterRegistry meterRegistry;
private final Map<String, SyncMetrics> metricsMap = new ConcurrentHashMap<>();
/**
* 监控事务同步执行
*/
public void monitorSynchronizationExecution(
String synchronizationName,
String phase,
long duration,
boolean success) {
// 记录指标
Timer timer = Timer.builder("transaction.sync.duration")
.tag("sync", synchronizationName)
.tag("phase", phase)
.tag("success", String.valueOf(success))
.publishPercentiles(0.5, 0.95, 0.99)
.register(meterRegistry);
timer.record(duration, TimeUnit.MILLISECONDS);
// 更新本地统计
SyncMetrics metrics = metricsMap.computeIfAbsent(
synchronizationName, k -> new SyncMetrics());
metrics.recordExecution(phase, duration, success);
// 检测性能问题
if (duration > 1000) { // 超过1秒
warnSlowSynchronization(synchronizationName, phase, duration);
}
}
/**
* 生成性能报告
*/
public PerformanceReport generateReport(Duration period) {
PerformanceReport report = new PerformanceReport();
report.setPeriod(period);
report.setGeneratedAt(new Date());
List<SyncPerformance> performances = new ArrayList<>();
for (Map.Entry<String, SyncMetrics> entry : metricsMap.entrySet()) {
SyncPerformance perf = new SyncPerformance();
perf.setSyncName(entry.getKey());
perf.setMetrics(entry.getValue().getSummary(period));
performances.add(perf);
}
report.setPerformances(performances);
// 分析性能瓶颈
report.setBottlenecks(identifyBottlenecks(performances));
// 提供优化建议
report.setRecommendations(generateRecommendations(report));
return report;
}
/**
* 同步器指标
*/
@Data
public static class SyncMetrics {
private Map<String, PhaseMetrics> phaseMetrics = new ConcurrentHashMap<>();
private long totalExecutions;
private long failedExecutions;
private long totalDuration;
public void recordExecution(String phase, long duration, boolean success) {
PhaseMetrics metrics = phaseMetrics.computeIfAbsent(
phase, k -> new PhaseMetrics());
metrics.recordExecution(duration, success);
totalExecutions++;
totalDuration += duration;
if (!success) {
failedExecutions++;
}
}
public MetricsSummary getSummary(Duration period) {
MetricsSummary summary = new MetricsSummary();
summary.setTotalExecutions(totalExecutions);
summary.setFailedExecutions(failedExecutions);
summary.setSuccessRate(
totalExecutions > 0 ?
(double) (totalExecutions – failedExecutions) / totalExecutions : 1.0);
summary.setAverageDuration(
totalExecutions > 0 ?
(double) totalDuration / totalExecutions : 0);
// 按阶段统计
Map<String, PhaseSummary> phaseSummaries = new HashMap<>();
for (Map.Entry<String, PhaseMetrics> entry : phaseMetrics.entrySet()) {
phaseSummaries.put(entry.getKey(), entry.getValue().getSummary());
}
summary.setPhaseSummaries(phaseSummaries);
return summary;
}
}
/**
* 事务同步优化建议器
*/
@Component
public static class SyncOptimizationAdvisor {
public List<OptimizationSuggestion> analyzeAndSuggest(
PerformanceReport report) {
List<OptimizationSuggestion> suggestions = new ArrayList<>();
for (SyncPerformance perf : report.getPerformances()) {
MetricsSummary summary = perf.getMetrics();
// 检查失败率
if (summary.getSuccessRate() < 0.95) {
suggestions.add(createSuggestion(
perf.getSyncName(),
"HIGH_FAILURE_RATE",
String.format("同步器失败率较高: %.1f%%,建议检查错误处理",
(1 – summary.getSuccessRate()) * 100),
"检查同步器的异常处理逻辑,确保不会影响事务提交"
));
}
// 检查执行时间
if (summary.getAverageDuration() > 100) { // 超过100ms
suggestions.add(createSuggestion(
perf.getSyncName(),
"HIGH_EXECUTION_TIME",
String.format("同步器执行时间较长: %.1fms",
summary.getAverageDuration()),
"考虑将耗时操作移到afterCommit中异步执行"
));
}
// 检查特定阶段的性能
for (Map.Entry<String, PhaseSummary> entry :
summary.getPhaseSummaries().entrySet()) {
if (entry.getValue().getAverageDuration() > 50) {
suggestions.add(createSuggestion(
perf.getSyncName(),
"PHASE_PERFORMANCE_ISSUE",
String.format("阶段[%s]执行时间较长: %.1fms",
entry.getKey(), entry.getValue().getAverageDuration()),
"优化该阶段的处理逻辑,或考虑延迟执行"
));
}
}
}
return suggestions;
}
}
}
篇幅限制下面就只能给大家展示小册部分内容了。整理了一份核心面试笔记包括了:Java面试、Spring、JVM、MyBatis、Redis、MySQL、并发编程、微服务、Linux、Springboot、SpringCloud、MQ、Kafc
需要全套面试笔记及答案 【点击此处即可/免费获取】
2. 线程池优化配置
java
复制
下载
/**
* 事务同步线程池配置
*/
@Configuration
public class TransactionSyncExecutorConfiguration {
/**
* 异步执行器 – 用于afterCommit等异步操作
*/
@Bean(name = "transactionSyncExecutor")
public Executor transactionSyncExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
// 核心配置
executor.setCorePoolSize(10);
executor.setMaxPoolSize(50);
executor.setQueueCapacity(1000);
executor.setKeepAliveSeconds(60);
// 线程配置
executor.setThreadNamePrefix("tx-sync-");
executor.setThreadPriority(Thread.NORM_PRIORITY);
// 拒绝策略 – 使用调用者运行,避免事务等待超时
executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
// 等待任务完成
executor.setWaitForTasksToCompleteOnShutdown(true);
executor.setAwaitTerminationSeconds(30);
// 增强的线程工厂
executor.setThreadFactory(new TransactionAwareThreadFactory());
executor.initialize();
return executor;
}
/**
* 事务感知的线程工厂
*/
private static class TransactionAwareThreadFactory implements ThreadFactory {
private final AtomicInteger threadNumber = new AtomicInteger(1);
private final String namePrefix = "tx-sync-thread-";
@Override
public Thread newThread(Runnable r) {
Thread t = new Thread(r, namePrefix + threadNumber.getAndIncrement());
// 设置线程上下文
t.setUncaughtExceptionHandler(new TransactionSyncExceptionHandler());
// 传递事务上下文
if (r instanceof TransactionAwareRunnable) {
((TransactionAwareRunnable) r).propagateContext();
}
return t;
}
}
/**
* 事务感知的Runnable
*/
public abstract static class TransactionAwareRunnable implements Runnable {
private final Map<Object, Object> resources;
private final Set<TransactionSynchronization> synchronizations;
private final String transactionName;
private final boolean readOnly;
public TransactionAwareRunnable() {
// 捕获当前线程的事务上下文
this.resources = TransactionSynchronizationManager.getResourceMap();
this.synchronizations = TransactionSynchronizationManager.getSynchronizations();
this.transactionName = TransactionSynchronizationManager.getCurrentTransactionName();
this.readOnly = TransactionSynchronizationManager.isCurrentTransactionReadOnly();
}
protected void propagateContext() {
// 在新线程中恢复事务上下文
if (resources != null) {
TransactionSynchronizationManager.setResources(resources);
}
if (synchronizations != null) {
TransactionSynchronizationManager.setSynchronizations(synchronizations);
}
if (transactionName != null) {
TransactionSynchronizationManager.setCurrentTransactionName(transactionName);
}
TransactionSynchronizationManager.setCurrentTransactionReadOnly(readOnly);
}
@Override
public abstract void run();
@Override
protected void finalize() throws Throwable {
// 确保清理线程本地变量
TransactionSynchronizationManager.clear();
super.finalize();
}
}
}
七、常见问题与解决方案
1. 内存泄漏预防
java
复制
下载
/**
* 事务同步内存泄漏检测器
*/
@Component
public class TransactionSyncMemoryLeakDetector {
private final ScheduledExecutorService detector =
Executors.newSingleThreadScheduledExecutor();
@PostConstruct
public void startDetection() {
// 每10分钟检测一次
detector.scheduleAtFixedRate(this::detectLeaks, 10, 10, TimeUnit.MINUTES);
}
/**
* 检测内存泄漏
*/
private void detectLeaks() {
try {
// 1. 检测未清理的ThreadLocal
detectThreadLocalLeaks();
// 2. 检测长时间运行的事务
detectLongRunningTransactions();
// 3. 检测未释放的资源
detectUnreleasedResources();
// 4. 检测同步器泄漏
detectSynchronizationLeaks();
} catch (Exception e) {
logger.error("Memory leak detection failed", e);
}
}
/**
* ThreadLocal泄漏检测
*/
private void detectThreadLocalLeaks() {
// 获取所有活动线程
Set<Thread> threads = Thread.getAllStackTraces().keySet();
for (Thread thread : threads) {
if (thread.isAlive() && !thread.isDaemon()) {
try {
// 检查线程是否持有事务资源
Map<Object, Object> resources = getThreadResources(thread);
if (resources != null && !resources.isEmpty()) {
logPotentialLeak(thread, resources);
}
} catch (Exception e) {
logger.debug("Failed to check thread: {}", thread.getName(), e);
}
}
}
}
/**
* 长时间运行事务检测
*/
private void detectLongRunningTransactions() {
Map<Thread, TransactionInfo> activeTransactions =
TransactionMonitor.getActiveTransactions();
long now = System.currentTimeMillis();
long warningThreshold = 30 * 1000; // 30秒
long criticalThreshold = 60 * 1000; // 60秒
for (Map.Entry<Thread, TransactionInfo> entry : activeTransactions.entrySet()) {
long duration = now – entry.getValue().getStartTime();
if (duration > criticalThreshold) {
logger.error("Critical: Transaction running for {}ms on thread {}",
duration, entry.getKey().getName());
handleCriticalTransaction(entry.getKey(), entry.getValue());
} else if (duration > warningThreshold) {
logger.warn("Warning: Transaction running for {}ms on thread {}",
duration, entry.getKey().getName());
}
}
}
/**
* 预防性清理
*/
@Component
public static class TransactionCleanupInterceptor implements HandlerInterceptor {
@Override
public void afterCompletion(HttpServletRequest request,
HttpServletResponse response,
Object handler,
Exception ex) {
// 确保清理ThreadLocal
TransactionSynchronizationManager.clear();
// 清理自定义ThreadLocal
CustomThreadLocalManager.clearAll();
// 记录清理日志
if (logger.isDebugEnabled()) {
logger.debug("Cleaned up transaction context for request: {}",
request.getRequestURI());
}
}
}
}
总结
核心机制总结:
线程绑定机制:
-
使用ThreadLocal实现线程隔离的事务上下文
-
支持嵌套事务和事务挂起/恢复
-
确保事务资源的线程安全访问
回调生命周期:
-
完整的事务生命周期回调(beforeCommit, afterCommit, afterCompletion等)
-
支持优先级排序的同步器执行
-
异常处理与事务状态传播
资源管理:
-
资源绑定与自动释放
-
连接持有器模式管理数据库连接
-
多数据源事务同步支持
生产级特性:
-
性能监控与优化
-
内存泄漏检测与预防
-
异步处理支持
-
完整的事件发布机制
最佳实践建议:
同步器设计原则:
-
保持同步器轻量级,避免在beforeCommit中执行耗时操作
-
将耗时操作移到afterCommit中异步执行
-
确保同步器的异常处理不会影响事务提交
资源管理:
-
始终确保资源在事务完成后正确释放
-
使用ConnectionHolder管理数据库连接
-
实现ResourceHolder接口支持自定义资源
性能优化:
-
监控同步器执行时间,优化慢速同步器
-
合理配置线程池,避免线程饥饿
-
使用异步执行减少事务提交延迟
故障排查:
-
实现完整的事务追踪和日志记录
-
检测长时间运行的事务
-
定期检查ThreadLocal泄漏
TransactionSynchronizationManager是Spring事务管理的核心组件,通过合理利用其机制,可以实现强大的事务扩展功能,同时保证系统的稳定性和性能。




