transactionSynchronizationFactory结合数据库事务失效,FTP文件未归档
问题描述
期望流程:远程目录创建文件后,通过Poller拉取文件,执行数据库插入操作,再将远程源文件移动到另一远程归档目录。
当前问题:文件已成功拉取,数据库插入完成,但远程源目录的文件被删除,却未移动到目标归档目录。
环境信息
- Spring Boot版本:3.2.1
- 数据库:Postgres
Java配置代码
@Bean public DefaultFtpsSessionFactory ftpsSessionFactory() { DefaultFtpsSessionFactory sessionFactory = new DefaultFtpsSessionFactory(); sessionFactory.setHost(ftpHost); sessionFactory.setUsername(ftpUsername); sessionFactory.setPassword(ftpPassword); sessionFactory.setPort(ftpPort); sessionFactory.setControlEncoding(StandardCharsets.UTF_8.name()); sessionFactory.setImplicit(true); sessionFactory.setProtocol("TLS"); return sessionFactory; } @Bean public FtpInboundFileSynchronizer ftpInboundFileSynchronizer() { FtpInboundFileSynchronizer fileSynchronizer = new FtpInboundFileSynchronizer(ftpsSessionFactory()); fileSynchronizer.setDeleteRemoteFiles(true); fileSynchronizer.setRemoteDirectory("/remote-in"); fileSynchronizer.setFilter(new FtpSimplePatternFileListFilter("*.txt")); return fileSynchronizer; } @Bean @InboundChannelAdapter(channel = "ftpChannel", poller = @Poller(value = "pollerMetadata")) public MessageSource<File> ftpMessageSource() { FtpInboundFileSynchronizingMessageSource source = new FtpInboundFileSynchronizingMessageSource(ftpInboundFileSynchronizer()); source.setLocalDirectory(new File("/tmp")); source.setAutoCreateLocalDirectory(true); source.setLocalFilter(new AcceptOnceFileListFilter<>()); source.setMaxFetchSize(Integer.MIN_VALUE); return source; } @Bean public PollerMetadata pollerMetadata() { return Pollers.fixedDelay(ftpPollerDelay) .advice(transactionInterceptor()) .transactionSynchronizationFactory(transactionSynchronizationFactory()) .getObject(); } @Bean public JpaTransactionManager transactionManager() { return new JpaTransactionManager(); } @Bean public TransactionInterceptor transactionInterceptor() { return new TransactionInterceptorBuilder() .transactionManager(transactionManager()) .build(); } public TransactionSynchronizationFactory transactionSynchronizationFactory() { ExpressionEvaluatingTransactionSynchronizationProcessor processor = new ExpressionEvaluatingTransactionSynchronizationProcessor(); SpelExpressionParser spelParser = new SpelExpressionParser(); processor.setAfterCommitExpression( spelParser.parseExpression("payload.renameTo(new java.io.File" + "(#systemProperties.get('ftp.directory.archive' ) + " + "T(java.io.File).separator + payload.name))")); return new DefaultTransactionSynchronizationFactory(processor); } @Bean @ServiceActivator(inputChannel = "ftpChannel") public MessageHandler ftpInboundMessageHandler() { return message -> { Object payload = message.getPayload(); if (payload instanceof File file) { // save record in db } else { log.error("Received invalid message payload {}", message.getPayload()); } }; }
问题根因
- 自动删除远程文件:
FtpInboundFileSynchronizer设置了setDeleteRemoteFiles(true),拉取完成后直接删除远程源文件,导致后续无法执行归档移动。 - 错误操作本地文件:事务同步的
afterCommitExpression针对的是本地临时文件做重命名,而非远程文件,完全不符合远程归档的需求。
解决方案
步骤1:关闭自动删除远程文件
修改ftpInboundFileSynchronizer配置,禁止拉取后自动删除远程源文件:
@Bean public FtpInboundFileSynchronizer ftpInboundFileSynchronizer() { FtpInboundFileSynchronizer fileSynchronizer = new FtpInboundFileSynchronizer(ftpsSessionFactory()); fileSynchronizer.setDeleteRemoteFiles(false); // 改为false,保留远程源文件用于归档 fileSynchronizer.setRemoteDirectory("/remote-in"); fileSynchronizer.setFilter(new FtpSimplePatternFileListFilter("*.txt")); return fileSynchronizer; }
步骤2:实现远程文件归档逻辑
重构transactionSynchronizationFactory,通过FTPS会话执行远程文件的移动操作(FTP的rename方法可实现跨目录移动):
@Bean public TransactionSynchronizationFactory transactionSynchronizationFactory(FtpsSessionFactory ftpsSessionFactory) { ExpressionEvaluatingTransactionSynchronizationProcessor processor = new ExpressionEvaluatingTransactionSynchronizationProcessor(); SpelExpressionParser spelParser = new SpelExpressionParser(); // 事务提交后,将远程源文件移动到归档目录 processor.setAfterCommitExpression(spelParser.parseExpression( "@ftpsSessionFactory.getSession().rename(" + "'/remote-in/' + payload.name, " + "'/remote-archive/' + payload.name)" )); return new DefaultTransactionSynchronizationFactory(processor); }
注意:确保远程归档目录
/remote-archive/已提前创建,否则FTP的rename操作会失败。
步骤3:补充本地文件清理(可选)
如果需要清理本地临时文件,可在afterCommitExpression中追加删除逻辑:
processor.setAfterCommitExpression(spelParser.parseExpression( "@ftpsSessionFactory.getSession().rename('/remote-in/' + payload.name, '/remote-archive/' + payload.name) && payload.delete()" ));
步骤4:异常处理优化
在数据库插入操作中添加异常捕获,确保事务回滚时不会执行远程文件移动;若远程移动失败,可根据业务需求添加重试逻辑或告警。
内容的提问来源于stack exchange,提问作者riteshmaurya
相关产品推荐
相关产品推荐

