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

Spring Integration SFTP轮询:仅获文件名无需同步本地的入站适配器方案

方案可行性反馈:SftpStreamingMessageSource完全符合需求

核心结论

Sftp.inboundStreamingAdapter生成的SftpStreamingMessageSource是完美匹配你需求的开箱即用方案,无需大量自定义代码,也不用修改现有Spring Batch任务。

为什么可行

  • 不会同步远程文件到本地,仅返回包含远程文件元数据和输入流的消息,完全满足你“仅获取文件名/路径触发任务”的要求
  • 可以直接从消息中提取远程文件路径,替换原本地文件路径参数,无缝对接现有Spring Batch任务
  • 输入流可直接关闭,不会占用不必要的资源

修改示例代码

1. 调整Integration Flow配置

替换原本地文件入站适配器为SFTP流适配器:

@Bean
public IntegrationFlow sftpPollingFlow(SftpRemoteFileTemplate sftpRemoteFileTemplate,
                                       JobLaunchingGateway jobLaunchingGateway,
                                       FileMessageToJobRequest fileMessageToJobRequest,
                                       Executor jobsTaskExecutor) {
    return IntegrationFlows.from(
            Sftp.inboundStreamingAdapter(sftpRemoteFileTemplate)
                    .remoteDirectory(properties.getSftpInputDir())
                    .filter(new AcceptOnceFileListFilter<>()),
            c -> c.poller(Pollers.fixedRate(300, TimeUnit.SECONDS).maxMessagesPerPoll(50)))
            .log(LoggingHandler.Level.INFO, "Inb_GW_Msg", m -> m.getHeaders().get(FileHeaders.REMOTE_FILE))
            .channel(c -> c.executor(jobsTaskExecutor))
            .transform(fileMessageToJobRequest)
            .handle(jobLaunchingGateway, e -> e.advice(jobExecutionAdvice()))
            .log(LoggingHandler.Level.INFO, "Inb_GW_Msg_Processing_Result")
            .get();
}

注:sftpRemoteFileTemplate需要提前配置好SFTP会话工厂(包含主机、端口、认证信息等)

2. 修改转换器适配SFTP元数据

仅需调整转换器获取远程文件路径,无需修改Spring Batch任务逻辑:

@Transformer
public JobLaunchRequest toRequest(Message<InputStream> message) {
    // 从消息头获取远程文件元数据
    SftpFileInfo fileInfo = message.getHeaders().get(FileHeaders.REMOTE_FILE_INFO, SftpFileInfo.class);
    String remoteFilePath = fileInfo.getRemoteDirectory() + "/" + fileInfo.getFilename();

    // 关闭输入流(避免资源泄漏)
    try (InputStream ignored = message.getPayload()) {
        // 复用原Job参数构建逻辑
        JobParametersBuilder jobParametersBuilder = new JobParametersBuilder();
        jobParametersBuilder.addString(fileParameterName, remoteFilePath);
        return new JobLaunchRequest(job, jobParametersBuilder.toJobParameters());
    } catch (IOException e) {
        throw new RuntimeException("Failed to close SFTP input stream", e);
    }
}

关键注意事项

  • 过滤器选择:如果需要重启后不重复处理已扫描文件,建议替换AcceptOnceFileListFilter为SftpPersistentAcceptOnceFileListFilter,配合持久化存储(如Redis、数据库)记录已处理文件状态
  • 流资源管理:必须确保输入流被关闭,示例中用try-with-resources自动关闭,也可以通过Spring Integration的StreamTransformer或PayloadTypeRouter配合自动关闭机制
  • SFTP权限:确保SFTP账号有权限读取目标目录的文件列表,避免权限不足导致的扫描失败

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 03:18:15