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

Apache NiFi 2.2.0 ExecuteGroovyScript聚合FlowFile遇MissingMethodException求助

Apache NiFi 2.2.0 ExecuteGroovyScript聚合FlowFile报错解决

错误原因分析

报错核心是MissingMethodException,说明session.read()方法的参数不匹配:

  • 传入了StreamCallback类型,但NiFi的GroovyProcessSessionWrap.read()方法需要的是InputStreamCallback类型
  • 代码存在变量名错误:定义了def data = [],却尝试调用codes.add(),会导致后续空指针异常

修正后的完整代码

import groovy.json.JsonSlurper
import groovy.json.JsonOutput
import org.apache.nifi.processor.io.InputStreamCallback
import org.apache.nifi.processor.io.OutputStreamCallback

// 获取最多1000个FlowFile
def flowFileList = session.get(1000)
if (flowFileList.isEmpty()) {
    return
}

def jsonSlurper = new JsonSlurper()
def data = []

flowFileList.each { flowFile ->
    // 使用InputStreamCallback替代StreamCallback,匹配方法签名
    session.read(flowFile, { inputStream ->
        def json = jsonSlurper.parse(inputStream)
        // 修正变量名:用data而非codes
        data.add(json.cis)
    } as InputStreamCallback)
    // 清理已处理的原始FlowFile,避免滞留队列
    session.remove(flowFile)
}

// 构造聚合后的JSON结构
def jsonOutput = JsonOutput.toJson([
    'data': data,
    'x': 'y',
    'z': 't',
    'v': 'f'
])

// 创建新FlowFile并写入聚合内容
def newFlowFile = session.create()
newFlowFile = session.write(newFlowFile, { outputStream ->
    outputStream.write(jsonOutput.bytes)
} as OutputStreamCallback)

// 设置MIME类型并转移至成功关系
newFlowFile.'mime.type' = 'application/json'
session.transfer(newFlowFile, REL_SUCCESS)

关键修正点说明

  • 回调类型替换:导入InputStreamCallback,并将回调声明改为as InputStreamCallback,匹配session.read()的方法要求
  • 变量名修正:将codes.add(json.cis)改为data.add(json.cis),对应之前定义的集合变量
  • 原始文件清理:添加session.remove(flowFile),处理完每个原始FlowFile后从会话中移除,避免文件滞留
  • 写入回调规范:写入新FlowFile时显式使用OutputStreamCallback,代码更符合NiFi API规范

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 15:42:15