Spring SFTP Integration问题:重启删目录+重复处理已移动文件
Spring Integration SFTP 文件处理异常问题排查与解决
问题描述
我是Spring Integration新手,需求是从SFTP服务器的user1/upload目录处理文件,处理完成后移动到user1/processed目录。当前代码整体可运行,但存在两个异常:
- 重启应用时,
user1/processed目录及其中原有文件会被删除,期望仅追加新文件而非清空目录; - 每次启动应用都会接收到已移动到
processed目录的旧文件(控制台可打印文件名),但WinSCP等工具连接SFTP时看不到这些文件,怀疑存在旧文件列表缓存。
相关配置代码
@Value("${cielo.sftp.host}") private String sftpHost; @Value("${cielo.sftp.port}") private int sftpPort; @Value("${cielo.sftp.user}") private String sftpUser; @Value("${cielo.sftp.pass}") private String sftpPasword; @Value("${cielo.sftp.remotedir}") private String sftpRemoteDirectoryDownload; @Value("${cielo.sftp.localdir}") private String sftpLocalDirectoryDownload; @Value("${cielo.sftp.filter}") private String sftpRemoteDirectoryDownloadFilter; @Bean @Order(Ordered.HIGHEST_PRECEDENCE) public SessionFactory<LsEntry> sftpSessionFactory() { DefaultSftpSessionFactory factory = new DefaultSftpSessionFactory(true); factory.setHost(sftpHost); factory.setPort(sftpPort); factory.setUser(sftpUser); factory.setPassword(sftpPasword); factory.setAllowUnknownKeys(true); //Set to true to allow connections to hosts with unknown (or changed) keys. Its default is 'false'. If false, a pre-populated knownHosts file is required. return new CachingSessionFactory<>(factory); } @Bean @Order(Ordered.HIGHEST_PRECEDENCE - 1) public SftpInboundFileSynchronizer sftpInboundFileSynchronizer(final SessionFactory<LsEntry> sftpSessionFactory) { SftpInboundFileSynchronizer fileSynchronizer = new SftpInboundFileSynchronizer(sftpSessionFactory); fileSynchronizer.setDeleteRemoteFiles(false); fileSynchronizer.setRemoteDirectory(sftpRemoteDirectoryDownload); fileSynchronizer .setFilter(new SftpSimplePatternFileListFilter(sftpRemoteDirectoryDownloadFilter)); //todo maybe use RegexPatternFileListFilter? return fileSynchronizer; } @Bean @Order(Ordered.HIGHEST_PRECEDENCE - 2) @InboundChannelAdapter(channel = "fromSftpChannel", poller = @Poller(fixedDelay = "1000")) //@InboundChannelAdapter(channel = "fromSftpChannel", poller = @Poller(cron = "${cielo.sftp.poller.cron}")) public MessageSource<File> sftpMessageSource(final SftpInboundFileSynchronizer sftpInboundFileSynchronizer) { SftpInboundFileSynchronizingMessageSource source = new SftpInboundFileSynchronizingMessageSource(sftpInboundFileSynchronizer); source.setLocalDirectory(new File("/tmp/local")); source.setAutoCreateLocalDirectory(true); source.setLocalFilter(new AcceptOnceFileListFilter<File>()); return source; } @Bean @ServiceActivator( inputChannel = "fromSftpChannel") public MessageHandler resultFileHandler() { return new MessageHandler() { @Override public void handleMessage(final Message<?> message) throws MessagingException { String payload = String.valueOf(message.getPayload()); System.err.println(payload); } }; } private static final SpelExpressionParser PARSER = new SpelExpressionParser(); @Bean(name="fromSftpChannel") public MessageChannel fromSftpChannel() { return new PublishSubscribeChannel(); } @Bean @ServiceActivator(inputChannel = "fromSftpChannel") @Order(Ordered.LOWEST_PRECEDENCE) public MessageHandler moveFile() { SftpOutboundGateway sftpOutboundGateway = new SftpOutboundGateway(sftpSessionFactory(), AbstractRemoteFileOutboundGateway.Command.MV.getCommand(), "'/user1/upload/'.concat(" + PARSER.parseExpression("payload.getName()").getExpressionString() + ")"); sftpOutboundGateway.setRenameExpressionString("'/user1/processed/'.concat(" + PARSER.parseExpression("payload.getName()").getExpressionString() + ")"); sftpOutboundGateway.setRequiresReply(false); sftpOutboundGateway.setOutputChannelName("nullChannel"); sftpOutboundGateway.setOrder(Ordered.LOWEST_PRECEDENCE); sftpOutboundGateway.setAsync(true); return sftpOutboundGateway; }
问题解决方法
1. 修复重启时processed目录被清空的问题
问题根源在于文件移动逻辑的路径拼接风险,以及缺少目录存在性校验:
- 优化文件移动路径表达式:避免硬编码路径拼接可能导致的目录覆盖,改用从消息头获取远程文件的原始目录,确保移动路径准确:
// 替换moveFile方法中的网关初始化代码 SftpOutboundGateway sftpOutboundGateway = new SftpOutboundGateway(sftpSessionFactory(), AbstractRemoteFileOutboundGateway.Command.MV.getCommand(), "headers['file_remoteDirectory'] + '/' + payload.getName()"); sftpOutboundGateway.setRenameExpressionString("'/user1/processed/' + payload.getName()"); - 确保
processed目录存在:在应用启动时,通过SFTP会话手动创建user1/processed目录(如果不存在),避免因目录不存在导致的文件移动失败或目录被误创建覆盖:// 在配置类中添加初始化方法 @PostConstruct public void initRemoteDirectories() { try (Session<LsEntry> session = sftpSessionFactory().getSession()) { String processedDir = "/user1/processed"; if (!session.exists(processedDir)) { session.mkdir(processedDir); } } catch (Exception e) { e.printStackTrace(); } }
2. 修复重启后重复接收已处理文件的问题
核心是替换内存型过滤器为持久化过滤器,跟踪已处理的远程文件:
- 使用持久化远程文件过滤器:将
SftpSimplePatternFileListFilter替换为SftpPersistentAcceptOnceFileListFilter,它会把已处理文件的元数据持久化,重启后不会重复扫描:// 修改sftpInboundFileSynchronizer的过滤器配置 @Bean @Order(Ordered.HIGHEST_PRECEDENCE - 1) public SftpInboundFileSynchronizer sftpInboundFileSynchronizer(final SessionFactory<LsEntry> sftpSessionFactory) { SftpInboundFileSynchronizer fileSynchronizer = new SftpInboundFileSynchronizer(sftpSessionFactory); fileSynchronizer.setDeleteRemoteFiles(false); fileSynchronizer.setRemoteDirectory(sftpRemoteDirectoryDownload); // 组合文件名过滤和持久化已处理过滤 CompositeFileListFilter<LsEntry> compositeFilter = new CompositeFileListFilter<>(); compositeFilter.addFilter(new SftpSimplePatternFileListFilter(sftpRemoteDirectoryDownloadFilter)); // 使用SimpleMetadataStore持久化,也可以替换为RedisMetadataStore等分布式存储 compositeFilter.addFilter(new SftpPersistentAcceptOnceFileListFilter(new SimpleMetadataStore(), "sftp_processed_tracker")); fileSynchronizer.setFilter(compositeFilter); return fileSynchronizer; } - 移除异步移动配置:当前
moveFile处理器设置了setAsync(true),可能导致文件未完成移动就被下一次轮询扫描到,建议去掉该配置,确保文件移动完成后再进行下一次轮询:// 移除这行代码 // sftpOutboundGateway.setAsync(true); - 调整本地目录:
/tmp/local目录可能在系统重启后被清理,建议改用自定义的持久化本地目录,避免本地缓存丢失导致的重复同步。
内容的提问来源于stack exchange,提问作者Nataliya K
相关产品推荐
相关产品推荐

