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

