基于PostgreSQL配置的Apache NiFi动态文件传输方案问询
基于PostgreSQL配置的Apache NiFi动态文件传输解决方案
针对你遇到的ListFile无法接收上游FlowFile、FetchFile不支持匹配模式的问题,可通过脚本处理器动态扫描文件的方式实现需求,具体流程及处理器配置如下:
完整流程链路
QueryDatabaseTableRecord → SplitRecord → EvaluateJsonPath → ExecuteGroovyScript(或ExecutePythonScript) → FetchFile → PutFile
1. 保留原有基础流程
QueryDatabaseTableRecord:读取PostgreSQL中的配置记录,输出包含所有配置的FlowFile。SplitRecord:将多记录的FlowFile拆分为单条配置的FlowFile,每条对应一组源路径、匹配模式、目标路径。EvaluateJsonPath:提取配置字段为FlowFile属性,建议设置:- 源路径 → 属性
source_dir - 文件匹配模式 → 属性
file_pattern(注意:若数据库中是通配符如*.txt,后续需转成正则格式) - 目标路径 → 属性
target_dir
- 源路径 → 属性
2. 核心:用脚本处理器动态扫描符合模式的文件
使用ExecuteGroovyScript处理器,基于上游传入的配置属性,扫描指定目录下的匹配文件,并为每个文件生成独立的FlowFile:
配置参数
- 脚本内容如下(直接复制到
Script属性):
import groovy.io.FileType import org.apache.nifi.processor.io.StreamCallback import java.nio.charset.StandardCharsets def flowFile = session.get() if (!flowFile) return // 读取上游传入的配置属性 def sourceDir = flowFile.getAttribute('source_dir') def filePattern = flowFile.getAttribute('file_pattern') def targetDir = flowFile.getAttribute('target_dir') def dir = new File(sourceDir) // 校验源目录合法性 if (!dir.exists() || !dir.isDirectory()) { session.transfer(flowFile, REL_FAILURE) return } // 转换通配符为正则(如果数据库中用的是*.txt这类通配符) def regexPattern = filePattern.replace('.', '\\.').replace('*', '.*') def matcher = ~regexPattern def matchedFiles = [] dir.eachFileMatch(FileType.FILES, matcher) { file -> matchedFiles.add(file.absolutePath) } // 无匹配文件时直接流转原FlowFile if (matchedFiles.isEmpty()) { session.transfer(flowFile, REL_SUCCESS) return } // 为每个匹配文件生成新FlowFile,携带文件路径和目标路径属性 matchedFiles.each { filePath -> def newFlowFile = session.create(flowFile) newFlowFile = session.putAttribute(newFlowFile, 'file_path', filePath) newFlowFile = session.putAttribute(newFlowFile, 'target_dir', targetDir) session.transfer(newFlowFile, REL_SUCCESS) } // 移除原配置FlowFile session.remove(flowFile)
- 配置关系:默认
REL_SUCCESS和REL_FAILURE,分别流转到后续处理器和失败队列。
3. 读取并传输文件
FetchFile:配置File属性为${file_path},动态引用脚本生成的文件绝对路径,实现精确读取单个文件。PutFile:配置Directory属性为${target_dir},将读取到的文件写入目标路径。
关键注意事项
- 匹配模式转换:如果数据库中存储的是通配符(如
*.csv),脚本中已做简单转换,若有复杂匹配规则,需自行调整正则转换逻辑。 - 权限控制:确保NiFi运行用户对源目录有读取权限,对目标目录有写入权限。
- 错误处理:为
ExecuteGroovyScript、FetchFile、PutFile配置REL_FAILURE分支,用于处理目录不存在、文件读取失败等异常,可搭配Notify处理器发送告警。
内容的提问来源于stack exchange,提问作者Tvitmsvleli
相关产品推荐
相关产品推荐

