如何通过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
相关产品推荐
相关产品推荐

