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

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.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 20:44:52