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

如何在SFTP文件处理完毕后停止IntegrationFlow并每日定时重启?

解决Spring Integration SFTP定时触发后续无操作问题

问题描述

我实现了一个StandardIntegrationFlow,从SFTP服务器的/home/{userId}/files/...路径递归读取文件,过滤后处理符合条件的文件。已创建RecursiveSftpRemoteFileTemplate,可成功解析所有用户的文件。通过@Scheduled注解每日0点触发Flow启动,但首次触发能处理所有文件,后续触发无任何操作。需求是处理完所有文件后停止Flow,让下次定时触发时能重新遍历处理所有文件。使用Spring Integration版本为5.5.15。

原实现代码:

集成Flow定义

@Bean
public StandardIntegrationFlow readingFilesAndProcess(
    SessionFactory<ChannelSftp.LsEntry> sessionFactory,
    PollerMetadata sftpStreamingPoller,
    RecursiveSftpRemoteFileTemplate recursiveReceiptsRemoteFileTemplate,
    FileeValidator fileValidator) {
  return IntegrationFlows.from(
          Sftp.inboundStreamingAdapter(recursiveReceiptsRemoteFileTemplate)
              .remoteDirectory("/home"),
          e -> e.poller(sftpStreamingPoller).autoStartup(false))
      .filter(Message.class, fileValidator::isFileValid)
      .handle(
          message -> {
            log.info("Handling message at: {}", LocalDateTime.now());
            handlingFile(message);
          })
      .get();
}

递归SFTP模板定义

@Bean
public RecursiveSftpRemoteFileTemplate recursiveReceiptsRemoteFileTemplate(
    SessionFactory<ChannelSftp.LsEntry> sessionFactory,
    SftpPersistentAcceptOnceFileListFilter sftpPersistentAcceptOnceFileListFilter) {
  final RecursiveSftpRemoteFileTemplate template =
      new RecursiveSftpRemoteFileTemplate(
          sessionFactory,
          1,
          file -> sftpPersistentAcceptOnceFileListFilter.accept(file),
          List.of("files"));
  template.setRemoteDirectoryExpression(new LiteralExpression("/home"));
  return template;
}

定时触发方法

@Scheduled(cron = "0 0 0 * * *", zone = "ECT")
public void triggeringFlow(){
    log.info("Triggering the IntegrationFlow");
    readingFilesAndProcess.start();
}

问题根源

  1. SftpPersistentAcceptOnceFileListFilter会持久化已处理文件的记录,后续触发时会自动跳过所有已处理过的文件
  2. 当前流程启动后不会自动停止,持续轮询但无新文件可处理,导致后续定时触发看似无操作

解决方案

1. 重置文件过滤器记录

每次定时触发前,清除过滤器的已处理文件记录,确保重新扫描所有文件:

@Autowired
private SftpPersistentAcceptOnceFileListFilter sftpPersistentAcceptOnceFileListFilter;
@Autowired
private StandardIntegrationFlow readingFilesAndProcess;

@Scheduled(cron = "0 0 0 * * *", zone = "ECT")
public void triggeringFlow(){
    log.info("Triggering the IntegrationFlow at: {}", LocalDateTime.now());
    // 清除已处理文件的持久化记录
    sftpPersistentAcceptOnceFileListFilter.clear();
    // 仅当流程未运行时启动,避免重复启动
    if (!readingFilesAndProcess.isRunning()) {
        readingFilesAndProcess.start();
    }
}

2. 处理完文件后自动停止流程

通过给Poller添加Advice,检测到轮询无新文件时自动停止Flow:

@Bean
public PollerMetadata sftpStreamingPoller() {
    return Pollers.fixedDelay(5000)
            .advice(stopFlowAdvice())
            .maxMessagesPerPoll(1) // 每次轮询处理1个文件,可按需调整
            .get();
}

@Bean
public Advice stopFlowAdvice() {
    return new AbstractMessageAdvice() {
        @Override
        protected Object doInvoke(MethodInvocation invocation, Message<?> message) throws Throwable {
            Object result = invocation.proceed();
            // 轮询返回null表示无更多文件,停止流程
            if (result == null) {
                readingFilesAndProcess.stop();
                log.info("IntegrationFlow stopped after processing all files");
            }
            return result;
        }
    };
}

关键说明

  • SftpPersistentAcceptOnceFileListFilter.clear()会清除所有持久化的文件标记,确保每次触发都能重新扫描全量文件
  • maxMessagesPerPoll控制单次轮询处理的文件数量,可根据业务场景调整
  • 添加isRunning()判断避免重复启动流程,防止资源浪费

内容的提问来源于stack exchange,提问作者Vasilis Iak

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 20:01:05