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

Nifi1.6版本中List SFTP获取多文件后如何仅处理最新修改文件

NiFi 1.6 仅处理SFTP最新文件的实现方案

你当前遇到的属性丢失问题是MergeContent的默认属性保留策略导致的,且你本身的需求仅需要处理最新文件不需要合并文件,因此不需要走MergeContent流程,直接通过属性筛选即可实现需求,以下是适配NiFi 1.6版本的可行方案:

前置确认

ListSFTP处理器默认开启List File Attributes配置后,会为每个输出流文件携带ftp.lastModifiedTime属性,值为对应源文件最后修改时间的长整型时间戳,该属性会全程保留不需要额外处理。

推荐方案:自定义脚本实现(配置简单无额外依赖)

使用ExecuteScript处理器批量筛选最新文件,无需修改现有ListSFTP配置:

  1. 在ListSFTP下游连接ExecuteScript处理器
  2. 处理器基础配置:
    • 脚本语言选择Groovy
    • 脚本内容填写如下代码:
    // 单次拉取的最大文件数,可根据实际业务单次ListSFTP返回的文件量调整
    def flowFiles = session.get(1000)
    if (!flowFiles) return
    
    // 按ftp.lastModifiedTime升序排序后取最后一个,即为最新文件
    def latestFile = flowFiles.sort{ it.getAttribute('ftp.lastModifiedTime') as long }.last()
    
    // 最新文件路由到success分支进入后续处理
    session.transfer(latestFile, REL_SUCCESS)
    // 其余非最新文件可根据需求选择删除或路由到其他分支
    flowFiles.findAll{ it != latestFile }.each{ session.remove(it) }
    
  3. 调整处理器的Run Schedule,建议和ListSFTP的轮询间隔保持一致,确保单次轮询拉取的所有文件能被批量处理,避免分批拉取导致筛选逻辑错误。

可选无脚本方案(仅当不允许使用自定义脚本时使用)

如果环境不允许运行自定义脚本,可通过以下处理器组合实现:

  • 第一步:通过UpdateAttribute给所有流文件新增两个属性:
    • group_id:值为固定值sftp_batch,用于标记同一次轮询拉取的文件
    • reverse_modify_time:值为${ftp.lastModifiedTime:multiply(-1)},用于后续排序
  • 第二步:给ListSFTP到下一级处理器的队列配置PriorityAttributePrioritizer,排序依据设为reverse_modify_time升序,这样最新的文件会排在队列最前面
  • 第三步:连接DetectDuplicate处理器,Entry Identifier设为${group_id},缓存有效期设为和ListSFTP的轮询间隔一致,这样只有第一个进入处理器的最新文件会被放行,其余同批次的旧文件都会被识别为重复路由到duplicate分支。

注意事项

  • 如果需要保留非最新文件,可修改脚本逻辑将非最新文件路由到其他分支存储,而非直接删除
  • 筛选出的最新流文件完整保留了原始的ftp.lastModifiedTime属性,可直接用于后续的业务逻辑判断

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 12:15:03