Spring Integration SFTP 从消息头提取SessionFactory创建出站网关删除文件
问题背景
当前正在实现多SFTP源文件由同一集成流统一处理的工作流,初始流程已可正常运行:
- 可为每个SFTP源(
SFTP_DEDUPED_FILES_0、SFTP_DEDUPED_FILES_1)使用独立的SFTP SessionFactory列出远端目录文件 - 可将携带文件名信息的消息转发至统一通道
SFTP_FILES_CHANNEL - 可通过
FILE_PROCESSING_FLOW完成文件信息的处理逻辑
现有配置代码如下:
@Bean(SFTP_DEDUPED_FILES_0) public IntegrationFlow getDedupedFilesFromSftp(List<SftpConnectionInfo> connectionInfos) { return createSftpGetFlow(connectionInfos.get(0)); } @Bean(SFTP_DEDUPED_FILES_1) public IntegrationFlow getDedupedFilesFromSftp1(List<SftpConnectionInfo> connectionInfos) { return createSftpGetFlow(connectionInfos.get(1)); } private IntegrationFlow createSftpGetFlow(SftpConnectionInfo connectionInfo) { return IntegrationFlows .from(PROVIDE_POLL_FREQUENCY) .handle(listRemoteFiles(connectionInfo.sessionFactory(), connectionInfo.dedupingKeyPrefix(), connectionInfo.sftpRemotePath())) .split() .enrichHeaders(addFilePathMessageHeader(connectionInfo.sftpRemotePath())) .enrichHeaders(addSessionFactoryMessageHeader(connectionInfo.sessionFactory())) .channel(SFTP_FILES_CHANNEL) .get(); } @Bean(SFTP_FILES_CHANNEL) public SubscribableChannel sftpFilesChannel() { return new PublishSubscribeChannel(); } @Bean(FILE_PROCESSING_FLOW) public IntegrationFlow sftpGetFlow() { return IntegrationFlows .from(SFTP_FILES_CHANNEL) .log(logTheFilePath()) .log(logTheMessageHeaders()) .handle(deleteTheFile()) .get(); } private SftpOutboundGatewaySpec listRemoteFiles(SessionFactory<ChannelSftp.LsEntry> sessionFactory, String dedupingKeyPrefix, String sftpRemotePath) { return Sftp.outboundGateway(sessionFactory, LS, asExpression(sftpRemotePath)) .options(RECURSIVE, NAME_ONLY) .filter(onlyFilesWeHaveNotSeenYet(dedupingKeyPrefix)) .filter(onlyFiles()); }
现存问题
文件处理完成后需要删除对应SFTP服务器上的源文件,该操作需要匹配使用源对应的SessionFactory。目前已经将对应SessionFactory提前存入消息头,但无法正确提取该参数来创建执行删除操作的出站网关。
之前尝试实现的deleteTheFile()方法存在两个缺陷:每次处理消息都会新建网关,性能较差;且手动创建的网关未经过Spring容器初始化,触发了No beanFactory异常。
错误实现代码:
private MessageHandler deleteTheFile() { return message -> { SessionFactory<ChannelSftp.LsEntry> sessionFactory = (SessionFactory<ChannelSftp.LsEntry>) message.getHeaders().get(CUSTOM_HEADER_SESSION_FACTORY); SftpOutboundGatewaySpec gatewaySpec = Sftp.outboundGateway(sessionFactory, RM, "headers['" + CUSTOM_HEADER_REMOTE_FILE + "']"); gatewaySpec.get().handleMessage(message); }; }
异常堆栈:
java.lang.RuntimeException: No beanFactory at org.springframework.integration.expression.ExpressionUtils.createStandardEvaluationContext(ExpressionUtils.java:90) ~[spring-integration-core-5.5.10.jar:5.5.10] at org.springframework.integration.util.AbstractExpressionEvaluator.getEvaluationContext(AbstractExpressionEvaluator.java:111) ~[spring-integration-core-5.5.10.jar:5.5.10] at org.springframework.integration.util.AbstractExpressionEvaluator.getEvaluationContext(AbstractExpressionEvaluator.java:97) ~[spring-integration-core-5.5.10.jar:5.5.10] at org.springframework.integration.util.AbstractExpressionEvaluator.evaluateExpression(AbstractExpressionEvaluator.java:169) ~[spring-integration-core-5.5.10.jar:5.5.10] at org.springframework.integration.util.AbstractExpressionEvaluator.evaluateExpression(AbstractExpressionEvaluator.java:127) ~[spring-integration-core-5.5.10.jar:5.5.10] at org.springframework.integration.handler.ExpressionEvaluatingMessageProcessor.processMessage(ExpressionEvaluatingMessageProcessor.java:109) ~[spring-integration-core-5.5.10.jar:5.5.10] at org.springframework.integration.file.remote.gateway.AbstractRemoteFileOutboundGateway.obtainRemoteFilePath(AbstractRemoteFileOutboundGateway.java:760) ~[spring-integration-file-5.5.10.jar:5.5.10] at org.springframework.integration.file.remote.gateway.AbstractRemoteFileOutboundGateway.doRm(AbstractRemoteFileOutboundGateway.java:709) ~[spring-integration-file-5.5.10.jar:5.5.10] at org.springframework.integration.file.remote.gateway.AbstractRemoteFileOutboundGateway.handleRequestMessage(AbstractRemoteFileOutboundGateway.java:592) ~[spring-integration-file-5.5.10.jar:5.5.10] at org.springframework.integration.handler.AbstractReplyProducingMessageHandler.handleMessageInternal(AbstractReplyProducingMessageHandler.java:136) ~[spring-integration-core-5.5.10.jar:5.5.10] at org.springframework.integration.handler.AbstractMessageHandler.handleMessage(AbstractMessageHandler.java:56) ~[spring-integration-core-5.5.10.jar:5.5.10] at com.example.sftp.incoming.SftpMergedIncomingRecursiveConfiguration.lambda$deleteTheFile$4(SftpMergedIncomingRecursiveConfiguration.java:136) ~[classes/:na]
实现方案
不要使用每次手动实例化网关的写法,将删除操作的出站网关注册为Spring管理的单例Bean,通过SpEL表达式配置网关动态从消息头读取对应SFTP源的SessionFactory,既解决容器上下文缺失的报错,也避免重复创建实例的性能损耗。
- 首先定义删除操作的出站网关Bean,配置动态SessionFactory读取逻辑:
@Bean public SftpOutboundGateway sftpDeleteHandler() { // 构造时传入空的SessionFactory占位,后续通过表达式动态获取 SftpOutboundGateway deleteGateway = new SftpOutboundGateway( (SessionFactory<ChannelSftp.LsEntry>) null, AbstractRemoteFileOutboundGateway.Command.RM.getCommand(), "headers['" + CUSTOM_HEADER_REMOTE_FILE + "']" ); // 配置SpEL表达式,执行删除时从消息头取出对应SFTP源的SessionFactory SpelExpressionParser parser = new SpelExpressionParser(); deleteGateway.setSessionFactoryExpression( parser.parseExpression("headers['" + CUSTOM_HEADER_SESSION_FACTORY + "']") ); // 可选配置:删除时忽略文件不存在等非核心异常 deleteGateway.setOption(AbstractRemoteFileOutboundGateway.Option.SUPPRESS_EXCEPTIONS); return deleteGateway; }
- 修改文件处理流的配置,直接引用注册好的网关Bean作为处理器:
@Bean(FILE_PROCESSING_FLOW) public IntegrationFlow sftpGetFlow() { return IntegrationFlows .from(SFTP_FILES_CHANNEL) .log(logTheFilePath()) .log(logTheMessageHeaders()) .handle(sftpDeleteHandler()) .get(); }
注意:消息头中存储的SessionFactory直接使用对象引用即可,同应用内流转的消息不需要对SessionFactory做序列化处理。
内容的提问来源于stack exchange,提问作者RobertG
相关产品推荐
相关产品推荐

