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

Spring Integration运行时向MessageHandler传参及动态SFTP模板选择问题

问题描述

尝试使用Spring Integration搭建SFTP文件监听器,要求按规则轮询文件并在处理完成后重命名文件。由于需要运行时动态创建轮询器,无法硬编码主机、端口、用户名、密码及文件匹配规则等信息。

当前文件处理与重命名功能基本可用,但无法为入站轮询器和OutboundGateway在运行时动态选择SftpRemoteFileTemplate或sftpSessionFactory。曾尝试DelegatingSessionFactory但未与轮询器配合成功,遂自行实现内存缓存存储每个监听器对应的SftpRemoteFileTemplate,期望重命名时复用该模板。

核心问题:

  • 无法从消息头中获取值来匹配对应的SFTPRemoteTemplate,报错Property or field 'headers' cannot be found on null,推测是创建IntegrationFlow时消息头尚未生成。
  • 若通过handle方法传递消息头参数,SFTP文件重命名功能无法正常工作,且不清楚如何复用会话。

现有代码:

@Configuration
@Slf4j
public class InboundFtpConfiguration {

HashMap<Integer, SftpRemoteFileTemplate> sftpRemoteFileTemplateMap = new HashMap<>();

private static SessionFactory<ChannelSftp.LsEntry> createSFtpSessionFactory(String host, int port, String username, String password) {
    DefaultSftpSessionFactory sftpSessionFactory = new DefaultSftpSessionFactory();
    sftpSessionFactory.setHost(host);
    sftpSessionFactory.setPort(port);
    sftpSessionFactory.setUser(username);
    sftpSessionFactory.setPassword(password);
    
    java.util.Properties config = new java.util.Properties();
    config.put("StrictHostKeyChecking", "no");
    sftpSessionFactory.setSessionConfig(config);
    return sftpSessionFactory;
}


public String createInboundChannelAdapter(FTPObject ftpObject) {

    if(sftpRemoteFileTemplateMap.get(ftpObject.getKey()) == null){
        sftpRemoteFileTemplateMap.put(ftpObject.getKey(),
                new SftpRemoteFileTemplate(createSFtpSessionFactory(ftpObject.getHost(), ftpObject.getPort(), ftpObject.getUsername(), ftpObject.getPassword())
        ));
    }


    var flow = IntegrationFlows
            .from(Sftp.inboundStreamingAdapter(sftpRemoteFileTemplateMap.get(ftpObject.getKey()))
                            .remoteDirectory(ftpObject.getRemoteDirectory())
                            .regexFilter("(.*).txt"),
                    sourcePollingChannelAdapterSpec -> {
                        sourcePollingChannelAdapterSpec.poller(pollerFactory -> pollerFactory.fixedDelay(5000));
                    }
            )
            .transform(Transformers.fromStream())
            .enrichHeaders(headerEnricherSpec -> {
               headerEnricherSpec.headerExpression("TENANTID", String.valueOf(TenantContext.getTenantId()));
               headerEnricherSpec.headerExpression("key", String.valueOf(ftpObject.getKey()));
            })
            .publishSubscribeChannel(subFlow -> subFlow
                    .subscribe(flow1 -> flow1.handle(h ->{
                        System.out.println("(1) Do Something First");
                    }))
                
                   //.subscribe(flow2 -> flow2.handle((p,h) -> renameFile(m.getHeaders())))
                   // .subscribe(flow2 -> flow2.handle(renameFile(m.getHeaders())))
                    .subscribe(flow2 -> flow2.handle(renameFile()))
            )
            .log(LoggingHandler.Level.INFO)
            .get();


    return flowContext
            .registration(flow)
            .autoStartup(false)
            .register()
            .getId();
}

public MessageHandler renameFile(){
  
    //System.out.println("In rename file key is: "+new SpelExpressionParser().parseExpression("headers['key']").getValue());

    //ERROR generating line
    SftpOutboundGateway sftpOutboundGateway = new  SftpOutboundGateway(sftpRemoteFileTemplateMap.get((Integer) new SpelExpressionParser().parseExpression("headers['key']").getValue()),
            AbstractRemoteFileOutboundGateway.Command.MV.getCommand(),"headers['file_remoteDirectory'] + headers['file_remoteFile']");



    sftpOutboundGateway.setRenameExpressionString( "headers['file_remoteDirectory'] + headers['file_remoteFile']+'.processed'");
    sftpOutboundGateway.setRequiresReply(false);
    sftpOutboundGateway.setOutputChannelName("nullChannel");
    sftpOutboundGateway.setOrder(Ordered.LOWEST_PRECEDENCE);
    sftpOutboundGateway.setAsync(true);
    return sftpOutboundGateway;

   /* return Sftp.outboundGateway(sftpRemoteFileTemplateMap.get(1),
            AbstractRemoteFileOutboundGateway.Command.MV.getCommand(),"headers['file_remoteDirectory'] + headers['file_remoteFile']")
            .renameExpression("headers['file_remoteDirectory'] + headers['file_remoteFile']+'.processed'")
            .get();*/

}
}
解决思路

1. 根本错误原因

renameFile()方法在IntegrationFlow创建阶段被调用,此时还没有生成任何消息,直接解析SpEL表达式headers['key']必然会因为没有消息上下文而报错headers cannot be found on null。必须将消息头的获取逻辑延迟到消息处理阶段,而非流程初始化阶段。

2. 方案一:改用DelegatingSessionFactory实现动态会话选择

放弃手动维护SftpRemoteFileTemplate缓存,改用DelegatingSessionFactory统一管理会话工厂,并通过消息头动态切换:

@Bean
public DelegatingSessionFactory<ChannelSftp.LsEntry> delegatingSftpSessionFactory() {
    DelegatingSessionFactory<ChannelSftp.LsEntry> delegatingSessionFactory = new DelegatingSessionFactory<>();
    // 设置从消息头获取会话标识的策略
    delegatingSessionFactory.setThreadKeyStrategy(message -> message.getHeaders().get("key", Integer.class));
    return delegatingSessionFactory;
}

// 注册会话工厂到DelegatingSessionFactory
public void registerSftpSessionFactory(Integer key, String host, int port, String username, String password) {
    DefaultSftpSessionFactory sftpSessionFactory = new DefaultSftpSessionFactory();
    sftpSessionFactory.setHost(host);
    sftpSessionFactory.setPort(port);
    sftpSessionFactory.setUser(username);
    sftpSessionFactory.setPassword(password);
    
    java.util.Properties config = new java.util.Properties();
    config.put("StrictHostKeyChecking", "no");
    sftpSessionFactory.setSessionConfig(config);
    
    delegatingSftpSessionFactory.addSessionFactory(key, sftpSessionFactory);
}

修改流程创建逻辑,统一使用DelegatingSessionFactory:

public String createInboundChannelAdapter(FTPObject ftpObject) {
    if (!delegatingSftpSessionFactory.getSessionFactories().containsKey(ftpObject.getKey())) {
        registerSftpSessionFactory(ftpObject.getKey(), ftpObject.getHost(), ftpObject.getPort(), ftpObject.getUsername(), ftpObject.getPassword());
    }

    SftpRemoteFileTemplate template = new SftpRemoteFileTemplate(delegatingSftpSessionFactory);

    var flow = IntegrationFlows
            .from(Sftp.inboundStreamingAdapter(template)
                            .remoteDirectory(ftpObject.getRemoteDirectory())
                            .regexFilter("(.*).txt"),
                    sourcePollingChannelAdapterSpec -> pollerFactory.fixedDelay(5000)
            )
            .transform(Transformers.fromStream())
            .enrichHeaders(headerEnricherSpec -> {
               headerEnricherSpec.header("TENANTID", String.valueOf(TenantContext.getTenantId()));
               headerEnricherSpec.header("key", ftpObject.getKey()); // 直接设置标识,无需表达式
            })
            .publishSubscribeChannel(subFlow -> subFlow
                    .subscribe(flow1 -> flow1.handle(h -> System.out.println("(1) Do Something First")))
                    .subscribe(flow2 -> flow2.handle(Sftp.outboundGateway(delegatingSftpSessionFactory, "mv", "headers['file_remoteDirectory'] + headers['file_remoteFile']")
                            .renameExpression("headers['file_remoteDirectory'] + headers['file_remoteFile']+'.processed'")
                            .requiresReply(false)
                            .async(true)
                    ))
            )
            .log(LoggingHandler.Level.INFO)
            .get();

    return flowContext
            .registration(flow)
            .autoStartup(false)
            .register()
            .getId();
}

3. 方案二:手动在消息处理阶段获取Template并重命名

如果坚持使用自己的Template缓存,可将重命名逻辑改为方法调用,在消息处理时动态获取Template:

// 修改publishSubscribeChannel的订阅逻辑
.publishSubscribeChannel(subFlow -> subFlow
        .subscribe(flow1 -> flow1.handle(h -> System.out.println("(1) Do Something First")))
        .subscribe(flow2 -> flow2.handle(this::handleFileRename))
)

// 实现重命名方法
public void handleFileRename(String payload, MessageHeaders headers) {
    Integer key = headers.get("key", Integer.class);
    SftpRemoteFileTemplate template = sftpRemoteFileTemplateMap.get(key);
    
    String remoteDir = headers.get("file_remoteDirectory", String.class);
    String remoteFile = headers.get("file_remoteFile", String.class);
    String sourcePath = remoteDir + remoteFile;
    String targetPath = remoteDir + remoteFile + ".processed";
    
    // 调用Template的API完成重命名,自动复用会话
    template.rename(sourcePath, targetPath);
}

4. 会话复用说明

SftpRemoteFileTemplate内部会维护会话池,只要是同一个Template实例,多次调用API会自动复用活跃会话;使用DelegatingSessionFactory时,只要消息头的key一致,就会选择同一个底层会话工厂,同样能实现会话复用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 22:35:01