NiFi中是否存在AttributesToJSON处理器的反向功能?
实现NiFi中JSON转FlowFile属性(AttributesToJSON反向功能)
对于NiFi 2.0,你可以通过ExecuteScript处理器配合脚本快速实现JSON到属性的转换,以下是具体操作:
一、Groovy脚本实现(推荐)
- 拖拽
ExecuteScript处理器到画布,双击进入配置页面。 - 在
Script Language下拉菜单选择Groovy。 - 替换默认脚本为以下代码:
import groovy.json.JsonSlurper import org.apache.nifi.processor.io.StreamCallback import java.nio.charset.StandardCharsets def flowFile = session.get() if (!flowFile) return flowFile = session.write(flowFile, { inputStream, outputStream -> // 读取并解析JSON内容 def jsonContent = inputStream.text def json = new JsonSlurper().parseText(jsonContent) // 遍历JSON顶层键值对,设置为FlowFile属性 json.each { key, value -> // 跳过空值,避免无效属性;属性名需符合NiFi规则(无特殊字符/空格) if (value != null) { flowFile = session.putAttribute(flowFile, key.toString(), value.toString()) } } // 将原JSON内容写回(如果不需要保留内容,可省略此步骤) outputStream.write(jsonContent.getBytes(StandardCharsets.UTF_8)) } as StreamCallback) session.transfer(flowFile, REL_SUCCESS)
- 配置完成后启动处理器即可。
脚本说明
- 该脚本会解析FlowFile中的JSON内容,将顶层键值对转为FlowFile属性;
- 如果需要处理嵌套JSON(比如把
user.name转为属性名),可以扩展脚本添加扁平化逻辑; - 若不需要保留原JSON内容,可删除
outputStream.write(...)这一行,减少IO开销。
二、Clojure脚本实现
如果你偏好Clojure,也可以用以下脚本:
(require '[clojure.data.json :as json]) (require '[clojure.java.io :as io]) (def flow-file (.get session)) (when flow-file (let [updated-flow-file (.write session flow-file (proxy [org.apache.nifi.processor.io.StreamCallback] [] (process [in out] (let [json (json/read (io/reader in))] (doseq [[k v] json] (when v (.putAttribute session flow-file (str k) (str v)))) ;; 可选:写回原内容 (io/copy in out)))))] (.transfer session updated-flow-file (.getRelationshipByName processor "success"))))
注意事项
- NiFi属性名不能包含空格、
:、/等特殊字符,若JSON键名不符合规则,需在脚本中添加替换逻辑(比如把空格换成下划线); - 若JSON内容过大,脚本可能存在性能瓶颈,此时建议开发自定义处理器;
- 测试时可以先用
GenerateFlowFile生成包含JSON内容的测试数据,再连接ExecuteScript验证效果。
内容的提问来源于stack exchange,提问作者Trenin
相关产品推荐
相关产品推荐

