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
相关产品推荐
相关产品推荐

