You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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());
    }
  };
}

问题根因

  1. 自动删除远程文件:FtpInboundFileSynchronizer设置了setDeleteRemoteFiles(true),拉取完成后直接删除远程源文件,导致后续无法执行归档移动。
  2. 错误操作本地文件:事务同步的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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.01 02:10:55