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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 03:27:21