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

NiFi中使用Python实现极简ScriptedRecordSetWriter示例求助

NiFi ScriptedRecordSetWriter Python 极简实现示例

处理器基础配置步骤

  • 打开ScriptedRecordSetWriter处理器配置界面:
    • 将Record Reader设置为你已配置好的JsonTreeReader
    • Script Engine选择Python
    • 点击Script属性的编辑按钮,粘贴下方代码

Python 脚本代码

from org.apache.nifi.serialization import RecordSetWriterFactory, RecordWriter
from java.io import OutputStreamWriter, BufferedWriter
from java.nio.charset import StandardCharsets

# 自定义Writer工厂类,必须继承RecordSetWriterFactory接口
class HelloWorldWriterFactory(RecordSetWriterFactory):
    def createWriter(self, context, schema, outputStream):
        return HelloWorldWriter(outputStream, schema)

# 自定义Writer类,负责实际输出逻辑
class HelloWorldWriter(RecordWriter):
    def __init__(self, outputStream, schema):
        self.writer = BufferedWriter(OutputStreamWriter(outputStream, StandardCharsets.UTF_8))
        self.schema = schema

    def write(self, record):
        # 生成固定宽度的"Hello World",此处设为11字符,不足补空格
        fixed_output = "Hello World".ljust(11)
        self.writer.write(fixed_output)
        self.writer.newLine()
        return True

    def close(self):
        self.writer.flush()
        self.writer.close()

    def flush(self):
        self.writer.flush()

# 核心:必须返回RecordSetWriterFactory实例,解决"未定义RecordSetWriterFactory"报错的关键
return HelloWorldWriterFactory()

关键问题说明

  • 解决“未定义RecordSetWriterFactory”报错:NiFi要求ScriptedRecordSetWriter的脚本必须返回RecordSetWriterFactory类型的实例,因此必须自定义继承该接口的类,并在脚本末尾返回其实例,不能直接编写输出逻辑。
  • 固定宽度输出调整:示例中用ljust(11)保证输出字符串长度固定,若需要不同宽度,修改括号内的数值即可,多余位置会自动用空格填充。
  • 依赖类说明:脚本中导入的Java类均为NiFi运行环境自带,无需额外安装依赖包,直接调用即可。

内容的提问来源于stack exchange,提问作者Walker White

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 07:22:28