NiFi运行PyStreamCallback脚本报无length属性错误排查
NiFi ExecuteScript处理器运行Python脚本统计JSON数组长度报错修复
报错根因
抛出AttributeError: PyStreamCallback object has no attribute length是代码存在3处明确问题导致:
- 实例复用错误:
session.write()执行时创建了第一个PyStreamCallback实例,该实例执行process()方法后会生成length属性,但后续调用get_length_of_array()时是重新实例化的第二个全新PyStreamCallback对象,这个新对象从未执行process()方法,不存在length属性,调用直接报错。 - 流处理逻辑缺失:重写
StreamCallback.process()方法时,只读取了输入流内容,没有将内容写回outputStream,会导致处理后的FlowFile内容被清空。 - 变量名大小写错误:最后
session.transfer传入的是大写开头的flowFile,和前面赋值的小写flowfile变量名不匹配,会触发变量不存在的运行时错误。
另外原脚本中两层遍历取第一个key的逻辑是硬编码,仅适配每层只有一个key的两层嵌套JSON结构,碰到结构不符的内容会触发类型错误。
修正后可运行脚本
from org.apache.commons.io import IOUtils from java.nio.charset import StandardCharsets from org.apache.nifi.processor.io import StreamCallback import json class PyStreamCallback(StreamCallback): def __init__(self): self.length = "0" # 初始化属性默认值,从根源避免属性不存在报错 def process(self, inputStream, outputStream): # 读取FlowFile原始内容 raw_content = IOUtils.toString(inputStream, StandardCharsets.UTF_8) json_data = json.loads(raw_content) # 定位目标数组(保留原脚本两层取首个key的逻辑,可根据实际JSON结构调整为明确字段路径) target_node = json_data for _ in range(2): if isinstance(target_node, dict) and len(target_node.keys()) > 0: first_key = list(target_node.keys())[0] target_node = target_node[first_key] else: break # 统计数组长度 self.length = str(len(target_node)) if isinstance(target_node, list) else "0" # 必须将原始内容写回输出流,避免FlowFile内容丢失 outputStream.write(raw_content.encode(StandardCharsets.UTF_8)) def get_length_of_array(self): return self.length # 主流程逻辑 flowfile = session.get() if flowfile is not None: # 仅实例化一次callback,全程复用同一个实例 stream_callback = PyStreamCallback() flowfile = session.write(flowfile, stream_callback) # 从已执行process方法的实例中读取长度值写入属性 flowfile = session.putAttribute(flowfile, "length", stream_callback.get_length_of_array()) # 保持变量名大小写一致,传递FlowFile到成功关系 session.transfer(flowfile, REL_SUCCESS)
提示:如果你的JSON结构固定,建议不要用循环取第一个key的方式定位数组,直接写明确的字段路径即可,比如数组在
res.data字段下,直接写target_node = json_data["res"]["data"],容错性更高。
内容的提问来源于stack exchange,提问作者Vlad Vlad
相关产品推荐
相关产品推荐

