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

Spring Batch Integration轮询SFTP时实例空闲异常排查

问题根因

你遇到的两类异常并非完全由服务端导致,核心是分布式场景下配置缺失、逻辑冲突触发的,具体原因如下:

  • 初始过滤器顺序错误:最早使用CompositeFileListFilter时将SftpPersistentAcceptOnceFileListFilter放在正则匹配过滤器之前,会导致所有非目标格式文件都被写入Redis元数据,同时多实例并发扫描时,过滤器的“检查-标记”操作非原子,多个实例会同时拉取同一个文件,出现删远端文件报不存在的情况,拉取失败的实例会触发默认退避逻辑,导致长时间空闲。后续调整为ChainFileListFilter仅解决了非目标文件写入元数据的问题,没有解决原子性问题。
  • 缺失分布式锁:SftpInboundFileSynchronizer默认未配置全局锁,6个实例并发扫描远端SFTP目录时,会出现同文件被多实例同时拉取的情况:成功拉取并删除文件的实例正常执行任务,其余实例因远端文件已被删除触发异常,同时会将已处理文件写入元数据,导致后续轮询时误判无待处理文件,出现数十分钟空闲。
  • SFTP连接缓存配置不合理:CachingSessionFactory未配置连接有效性校验、空闲连接回收规则,服务端主动关闭空闲连接后,缓存中残留的死连接会触发断连、重连流程,重连后文件列表缓存刷新延迟叠加元数据脏数据,就会出现重连后10分钟以上不处理文件的问题。
  • 大文件场景参数不匹配:配置max-fetch-size:3、max-messages-per-poll:5,但单文件大小达50-200MB,拉取耗时长,且默认轮询为单线程阻塞执行,拉取、业务处理串行,很容易导致轮询线程被阻塞,出现实例空闲的假象。
修复方案

1. 增加分布式锁,修正过滤器配置

通过Redis锁保证同一时间只有一个实例执行远端目录扫描,同时给持久化元数据设置过期时间避免脏数据堆积:

// 配置Redis分布式锁,锁超时需大于单文件最大拉取+处理时长,这里示例设为10分钟
@Bean
public RedisLockRegistry redisLockRegistry(RedisConnectionFactory redisConnectionFactory) {
    return new RedisLockRegistry(redisConnectionFactory, "sftp-poll-lock", 600000);
}

// 修正文件过滤器逻辑
private FileListFilter<ChannelSftp.LsEntry> sftpFileListFilter() {
    return new ChainFileListFilter<ChannelSftp.LsEntry>()
            .addFilter(new SftpRegexPatternFileListFilter("^HELLO.*\\.xml$"))
            // 元数据过期时间设为3天,避免永久堆积
            .addFilter(new SftpPersistentAcceptOnceFileListFilter(
                    new RedisMetadataStore(redisConnectionFactory),
                    "prefix-",
                    259200000
            ));
}

// 给同步器配置分布式锁
@Bean
public SftpInboundFileSynchronizer sftpInboundFileSynchronizer(SessionFactory<ChannelSftp.LsEntry> sftpSessionFactory,
                                                               RedisLockRegistry redisLockRegistry) {
    SftpInboundFileSynchronizer synchronizer = new SftpInboundFileSynchronizer(sftpSessionFactory);
    synchronizer.setPreserveTimestamp(true);
    synchronizer.setDeleteRemoteFiles(true);
    synchronizer.setFilter(sftpFileListFilter());
    synchronizer.setRemoteDirectory(integrationProperties.getSftpSourceDirectory());
    synchronizer.setLocalFilenameGeneratorExpression(PARSER.parseExpression("#this.substring(#this.lastIndexOf('/')+1)"));
    synchronizer.setLockRegistry(redisLockRegistry);
    // 大文件场景每次只拉1个文件,避免长连接占用超时
    synchronizer.setMaxFetchSize(1);
    // 关闭本地文件列表缓存,每次轮询实时检查
    synchronizer.setLocalFilter(new AcceptAllFileListFilter<>());
    return synchronizer;
}

2. 优化SFTP连接缓存配置

增加连接有效性校验、空闲连接回收规则,避免使用失效死连接:

@Bean
public SessionFactory<ChannelSftp.LsEntry> sftpSessionFactory() {
    DefaultSftpSessionFactory factory = new DefaultSftpSessionFactory();
    factory.setHost(integrationProperties.getSftpSourceHost());
    factory.setPort(integrationProperties.getSftpSourcePort());
    factory.setUser(integrationProperties.getSftpSourceUser());
    factory.setPrivateKey(new FileSystemResource(integrationProperties.getSftpSourcePrivateKeyLocation()));
    factory.setAllowUnknownKeys(true);
    // 设置连接、读写超时为30秒
    factory.setConnectTimeout(30000);
    factory.setTimeout(30000);

    CachingSessionFactory<ChannelSftp.LsEntry> cachingFactory = new CachingSessionFactory<>(factory, 10);
    cachingFactory.setSessionWaitTimeout(10000);
    // 每次从缓存取连接时校验有效性,过滤已被服务端关闭的死连接
    cachingFactory.setTestSession(true);
    return cachingFactory;
}

3. 调整轮询线程模型与参数

使用独立线程池执行轮询任务,避免业务处理阻塞轮询线程,同时自定义错误Handler避免异常触发长退避:

@Bean
public IntegrationFlow myIntegrationFlow(
        SftpInboundFileSynchronizingMessageSource sftpInboundFileSynchronizingMessageSource,
        ThreadPoolTaskExecutor sftpPollExecutor) {
    return IntegrationFlows
            .from(sftpInboundFileSynchronizingMessageSource,
                    c -> c.poller(Pollers.fixedDelay(10000)
                            .maxMessagesPerPoll(1)
                            .taskExecutor(sftpPollExecutor)
                            .errorHandler(t -> log.warn("SFTP轮询异常,下一轮次自动重试", t))))
            .transform(myFileMessageToJobRequest())
            .handle(myJobLaunchingGateway())
            .get();
}

// 独立轮询线程池配置
@Bean
public ThreadPoolTaskExecutor sftpPollExecutor() {
    ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
    executor.setCorePoolSize(6);
    executor.setMaxPoolSize(10);
    executor.setQueueCapacity(0);
    executor.setThreadNamePrefix("sftp-poll-");
    executor.initialize();
    return executor;
}

对应配置参数调整为:

poller-delay: 10000
max-message-per-poll: 1
max-fetch-size: 1
服务端侧排查点

完成以上配置调整后如果仍存在连接被主动断开的情况,可联系服务端管理员确认以下限制,对应调整本地参数适配即可:

  • 单IP最大连接数限制
  • 空闲连接超时时间
  • 单连接最大传输时长限制

内容的提问来源于stack exchange,提问作者Guilherme Bernardi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 23:15:11