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

Spring Integration动态配置SFTP目录及会话问题求助

多SFTP服务器Spring Integration改造问题与解决方案

问题背景

正在改造一个基于Spring Integration实现SFTP文件收发的项目,需新增多个SFTP服务器,实现按条件选择对应服务器并统一处理文件。改造过程中遇到两个核心问题:

  • 消息通道的消息头包含远程目录信息,需根据该信息获取对应的SFTP会话(依赖application.yml中定义的config1或config2配置),但不知如何将该信息传递给ServiceActivator;
  • 从多SFTP服务器拉取文件时,需根据对应SFTP会话动态设置本地存储目录(远程与本地路径不同,配置在config1和config2中)。

第一个问题的解决方案

已解决第一个问题,具体调整如下:

  • 在delegatingSf()方法中添加默认工厂;
  • 将delegatingSf()注入到ServiceActivator中,移除原sftpSession()方法;
  • 在发送消息到uploadChannel的方法中,通过delegatingSf.setThreadKey(mapSessionKey)动态选择工厂,发送完成后调用clearThreadKey()清理线程上下文。

第二个问题的现状

对于第二个问题,目前使用localFileNameExpression动态获取本地目录,但通过remoteDirectory识别服务器存在冲突风险,虽当前可正常运行,但方案不够理想。

初始代码

SftpConfig.Sftp1 sftpConfig1;
SftpConfig.Sftp2 sftpConfig2;

@Bean
@BridgeTo
public MessageChannel uploadChannel() {
    return new PublishSubscribeChannel();
}

@Bean
public ExpressionParser spelExpressionParser() {
    return new SpelExpressionParser();
}

public SessionFactory<ChannelSftp.LsEntry> getSftpSession(SftpConfig sftp) {

    DefaultSftpSessionFactory factory = new DefaultSftpSessionFactory(true);
    factory.setHost(sftp.getHost());
    factory.setPort(sftp.getPort());
    factory.setUser(sftp.getUser());
    factory.setPassword(sftp.getPassword());
    factory.setAllowUnknownKeys(true);
    factory.setTimeout(sftp.getTimeout());
    log.info("timeout is set to: {}", sftp.getTimeout());

    return new CachingSessionFactory<>(factory);
}

@Bean
public DelegatingSessionFactory<ChannelSftp.LsEntry> delegatingSf() {

    Map<Object, SessionFactory<ChannelSftp.LsEntry>> mapSession = new HashMap<>();
    mapSession.put("config1", getSftpSession(sftpConfig1));
    mapSession.put("config2", getSftpSession(sftpConfig2));

    SessionFactoryLocator<ChannelSftp.LsEntry> sessionFactoryLocator = new DefaultSessionFactoryLocator<>(mapSession);

    return new DelegatingSessionFactory<>(sessionFactoryLocator);
}

@Bean
public RotatingServerAdvice advice() {

    List<RotationPolicy.KeyDirectory> keyDirectories = sftpConfig1.getCodes().stream()
            .map(code -> new RotationPolicy.KeyDirectory("config1", sftpConfig1.getReaderDirectory() + SEPARATOR + code))
            .collect(Collectors.toList());

    keyDirectories.addAll(sftpConfig2.getCodes().stream()
            .map(code -> new RotationPolicy.KeyDirectory("config2", sftpConfig2.getReaderDirectory() + SEPARATOR + code))
            .collect(Collectors.toList()));

    return new RotatingServerAdvice(delegatingSf(), keyDirectories);
}

@Bean
public IntegrationFlow sftpIntegrationFlow() {
    return IntegrationFlows.from(
                    Sftp.inboundAdapter(delegatingSf())
                            .filter(new SftpSimplePatternFileListFilter("*.csv"))
                            .filter(new SftpPersistentAcceptOnceFileListFilter(new SimpleMetadataStore(), "rotate"))
                            .localFilter(new AbstractFileListFilter<File>() {
                                @Override
                                public boolean accept(final File file) {
                                    return file.getName().endsWith(".csv");
                                }
                            })
                            .deleteRemoteFiles(false)
                            .temporaryFileSuffix(".new")
                            .localDirectory(new File()) // TODO dynamic local directory based on sftp session
                            .remoteDirectory("."),
                                e -> e.poller(Pollers.fixedDelay(1, MINUTES).advice(advice()).advice(logNoFileFoundAdvice())))
            .log(LoggingHandler.Level.INFO, "[SFTP]", m -> "Received file: " + m.getHeaders().get(FileHeaders.FILENAME))
            .channel("filesReceptionChannel")
            .enrichHeaders(h -> h.header("errorChannel", "errorChannel"))
            .get();
}

@Bean
public MethodInterceptor logNoFileFoundAdvice() {
    return invocation -> {
        Object result = invocation.proceed();
        if (result == null) {
            log.info("[SFTP] No files found");
        }
        return result;
    };
}

@Bean
public SftpRemoteFileTemplate sftpTemplate() {
    return new SftpRemoteFileTemplate(sftpSession());
}

@Bean
public SessionFactory<ChannelSftp.LsEntry> sftpSession() {
    return getSftpSession(); // TODO dynamic sftp session based on message received in serviceActivator bellow
}

@Bean
@ServiceActivator(inputChannel = "uploadChannel")
public MessageHandler uploadHandler() {
    return getFtpMessageHandler(sftpSession());
}

public MessageHandler getFtpMessageHandler(SessionFactory<ChannelSftp.LsEntry> sftpSession) {
    SftpMessageHandler handler = new SftpMessageHandler(sftpSession);
    handler.setRemoteDirectoryExpressionString("headers['remoteDirectory']");
    handler.setFileNameGenerator(message -> {
        if (message.getPayload() instanceof File) {
            return ((File) message.getPayload()).getName();
        } else {
            throw new IllegalArgumentException("File expected as payload.");
        }
    });
    handler.setUseTemporaryFileName(false);
    return handler;
}

修改后代码

SftpConfig.Sftp1 sftpConfig1;
SftpConfig.Sftp2 sftpConfig2;

@Bean
@BridgeTo
public MessageChannel uploadChannel() {
    return new PublishSubscribeChannel();
}

@Bean
public ExpressionParser spelExpressionParser() {
    return new SpelExpressionParser();
}

public SessionFactory<ChannelSftp.LsEntry> getSftpSession(SftpConfig sftp) {

    DefaultSftpSessionFactory factory = new DefaultSftpSessionFactory(true);
    factory.setHost(sftp.getHost());
    factory.setPort(sftp.getPort());
    factory.setUser(sftp.getUser());
    factory.setPassword(sftp.getPassword());
    factory.setAllowUnknownKeys(true);
    factory.setTimeout(sftp.getTimeout());
    log.info("timeout is set to: {}", sftp.getTimeout());

    return new CachingSessionFactory<>(factory);
}

@Bean
public DelegatingSessionFactory<ChannelSftp.LsEntry> delegatingSf() {

    Map<Object, SessionFactory<ChannelSftp.LsEntry>> mapSession = new HashMap<>();
    mapSession.put("config1", getSftpSession(sftpConfig1));
    mapSession.put("config2", getSftpSession(sftpConfig2));

    SessionFactoryLocator<ChannelSftp.LsEntry> sessionFactoryLocator = new DefaultSessionFactoryLocator<>(mapSession, mapSession.get("config1"));

    return new DelegatingSessionFactory<>(sessionFactoryLocator);
}

@Bean
public RotatingServerAdvice advice() {

    List<RotationPolicy.KeyDirectory> keyDirectories = sftpConfig1.getCodes().stream()
            .map(code -> new RotationPolicy.KeyDirectory("config1", sftpConfig1.getReaderDirectory() + SEPARATOR + code))
            .collect(Collectors.toList());

    keyDirectories.addAll(sftpConfig2.getCodes().stream()
            .map(code -> new RotationPolicy.KeyDirectory("config2", sftpConfig2.getReaderDirectory() + SEPARATOR + code))
            .collect(Collectors.toList()));

    return new RotatingServerAdvice(delegatingSf(), keyDirectories);
}

@Bean
public IntegrationFlow sftpIntegrationFlow() {
    return IntegrationFlows.from(
                    Sftp.inboundAdapter(delegatingSf())
                            .filter(new SftpSimplePatternFileListFilter("*.csv"))
                            .filter(new SftpPersistentAcceptOnceFileListFilter(new SimpleMetadataStore(), "rotate"))
                            .localFilter(new AbstractFileListFilter<File>() {
                                @Override
                                public boolean accept(final File file) {
                                    return file.getName().endsWith(".csv");
                                }
                            })
                            .deleteRemoteFiles(false)
                            .temporaryFileSuffix(".new")
                            .localFilenameExpression("@sftpEIPConfig.getLocalDirectoryReader(#remoteDirectory) + #this")
                            .localDirectory(new File("/"))
                            .remoteDirectory("."),
                                e -> e.poller(Pollers.fixedDelay(1, MINUTES).advice(advice()).advice(logNoFileFoundAdvice())))
            .log(LoggingHandler.Level.INFO, "[SFTP]", m -> "Received file: " + m.getHeaders().get(FileHeaders.FILENAME))
            .channel("filesReceptionChannel")
            .enrichHeaders(h -> h.header("errorChannel", "errorChannel"))
            .get();
}

@Bean
public MethodInterceptor logNoFileFoundAdvice() {
    return invocation -> {
        Object result = invocation.proceed();
        if (result == null) {
            log.info("[SFTP] No files found");
        }
        return result;
    };
}

@Bean
public SftpRemoteFileTemplate sftpTemplate() {
    return new SftpRemoteFileTemplate(delegatingSf());
}

@Bean
@ServiceActivator(inputChannel = "uploadChannel")
public MessageHandler uploadHandler() {
    return getFtpMessageHandler(delegatingSf());
}

public MessageHandler getFtpMessageHandler(SessionFactory<ChannelSftp.LsEntry> sftpSession) {
    SftpMessageHandler handler = new SftpMessageHandler(sftpSession);
    handler.setRemoteDirectoryExpressionString("headers['remoteDirectory']");
    handler.setFileNameGenerator(message -> {
        if (message.getPayload() instanceof File) {
            return ((File) message.getPayload()).getName();
        } else {
            throw new IllegalArgumentException("File expected as payload.");
        }
    });
    handler.setUseTemporaryFileName(false);
    return handler;
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 10:47:03