Spring Integration同流程多事务同步工厂不生效问题咨询
问题根因
该问题是Spring Integration事务同步的注册机制导致的:<int:transactional> 子元素仅在对应组件触发事务开启时,才会将关联的事务同步工厂注册到当前事务上下文。轮询器已经开启了顶层事务,转换器上配置的事务因传播行为为REQUIRED加入已有事务,不会触发新的同步注册流程,因此synchFactory2的提交后逻辑永远不会执行。
最优解决方案
无需新增任何业务组件,仅调整少量配置即可实现需求,核心逻辑是复用轮询器的事务同步工厂,通过发布订阅通道同时触发两个提交后动作:
- 删除原转换器上的
<int:transactional>配置和synchFactory2,无需额外的事务配置 - 去掉
<int:chain>的output-channel="nullChannel"属性,转换后的消息会作为轮询器的处理结果,自动绑定到顶层事务的同步上下文 - 将轮询器关联的事务同步工厂的提交后通道改为发布订阅通道,配置两个订阅者:
- 原有SFTP重命名网关,使用消息头保留的SFTP元数据执行文件重命名
- 简易桥接器,直接将转换后的消息转发到
onCommitSecondFlowChannel
调整后的核心配置如下:
<!-- 新增发布订阅通道,用于事务提交后同时触发两个动作 --> <int:publish-subscribe-channel id="afterCommitChannel" /> <!-- 复用原有同步工厂,修改提交后通道为新增的发布订阅通道 --> <int:transaction-synchronization-factory id="synchFactory1"> <int:after-commit channel="afterCommitChannel" /> </int:transaction-synchronization-factory> <!-- 原有SFTP重命名网关仅修改request-channel为afterCommitChannel,其余配置不变 --> <int-sftp:outbound-gateway session-factory="mySftpSessionFactory" request-channel="afterCommitChannel" command="mv" expression="headers[T(org.springframework.integration.file.FileHeaders).REMOTE_DIRECTORY].concat('/'.concat(headers[T(org.springframework.integration.file.FileHeaders).REMOTE_FILE]))" rename-expression="headers[T(org.springframework.integration.file.FileHeaders).REMOTE_DIRECTORY].concat('/'.concat(headers[T(org.springframework.integration.file.FileHeaders).REMOTE_FILE].concat('.done')))" requires-reply="false" reply-channel="nullChannel" /> <!-- 新增桥接器,直接转发转换后的消息到第二个流程通道 --> <int:bridge input-channel="afterCommitChannel" output-channel="onCommitSecondFlowChannel" /> <!-- 调整chain配置,去掉output-channel,删除转换器上的事务配置 --> <int:chain input-channel="channel1"> <int:stream-transformer charset="UTF-8" /> <int:json-to-object-transformer type="com.example.MyDto" /> <int:transformer ref="myTransformer" /> </int:chain>
该方案完全满足需求:两个提交后动作都在顶层事务提交后执行,文件重命名和消息发送的触发时机一致,不需要新增任何Service Activator或Gateway组件。
备选方案(手动注册同步)
如果不想调整现有配置结构,可以直接在自定义转换器myTransformer中手动注册事务提交后回调:
import org.springframework.transaction.support.TransactionSynchronization; import org.springframework.transaction.support.TransactionSynchronizationManager; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.support.MessageBuilder; // 注入目标通道 private MessageChannel onCommitSecondFlowChannel; public Object transform(MyDto input) { // 原有转换逻辑 Object processedResult = doTransform(input); // 注册事务提交后回调 if (TransactionSynchronizationManager.isActualTransactionActive()) { TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() { @Override public void afterCommit() { onCommitSecondFlowChannel.send(MessageBuilder.withPayload(processedResult).build()); } }); } return processedResult; }
使用该方案仅需删除原有转换器上的<int:transactional>配置和synchFactory2即可,不需要修改其他XML配置。
内容的提问来源于stack exchange,提问作者Mauro Molinari
相关产品推荐
相关产品推荐

