Spring Integration SFTP频繁出现failed to send channel request问题求助
Spring Integration SFTP 频繁触发"failed to send channel request"异常的解决方案
问题概述
我通过Spring Integration SFTP实现远程目录监听、文件拉取及远程文件删除功能,但频繁遭遇failed to send channel request异常。尝试更换为mwiede/jsch库后问题仍未解决,寻求可行的解决方案。
核心配置代码
SFTP会话工厂
@Bean public SessionFactory<ChannelSftp.LsEntry> sessionFactory() { final DefaultSftpSessionFactory factory = new DefaultSftpSessionFactory(false); factory.setHost("a-host"); factory.setPort(22); factory.setUser("username"); factory.setTimeout(30000); factory.setPrivateKey("private-key-location"); factory.setPrivateKeyPassphrase("passphrase"); factory.setAllowUnknownKeys(true); return factory; }
过滤规则与调度器
@Bean public ChainFileListFilter<ChannelSftp.LsEntry> testIncomingChainFileListFilter(final ConcurrentMetadataStore metadataStore) { final ChainFileListFilter<ChannelSftp.LsEntry> chainFileListFilter = new ChainFileListFilter<>(); chainFileListFilter.addFilter(new SftpRegexPatternFileListFilter("a-regex")); chainFileListFilter.addFilter(new SftpPersistentAcceptOnceFileListFilter(metadataStore, "a-prefix")); return chainFileListFilter; } @Bean public ThreadPoolTaskScheduler testTaskScheduler() { final ThreadPoolTaskScheduler taskScheduler = new ThreadPoolTaskScheduler(); taskScheduler.setAwaitTerminationSeconds(60); taskScheduler.setWaitForTasksToCompleteOnShutdown(true); taskScheduler.setPoolSize(1); taskScheduler.setThreadNamePrefix("A-Task-Scheduler"); return taskScheduler; }
主集成流(拉取文件)
@Bean public IntegrationFlow testIncomingIntegrationFlow(final ChainFileListFilter<ChannelSftp.LsEntry> testIncomingChainFileListFilter, final SessionFactory<ChannelSftp.LsEntry> sessionFactory, final ThreadPoolTaskScheduler testTaskScheduler, final PollSkipAdvice testIncomingPollSkipAdvice) { final List<AbstractRemoteFileOutboundGateway.Option> optionList = new ArrayList<>(); optionList.add(AbstractRemoteFileOutboundGateway.Option.PRESERVE_TIMESTAMP); optionList.add(AbstractRemoteFileOutboundGateway.Option.DELETE); final PollerSpec pollerSpec = getTestPollerSpec(testTaskScheduler, testIncomingPollSkipAdvice); return IntegrationFlows.from(() -> new GenericMessage<>("/a-remote-directory"), e -> e.poller(pollerSpec)) .split() .log(message -> "1- LS.getPayload (Flow: test, Source: a-prefix) -> " + message.getPayload()) .handle(Sftp.outboundGateway(sessionFactory, AbstractRemoteFileOutboundGateway.Command.LS, "payload") .options(AbstractRemoteFileOutboundGateway.Option.RECURSIVE, AbstractRemoteFileOutboundGateway.Option.NAME_ONLY) ) .split() .log(message -> "2- MGET.getPayload (Flow: test, Source: a-prefix)-> " + message.getPayload()) .handle(Sftp .outboundGateway(sessionFactory, AbstractRemoteFileOutboundGateway.Command.MGET, "'" + "/a-remote-directory" + "' + payload") .options(optionList.toArray(new AbstractRemoteFileOutboundGateway.Option[0])) .filter(testIncomingChainFileListFilter) .localDirectoryExpression("'" + System.getProperty("java.io.tmpdir") + File.separator + "a-local-working-directory" + File.separator + "a-prefix" + "' + #remoteDirectory") .temporaryFileSuffix(FILE_SUFFIX) .autoCreateLocalDirectory(true) .fileExistsMode(REPLACE) ) .handle(message -> fileService.handleMessage("a-prefix", "test", message)) .get(); } private PollerSpec getTestPollerSpec(final ThreadPoolTaskScheduler testTaskScheduler, final PollSkipAdvice testIncomingPollSkipAdvice) { return Pollers.cron("a-cron-expression") .maxMessagesPerPoll(1) .advice(testIncomingPollSkipAdvice) .taskExecutor(testTaskScheduler); }
另一个触发异常的集成流(入站适配器)
public IntegrationFlow anOtherTestIncomingIntegrationFlow(final ChainFileListFilter<ChannelSftp.LsEntry> anOtherTestIncomingChainFileListFilter, final ChainFileListFilter<File> anOtherTestIncomingLocalChainFileListFilter, final ThreadPoolTaskScheduler anOtherTestIncomingTaskScheduler, final SessionFactory<ChannelSftp.LsEntry> sessionFactory, final PollSkipAdvice anOtherTestIncomingPollSkipAdvice) { final String remoteDirectory = "another-remote-directory"; final String localDirectory = "another-local-directory"; final PollerSpec pollerSpec = getPollerSpec(anOtherTestIncomingTaskScheduler, anOtherTestIncomingPollSkipAdvice); return IntegrationFlows.from( Sftp.inboundAdapter(sessionFactory) .preserveTimestamp(true) .deleteRemoteFiles(true) .remoteDirectory(remoteDirectory) .localDirectory(Paths .get(localDirectory) .toAbsolutePath() .normalize() .toFile() ) .autoCreateLocalDirectory(true) .temporaryFileSuffix(FILE_SUFFIX) .maxFetchSize(100) .filter(anOtherTestIncomingChainFileListFilter) .localFilter(anOtherTestIncomingLocalChainFileListFilter) .autoCreateLocalDirectory(true), e -> e.id("anOtherTestBInboundAdapter") .autoStartup(true) .poller(pollerSpec)) .handle(message -> fileService.handleMessage("a-source", "another-test", message)) .get(); }
异常栈信息
Caused by: org.springframework.messaging.MessagingException: Failed to execute on session; nested exception is java.lang.IllegalStateException: failed to create SFTP Session at org.springframework.integration.file.remote.RemoteFileTemplate.execute(RemoteFileTemplate.java:461) at org.springframework.integration.file.remote.gateway.AbstractRemoteFileOutboundGateway.doLs(AbstractRemoteFileOutboundGateway.java:607) at org.springframework.integration.file.remote.gateway.AbstractRemoteFileOutboundGateway.handleRequestMessage(AbstractRemoteFileOutboundGateway.java:584) at org.springframework.integration.handler.AbstractReplyProducingMessageHandler.handleMessageInternal(AbstractReplyProducingMessageHandler.java:136) at org.springframework.integration.handler.AbstractMessageHandler.handleMessage(AbstractMessageHandler.java:55) ... 51 more Caused by: java.lang.IllegalStateException: failed to create SFTP Session at org.springframework.integration.sftp.session.DefaultSftpSessionFactory.getSession(DefaultSftpSessionFactory.java:403) at org.springframework.integration.sftp.session.DefaultSftpSessionFactory.getSession(DefaultSftpSessionFactory.java:60) at org.springframework.integration.file.remote.RemoteFileTemplate.execute(RemoteFileTemplate.java:447) ... 55 more Caused by: java.lang.IllegalStateException: failed to connect at org.springframework.integration.sftp.session.SftpSession.connect(SftpSession.java:301) at org.springframework.integration.sftp.session.DefaultSftpSessionFactory.getSession(DefaultSftpSessionFactory.java:397) ... 57 more Caused by: com.jcraft.jsch.JSchException: failed to send channel request at com.jcraft.jsch.Request.write(Request.java:65) at com.jcraft.jsch.RequestSftp.request(RequestSftp.java:47) at com.jcraft.jsch.ChannelSftp.start(ChannelSftp.java:237) at com.jcraft.jsch.Channel.connect(Channel.java:152) at org.springframework.integration.sftp.session.SftpSession.connect(SftpSession.java:296) ... 58 more
JSch日志片段(异常非必现,约10次请求触发1次)
2024 07 16 10:54:34.867##UTC##INFO####an-application####com.jcraft.jsch##org.springframework.integration.sftp.session.JschLogger##log##Scheduler1##Authentications that can continue: gssapi-with-mic,publickey,keyboard-interactive,password 2024 07 16 10:54:34.867##UTC##INFO####an-application####com.jcraft.jsch##org.springframework.integration.sftp.session.JschLogger##log##Scheduler1##Next authentication method: gssapi-with-mic 2024 07 16 10:54:34.874##UTC##INFO####an-application####com.jcraft.jsch##org.springframework.integration.sftp.session.JschLogger##log##Scheduler1##Authentications that can continue: publickey,keyboard-interactive,password 2024 07 16 10:54:34.874##UTC##INFO####an-application####com.jcraft.jsch##org.springframework.integration.sftp.session.JschLogger##log##Scheduler1##Next authentication method: publickey 2024 07 16 10:54:34.892##UTC##INFO####an-application####com.jcraft.jsch##org.springframework.integration.sftp.session.JschLogger##log##Scheduler1##Authentication succeeded (publickey). 2024 07 16 10:54:35.398##UTC##INFO####an-application####com.jcraft.jsch##org.springframework.integration.sftp.session.JschLogger##log##Scheduler1##Disconnecting from a-host port 22 2024 07 16 10:54:35.398##UTC##DEBUG####an-application####com.application.test##Scheduler1: incoming message payload -> [] 2024 07 16 10:54:35.398##UTC##INFO####an-application####com.jcraft.jsch##org.springframework.integration.sftp.session.JschLogger##log##Connect thread a-host session##Caught an exception, leaving main loop due to Socket closed
解决方案
1. 启用SFTP会话池
频繁创建新会话是触发此异常的常见原因,通过CachingSessionFactory复用已建立的会话:
@Bean public SessionFactory<ChannelSftp.LsEntry> sessionFactory() { DefaultSftpSessionFactory factory = new DefaultSftpSessionFactory(false); // 保留原配置 factory.setHost("a-host"); factory.setPort(22); factory.setUser("username"); factory.setTimeout(30000); factory.setPrivateKey("private-key-location"); factory.setPrivateKeyPassphrase("passphrase"); factory.setAllowUnknownKeys(true); // 添加会话池,设置合理的池大小(根据并发需求调整) return new CachingSessionFactory<>(factory, 5); }
2. 优化JSch连接配置
禁用GSSAPI认证(避免不必要的认证流程),并添加会话保活机制:
@Bean public SessionFactory<ChannelSftp.LsEntry> sessionFactory() { DefaultSftpSessionFactory factory = new DefaultSftpSessionFactory(false); // 原配置不变... // 禁用GSSAPI,优先使用公钥认证 Properties config = new Properties(); config.put("PreferredAuthentications", "publickey,keyboard-interactive,password"); // 设置连接超时与会话保活 config.put("ConnectTimeout", "30000"); config.put("ServerAliveInterval", "15000"); // 每15秒发送一次保活包 config.put("ServerAliveCountMax", "3"); // 连续3次未响应则断开 factory.setConfig(config); return new CachingSessionFactory<>(factory, 5); }
3. 添加重试机制
在Poller中配置重试策略,针对连接异常自动重试:
private PollerSpec getTestPollerSpec(final ThreadPoolTaskScheduler testTaskScheduler, final PollSkipAdvice testIncomingPollSkipAdvice) { // 构建重试模板 RetryTemplate retryTemplate = new RetryTemplate(); SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy(); retryPolicy.setMaxAttempts(3); // 最多重试3次 // 指定仅对连接相关异常重试 Set<Class<? extends Throwable>> retryableExceptions = new HashSet<>(); retryableExceptions.add(JSchException.class); retryableExceptions.add(IllegalStateException.class); retryPolicy.setRetryableExceptions(retryableExceptions); retryTemplate.setRetryPolicy(retryPolicy); return Pollers.cron("a-cron-expression") .maxMessagesPerPoll(1) .advice(testIncomingPollSkipAdvice) .advice(new RetryAdvice(retryTemplate)) // 添加重试增强 .taskExecutor(testTaskScheduler); }
4. 检查SFTP服务器端配置
联系运维人员检查SFTP服务器的sshd_config配置:
- 调整
MaxSessions参数,允许更多并发会话 - 修改
MaxStartups参数,放宽新连接的限制 - 确认服务器端是否有防火墙或连接超时的限制
5. 调整线程池配置
当前调度器线程池大小为1,若并发请求较多,可适当增大线程池容量:
@Bean public ThreadPoolTaskScheduler testTaskScheduler() { final ThreadPoolTaskScheduler taskScheduler = new ThreadPoolTaskScheduler(); taskScheduler.setAwaitTerminationSeconds(60); taskScheduler.setWaitForTasksToCompleteOnShutdown(true); taskScheduler.setPoolSize(3); // 根据实际并发需求调整 taskScheduler.setThreadNamePrefix("A-Task-Scheduler"); return taskScheduler; }
内容的提问来源于stack exchange,提问作者st.
相关产品推荐
相关产品推荐

