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

基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 04:00:11