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

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,既解决容器上下文缺失的报错,也避免重复创建实例的性能损耗。

  1. 首先定义删除操作的出站网关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;
}
  1. 修改文件处理流的配置,直接引用注册好的网关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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 21:18:24