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

Spring Integration递归读取远程SFTP子目录文件问题求助

Spring Integration SFTP 递归读取文件问题及优化需求

初始需求与代码

已基于Spring Integration的入站适配器实现从远程SFTP服务器单个目录获取文件的工作流程,现在需要实现递归读取远程父目录下所有子目录中的文件。初始代码如下:

@Bean
public SessionFactory<SftpClient.DirEntry> sftpSessionFactory() {
    DefaultSftpSessionFactory factory = new DefaultSftpSessionFactory(true);
    factory.setHost("localhost");
    factory.setPort(port);
    factory.setUser("foo");
    factory.setPassword("foo");
    factory.setAllowUnknownKeys(true);
    factory.setTestSession(true);
    return new CachingSessionFactory<>(factory);
}

@Bean
public SftpInboundFileSynchronizer sftpInboundFileSynchronizer() {
    SftpInboundFileSynchronizer fileSynchronizer = new SftpInboundFileSynchronizer(sftpSessionFactory());
    fileSynchronizer.setDeleteRemoteFiles(false);
    fileSynchronizer.setRemoteDirectory("foo");
    fileSynchronizer.setFilter(new SftpSimplePatternFileListFilter("*.xml"));
    return fileSynchronizer;
}

@Bean
@InboundChannelAdapter(channel = "sftpChannel", poller = @Poller(fixedDelay = "5000"))
public MessageSource<File> sftpMessageSource() {
    SftpInboundFileSynchronizingMessageSource source =
            new SftpInboundFileSynchronizingMessageSource(sftpInboundFileSynchronizer());
    source.setLocalDirectory(new File("sftp-inbound"));
    source.setAutoCreateLocalDirectory(true);
    source.setLocalFilter(new AcceptOnceFileListFilter<File>());
    source.setMaxFetchSize(1);
    return source;
}

@Bean
@ServiceActivator(inputChannel = "sftpChannel")
public MessageHandler handler() {
    return new MessageHandler() {
        @Override
        public void handleMessage(Message<?> message) throws MessagingException {
            System.out.println(message.getPayload());
        }
    };
}

第一次修改后的问题

感谢及时回复。已根据建议代码修改,应用可启动,但文件未被拾取。修改后的代码片段如下:

CompositeFileListFilter<LsEntry> compositeFileListFilter = new CompositeFileListFilter<>();
SftpPersistentAcceptOnceFileListFilter fileListFilter =
    new SftpPersistentAcceptOnceFileListFilter(
        (JdbcMetadataStore) context.getBean("metadataStore"), "REMOTE");
if (Constants.APP1.equals(appName) || Constants.APP2.equals(appName)) {
  SftpRegexPatternFileListFilter regexPatternFileListFilter =
      new SftpRegexPatternFileListFilter(Pattern.compile("^IL.*"));
  compositeFileListFilter.addFilter(regexPatternFileListFilter);
}
compositeFileListFilter.addFilter(fileListFilter);

return IntegrationFlows.fromSupplier(
        () -> sftpEnvironment.getSftpGLSIncomingDir(), // remote dir
        e -> e.autoStartup(true).poller(pollerMetada()))
    .handle(
        Sftp.outboundGateway(sftpSessionFactory(), Command.MGET, "payload")
            .options(Option.RECURSIVE)
            .filter(compositeFileListFilter)
            .fileExistsMode(FileExistsMode.IGNORE)
            .localDirectoryExpression("'/tmp/' + #remoteDirectory")) // re-create tree locally
    .split()
    .log()
    .get(); 

新实现的问题与待解决需求

改用新代码实现后,出现文件仅被部分处理的问题,例如10个文件中仅能处理5-6个,无法定位问题根源。此外还有两个待解决需求:

  • 当前实现可读取远程子目录文件并存储到本地,但希望无需存储到本地,直接在sftpChannel中处理这些文件;
  • 希望基于数据库实现去重机制,避免重复处理文件。

完整代码如下:

public class SFTPPollerService {

  @Bean
  public SessionFactory<LsEntry> sftpSessionFactory() {
    DefaultSftpSessionFactory factory = new DefaultSftpSessionFactory(true);
    //code
    return factory;
  }

//OLD code 
  //    @Bean
  //    public SftpInboundFileSynchronizer sftpInboundFileSynchronizer() {
  //        SftpInboundFileSynchronizer fileSynchronizer =
  //                new SftpInboundFileSynchronizer(sftpSessionFactory());
  //        fileSynchronizer.setDeleteRemoteFiles(sftpEnvironment.isDeleteRemoteFiles());
  //        fileSynchronizer.setRemoteDirectory(sftpEnvironment.getSftpGLSIncomingDir());
  //        fileSynchronizer.setPreserveTimestamp(true);
  //        CompositeFileListFilter<LsEntry> compositeFileListFilter = new
  // CompositeFileListFilter<>();
  //        SftpPersistentAcceptOnceFileListFilter fileListFilter =
  //                new SftpPersistentAcceptOnceFileListFilter(
  //                        (JdbcMetadataStore) context.getBean("metadataStore"), "REMOTE");
  //        if (Constants.app2.equals(appName)
  //                || Constants.app1.equals(appName)) {
  //            SftpRegexPatternFileListFilter regexPatternFileListFilter =
  //                    new SftpRegexPatternFileListFilter(Pattern.compile("*.txt"));
  //            compositeFileListFilter.addFilter(regexPatternFileListFilter);
  //        }
  //        compositeFileListFilter.addFilter(fileListFilter);
  //        fileSynchronizer.setFilter(compositeFileListFilter);
  //        return fileSynchronizer;
  //    }
  //
  //    @Bean
  //    @InboundChannelAdapter(channel = "sftpChannel", poller = @Poller("pollerMetada"))
  //    public MessageSource<File> sftpMessageSource() {
  //        SftpInboundFileSynchronizingMessageSource source =
  //                new SftpInboundFileSynchronizingMessageSource(sftpInboundFileSynchronizer());
  //        source.setLocalDirectory(new File(sftpEnvironment.getSftpLocalDir()));
  //
  //        source.setAutoCreateLocalDirectory(true);
  //
//          try {
//              source.setLocalFilter(
//                      (FileSystemPersistentAcceptOnceFileListFilter)
//                              context.getBean("filelistFilter"));
//          } catch (Exception e) {
//              LOG.error(
//                      "Exception caught while setting local filter on
//   SftpInboundFileSynchronizingMessageSource",
//                      e);
//          }
  //        source.setMaxFetchSize(sftpEnvironment.getMaxFetchFileSize());
  //
  //        return source;
  //    }

//new Code 
  @Bean
  public IntegrationFlow sftpInboundFlow() {

    CompositeFileListFilter<LsEntry> compositeFileListFilter = new CompositeFileListFilter<>();
    SftpPersistentAcceptOnceFileListFilter fileListFilter =
        new SftpPersistentAcceptOnceFileListFilter(
            (JdbcMetadataStore) context.getBean("metadataStore"), "REMOTE");
    if (Constants.app2.equals(appName) || Constants.app1.equals(appName)) {
      SftpRegexPatternFileListFilter regexPatternFileListFilter =
          new SftpRegexPatternFileListFilter(Pattern.compile("(subDir | *.txt)"));

      compositeFileListFilter.addFilter(regexPatternFileListFilter);
    }

    fileListFilter.setForRecursion(true);
    FileSystemPersistentAcceptOnceFileListFilter fileSystemPersistentAcceptOnceFileListFilter = (FileSystemPersistentAcceptOnceFileListFilter) context.getBean(
        "filelistFilter");
    compositeFileListFilter.addFilter(fileListFilter);

//    IntegrationFlow ir =
//        IntegrationFlows.from(
//                Sftp.inboundAdapter(sftpSessionFactory())
//                    .preserveTimestamp(true)
//                    .remoteDirectory(sftpEnvironment.getSftpGLSIncomingDir())
//                    .deleteRemoteFiles(sftpEnvironment.isDeleteRemoteFiles())
//                    .filter(compositeFileListFilter)
//                    .autoCreateLocalDirectory(true)
//                    .localDirectory(new File(sftpEnvironment.getSftpLocalDir())),
//                e -> e.autoStartup(true).poller(pollerMetada()))
//            .handle(handler())
//            .get();

    return IntegrationFlows.fromSupplier(
            () -> sftpEnvironment.getSftpGLSIncomingDir(), // remote dir
            e -> e.autoStartup(true).poller(pollerMetada()))
        .handle(
            Sftp.outboundGateway(sftpSessionFactory(), Command.MGET, "payload")
                .options(Option.RECURSIVE)
                .fileExistsMode(FileExistsMode.IGNORE)
                .regexFileNameFilter("(dsv[0-9]|.*.xml)")

                //                .filter(compositeFileListFilter)
                .localDirectoryExpression("'user/localDir/test/'"))
        //        .handle(handler())
        // .patternFileNameFilter(".*\\.xml")) // re-create tree locally
        .split()
        .channel("sftpChannel")
        // .handle(handler())
        .log()
        .get();
  }

  @Bean
  public PollerMetadata pollerMetada() {
    PollerMetadata pm = new PollerMetadata();
    ExpressionEvaluatingTransactionSynchronizationProcessor processor =
        new ExpressionEvaluatingTransactionSynchronizationProcessor();
    ExpressionParser parser = new SpelExpressionParser();
    Expression exp = parser.parseExpression("payload.delete()");
    processor.setAfterRollbackExpression(exp);
    TransactionSynchronizationFactory tsf = new DefaultTransactionSynchronizationFactory(processor);
    pm.setTransactionSynchronizationFactory(tsf);
    List<Advice> advices = new ArrayList<>();
    advices.add(compoundTriggerAdvice());
    pm.setAdviceChain(advices);
    pm.setTrigger(compoundTrigger());
    pm.setMaxMessagesPerPoll(sftpEnvironment.getMaxMessagesPerPoll());

    return pm;
  }

  @Bean
  public CronTrigger cronTrigger() {
    if (LOG.isDebugEnabled()) {
      return new CronTrigger(sftpEnvironment.getPollerCronExpressionWhenDebugModeIsEnabled());
    } else {
      return new CronTrigger(sftpEnvironment.getPollerCronExpression());
    }
  }

  @Bean
  public PeriodicTrigger periodicTrigger() {
    return new PeriodicTrigger(sftpEnvironment.getPeriodicTriggerInMillis());
  }

  @Bean
  public CompoundTrigger compoundTrigger() {
    return new CompoundTrigger(cronTrigger());
  }

  @Bean
  public CompoundTriggerAdvice compoundTriggerAdvice() {
    return new CompoundTriggerAdvice(compoundTrigger(), periodicTrigger());
  }

  @Bean
  public FileSystemPersistentAcceptOnceFileListFilter filelistFilter(MetadataStore datastore) {
    return new FileSystemPersistentAcceptOnceFileListFilter((JdbcMetadataStore) datastore, "INT");
  }

  @Bean
  public PlatformTransactionManager transactionManager() {
    return new org.springframework.integration.transaction.PseudoTransactionManager();
  }

  @Bean
  DataSource dataSource() throws SQLException {

    OracleDataSource dataSource = new OracleDataSource();
    dataSource.setUser(databaseProperties.getOracleUsername());
    dataSource.setPassword(databaseProperties.getOraclePassword());
    dataSource.setURL(databaseProperties.getOracleUrl());
    dataSource.setImplicitCachingEnabled(true);
    dataSource.setFastConnectionFailoverEnabled(true);
    return dataSource;
  }

  /**
   * Creates a {@link JdbcMetadataStore} for the de-duplication logic.
   *
   * <p>This method uses the "REGION" column of the metadatastore table to differentiate between
   * multiple apps. The value of the "REGION" column is set equal to the app-name.
   *
   * @return a JDBC metadata store
   * @throws SQLException in case an exception occurs during connection to SQL database
   */
  @Bean
  public MetadataStore metadataStore() throws SQLException {
    JdbcMetadataStore jdbcMetadataStore = new JdbcMetadataStore(dataSource());

  
    if (!Constants.app2.equals(appName)) {
      jdbcMetadataStore.setRegion(appName);
    }

    return jdbcMetadataStore;
  }

  @Bean
  @ServiceActivator(inputChannel = "sftpChannel")
  public MessageHandler handler() {
    return message -> {
      File file = (File) message.getPayload();
    
      FileDto fileDto = new FileDto(file);
      fileHandler.handle(fileDto);
      LOG.info("controller is here ");
      try {
        if (sftpEnvironment.isDeleteLocalFiles()) {
          Files.deleteIfExists(Paths.get(file.toString()));
        }
      } catch (IOException e) {
        // TODO retry/report/handle gracefully
        LOG.error(String.format("MessageHandler had error message=%s", message), e);
      }
    };
  }
}

内容的提问来源于stack exchange,提问作者arpit singh

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 20:35:27