如何在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(); }
问题根源
SftpPersistentAcceptOnceFileListFilter会持久化已处理文件的记录,后续触发时会自动跳过所有已处理过的文件- 当前流程启动后不会自动停止,持续轮询但无新文件可处理,导致后续定时触发看似无操作
解决方案
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
相关产品推荐
相关产品推荐

