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
相关产品推荐
相关产品推荐

