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

如何在NiFi中按文件名逐个为每个FlowFile执行Groovy脚本文件

在NiFi中为每个FlowFile逐个执行对应Groovy脚本的实现方案

核心思路

通过FlowFile的文件名参数匹配并加载对应的Groovy脚本,利用NiFi的脚本执行类处理器,为每个FlowFile单独执行匹配到的脚本逻辑。

具体实现步骤

1. 规范Groovy脚本存储

  • 将所有Groovy脚本放在NiFi集群可访问的路径(本地目录或HDFS均可),确保脚本文件名和FlowFile的文件名一一对应(比如FlowFile名为data_001.json,对应脚本为data_001.groovy)。

2. 加载对应Groovy脚本到FlowFile

使用FetchFile处理器完成脚本加载:

  • 配置File Name属性为${filename}.groovy(${filename}会自动替换为当前FlowFile的文件名)。
  • 设置Destination为flowfile-content,将脚本内容写入FlowFile的内容区;或选择attribute,把脚本内容存入自定义属性(比如groovy_script_content)。

3. 逐个执行Groovy脚本

推荐用ExecuteScript处理器,配置如下:

  • 选择Script Language为Groovy。
  • 在Script Body中编写动态执行逻辑:
    // 读取脚本内容(存在属性里就取属性,存在内容里就读流)
    def scriptContent = flowFile.getAttribute('groovy_script_content') ?: new String(flowFile.getInputStream().readAllBytes())
    
    // 编译并执行脚本
    def shell = new GroovyShell()
    def script = shell.parse(scriptContent)
    
    // 把FlowFile和NiFi会话传递给脚本,方便脚本处理数据
    script.setProperty('flowFile', flowFile)
    script.setProperty('session', session)
    script.run()
    
    // 返回处理后的FlowFile到成功分支
    return session.transfer(flowFile, REL_SUCCESS)
    
  • 注意:如果脚本需要修改FlowFile的内容或属性,必须通过session对象操作,避免上下文冲突。

4. 容错处理

  • 添加RouteOnAttribute处理器,通过检查groovy_script_content属性是否存在、或捕获ExecuteScript的REL_FAILURE关系,将异常FlowFile路由到失败队列,用于重试或告警。

替代方案:InvokeScriptedProcessor

如果需要更模块化的管理,可使用InvokeScriptedProcessor:

  • 编写实现Processor接口的Groovy类,在类中根据FlowFile文件名动态加载外部脚本并执行。
  • 配置处理器时指定脚本类路径,同样通过${filename}参数匹配对应脚本文件。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 09:52:26