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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 03:49:59