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
相关产品推荐
相关产品推荐

