NiFi中如何将FlowFile属性传入ExecuteScript处理器的JS函数?
NiFi ExecuteScript(ECMAScript)中正确处理FlowFile属性的方法
核心错误原因
- 属性引用方式错误:ExecuteScript的JS环境不支持直接用属性名(如
foo(A))或NiFi表达式语言(如foo(${A}))调用函数,必须通过flowFile.getAttribute()方法显式读取属性值。 - FlowFile未流转:所有从
session.get()获取或通过session.create()创建的FlowFile,必须调用session.transfer()流转到指定关系(如REL_SUCCESS/REL_FAILURE),或调用session.remove()移除,否则会触发会话错误。
正确代码示例
1. 读取属性、调用自定义函数并修改属性
// 自定义业务处理函数 function processAttr(value) { // 示例逻辑:给属性值添加处理前缀 return "handled_" + value; } var flowFile = session.get(); if (flowFile) { try { // 读取EvaluateJSONPath设置的目标属性 var targetAttr = flowFile.getAttribute("A"); // 替换为你的实际属性名 // 调用自定义函数,传入读取到的属性值 var processedResult = processAttr(targetAttr); // 更新FlowFile属性(putAttribute返回新对象,建议重新赋值) flowFile = session.putAttribute(flowFile, "ProcessedResult", processedResult); // 必须将FlowFile流转到指定关系 session.transfer(flowFile, REL_SUCCESS); } catch (err) { // 异常处理:流转到失败关系并记录日志 session.transfer(flowFile, REL_FAILURE); log.error("属性处理失败: " + err.message); } }
2. 创建新FlowFile并处理属性
var flowFile = session.get(); if (flowFile) { try { var sourceAttr = flowFile.getAttribute("someAttribute"); // 创建新的FlowFile实例 var newFlowFile = session.create(); // 为新FlowFile设置属性 newFlowFile = session.putAttribute(newFlowFile, "NewAttribute", sourceAttr); // 流转新FlowFile到成功关系 session.transfer(newFlowFile, REL_SUCCESS); // 处理原FlowFile:可流转到自定义关系或直接移除 session.transfer(flowFile, REL_ORIGINAL); // 若不需要原FlowFile,可使用:session.remove(flowFile); } catch (err) { session.transfer(flowFile, REL_FAILURE); log.error("创建新FlowFile失败: " + err.message); } }
关键注意事项
- 所有FlowFile操作必须在
session上下文内完成,且每个FlowFile必须有明确的流转或移除动作,否则会触发会话异常。 - 修改属性时,
session.putAttribute()会返回更新后的FlowFile对象,重新赋值可避免潜在的状态不一致问题。 - 自定义函数的参数直接传入
getAttribute()获取的字符串值即可,无需额外语法转换。
内容的提问来源于stack exchange,提问作者edjm
相关产品推荐
相关产品推荐

