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
相关产品推荐
相关产品推荐

