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

如何通过Sftp.InboundStreamingAdapter实现处理后删除远程SFTP文件

问题解答

核心结论

Sftp.inboundStreamingAdapter本身没有像Sftp.inboundAdapter那样的自动删除配置属性,需要额外添加流程来实现远程文件的删除——因为流式适配器只是读取文件内容为流,不会自动管理远程文件的生命周期。

实现方案

你需要利用消息头中携带的远程文件信息(file_remoteDirectory和file_remoteFile),在整个文件处理完成后(所有CSV行都成功发送到Kafka),调用Sftp.outboundGateway执行删除命令。由于你使用了文件拆分器(Files.splitter())将文件拆分为多行处理,必须通过聚合器等待所有行处理完成并收到文件结束标记后,再触发删除操作。

具体代码修改

1. 定义处理完成后的通道

@Bean
public MessageChannel kafkaProcessedChannel() {
    return new DirectChannel();
}

2. 修改Kafka发布流程,将成功处理的消息导向聚合通道

@Bean
public IntegrationFlow publishToKafkaFlow(KafkaTemplate<String, String> kafkaTemplate,
                                          MessageChannel kafkaProducerErrorRecordChannel,
                                          QueueChannel kafkaPojoMessageChannel,
                                          MessageChannel kafkaProcessedChannel) {

    return IntegrationFlow.from(kafkaPojoMessageChannel)
                          .log(LoggingHandler.Level.DEBUG,
                               "DataSftpToKafkaIntegrationFlow", e -> "Payload: " + e.getPayload())
                          .handle(Kafka.outboundChannelAdapter(kafkaTemplate)
                                       .topic(KAFKA_TOPIC),
                                  e -> e.id("KafkaProducer"))
                          .routeByException(r -> r
                                  .channelMapping(KafkaProducerException.class, kafkaProducerErrorRecordChannel)
                                  .defaultOutputChannel("errorChannel"))
                          .channel(kafkaProcessedChannel) // 成功处理后发送到聚合通道
                          .get();
}

3. 添加聚合与远程文件删除流程

@Bean
public IntegrationFlow deleteRemoteFileFlow(SessionFactory<SftpClient.DirEntry> sftpSessionFactory) {
    return IntegrationFlow.from("kafkaProcessedChannel")
            // 按远程文件名聚合同一文件的所有消息
            .aggregate(a -> a
                    .correlationExpression("headers['file_remoteFile']")
                    // 收到文件结束标记时释放聚合结果
                    .releaseStrategy(g -> g.getMessages().stream()
                            .anyMatch(m -> m.getPayload() instanceof FileSplitter.FileMarker
                                    && ((FileSplitter.FileMarker) m.getPayload()).getMark() == FileSplitter.FileMarker.Mark.END))
                    .expireGroupsUponCompletion(true) // 处理完成后销毁分组
                    .sendPartialResultOnExpiry(false))
            // 调用SFTP网关执行删除命令
            .handle(Sftp.outboundGateway(sftpSessionFactory, AbstractRemoteFileOutboundGateway.Command.DELETE,
                    "headers['file_remoteDirectory'] + '/' + headers['file_remoteFile']"))
            .log(LoggingHandler.Level.INFO, "DataSftpToKafkaIntegrationFlow",
                 "Deleted remote file: headers['file_remoteDirectory'] + '/' + headers['file_remoteFile']")
            .get();
}

关键说明

  • 聚合器作用:因为你通过Files.splitter()将文件拆分为多行处理,聚合器会等待同一文件的所有行处理完成,直到收到FileMarker.Mark.END标记后,才执行删除操作,避免文件未处理完就被删除。
  • 消息头信息:Sftp.inboundStreamingAdapter会自动将远程文件的路径和名称放入消息头(file_remoteDirectory、file_remoteFile),直接用来构建删除命令的路径即可。
  • 异常处理:如果Kafka发送失败,消息会被路由到错误通道,不会触发删除,保证只有处理成功的文件才会被删除。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 13:43:30