如何在NiFi中合并多ExecuteScript处理器属性并汇总状态码
解决方案
针对你遇到的NiFi多ExecuteScript处理器属性合并并求和的问题,提供两种高效可行的方案:
方案一:MergeAttribute合并属性 + UpdateAttribute计算总和
适合四个ExecuteScript逻辑无法整合、必须独立执行的场景:
- 将四个ExecuteScript处理器的输出连接到MergeAttribute处理器
- 配置
Correlation Attribute Name为原始流文件的唯一标识(比如uuid),确保同一份原始文件的四个分支流文件被归为一组 - 设置
Merge Strategy为Merge Attributes,Attribute Strategy为Keep All Attributes,合并后的流文件会包含所有四个returnCodeXX属性
- 配置
- 接着连接UpdateAttribute处理器,用NiFi表达式语言计算总和:
添加属性finalReturnCode,值配置为:
(注:${returnCodeAB:toNumber():plus(${returnCodeAC:toNumber()}):plus(${returnCodeAL:toNumber()}):plus(${returnCodeAM:toNumber()})}toNumber()用于确保属性值转为数字类型,避免字符串拼接错误)
方案二:整合逻辑到单个ExecuteScript处理器
如果四个配置集的执行逻辑可以整合,这是最高效的方案,避免流文件拆分与合并的开销:
直接在一个ExecuteScript中完成所有四个配置集的执行、属性添加及求和,示例Python脚本如下:
from org.apache.nifi.processor.io import StreamCallback import java.io class Callback(StreamCallback): def process(self, inputStream, outputStream): # 执行AB配置集逻辑,替换为实际业务代码获取状态码 return_code_ab = 0 # 执行AC配置集逻辑 return_code_ac = 0 # 执行AL配置集逻辑 return_code_al = 0 # 执行AM配置集逻辑 return_code_am = 0 # 计算最终状态码 final_return_code = return_code_ab + return_code_ac + return_code_al + return_code_am # 为流文件添加属性 global flowFile flowFile = session.putAttribute(flowFile, "returnCodeAB", str(return_code_ab)) flowFile = session.putAttribute(flowFile, "returnCodeAC", str(return_code_ac)) flowFile = session.putAttribute(flowFile, "returnCodeAL", str(return_code_al)) flowFile = session.putAttribute(flowFile, "returnCodeAM", str(return_code_am)) flowFile = session.putAttribute(flowFile, "finalReturnCode", str(final_return_code)) flowFile = session.get() if flowFile is not None: flowFile = session.write(flowFile, Callback()) session.transfer(flowFile, REL_SUCCESS)
该方案直接在一个处理器内完成所有操作,无需分支,性能最优。
内容的提问来源于stack exchange,提问作者Shreya B
相关产品推荐
相关产品推荐

