请求提供NiFi ExecuteScript处理器的Python JSON转换脚本
NiFi ExecuteScript 处理器Python转换脚本
输入JSON示例
[ { "id": 3, "name": "ghi", "gender": null, "city": "" }, { "id": 5, "name": "", "gender": "", "city": "es" } ]
输出JSON示例
[ { "field": "gender", "remark": "testing", "id": 3 }, { "field": "city", "remark": "testing", "id": 3 }, { "field": "name", "remark": "testing", "id": 5 }, { "field": "gender", "remark": "testing", "id": 5 } ]
转换规则
- 遍历输入数组中的每个对象
- 提取对象中值为
null或空字符串""的键(跳过id字段本身) - 为每个符合条件的键生成新对象,包含三个固定结构字段:
field:对应提取的键名remark:固定值"testing"id:对应原对象的id值
Python脚本(用于NiFi ExecuteScript处理器)
import json from org.apache.commons.io import IOUtils from java.nio.charset import StandardCharsets from org.apache.nifi.processor.io import StreamCallback class JsonTransformCallback(StreamCallback): def process(self, inputStream, outputStream): # 读取输入流中的JSON数据 input_content = IOUtils.toString(inputStream, StandardCharsets.UTF_8) source_data = json.loads(input_content) result_list = [] for obj in source_data: current_id = obj.get("id") # 遍历对象的键值对,筛选目标字段 for key, val in obj.items(): if key == "id": continue if val is None or val == "": result_list.append({ "field": key, "remark": "testing", "id": current_id }) # 将转换结果写入输出流 outputStream.write(json.dumps(result_list, indent=2).encode('utf-8')) # 处理流文件并传递结果 flow_file = session.get() if flow_file is not None: flow_file = session.write(flow_file, JsonTransformCallback()) session.transfer(flow_file, REL_SUCCESS)
NiFi处理器配置说明
- 在ExecuteScript处理器的配置界面,将Script Language设置为
Python - 将上述脚本完整粘贴到Script Body区域
- 确保处理器的输入输出关系配置正确,处理完成的流文件会通过
REL_SUCCESS关系传递
内容的提问来源于stack exchange,提问作者just learner
相关产品推荐
相关产品推荐

