Nifi1.6版本中List SFTP获取多文件后如何仅处理最新修改文件
NiFi 1.6 仅处理SFTP最新文件的实现方案
你当前遇到的属性丢失问题是MergeContent的默认属性保留策略导致的,且你本身的需求仅需要处理最新文件不需要合并文件,因此不需要走MergeContent流程,直接通过属性筛选即可实现需求,以下是适配NiFi 1.6版本的可行方案:
前置确认
ListSFTP处理器默认开启List File Attributes配置后,会为每个输出流文件携带ftp.lastModifiedTime属性,值为对应源文件最后修改时间的长整型时间戳,该属性会全程保留不需要额外处理。
推荐方案:自定义脚本实现(配置简单无额外依赖)
使用ExecuteScript处理器批量筛选最新文件,无需修改现有ListSFTP配置:
- 在ListSFTP下游连接
ExecuteScript处理器 - 处理器基础配置:
- 脚本语言选择
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) } - 脚本语言选择
- 调整处理器的
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
相关产品推荐
相关产品推荐

