在Apache NiFi中合并两个JSON的技术实现问询
在Apache NiFi中合并两个JSON的技术实现问询
嗨,我目前遇到了一个在Apache NiFi中合并两个JSON文档的需求,以下是我的输入数据和期望的输出结果,想请教具体的实现方法:
输入JSON 1
{ "property_name" : "TEST REAL ", "x_read" : "f6d139", "y_read" : "348be3" }
输入JSON 2
{ "property_name": "TEST REAL", "property_link": "https://sffff/sadsadd/", "location": [ { "location_name": "01 green", "location_type_name": "Green", "location_conditions": { "timestamp": "2024-04-24 07:32:07", "num_samples": 10, "img_extent": { "lng": -5.32, "lnd": -15.5 }, "avg": 31.76, "avg_color": "#4EC0003DA", "rating": "High" } } ] }
期望输出JSON
{ "property_name": "TEST REAL", "property_link": "https://sffff/sadsadd/", "x_read": "f6d139", "y_read": "348be3", "location": [ { "location_name": "01 green", "location_type_name": "Green", "location_conditions": { "timestamp": "2024-04-24 07:32:07", "num_samples": 10, "img_extent": { "lng": -5.32, "lnd": -15.5 }, "avg": 31.76, "avg_color": "#4EC0003DA", "rating": "High" } } ] }
(注:补全了期望输出中缺失的部分,核心是将输入1的x_read、y_read字段合并到输入2中,同时统一property_name为无空格版本)
可行实现方案
方案一:使用JoltTransformJSON处理器(推荐)
Jolt是NiFi中处理JSON结构转换的工具,通过编写Jolt规范就能实现合并,步骤如下:
- 预处理输入1的字段:添加
JoltTransformJSON处理器,用ModifyOverwriteBeta规范去除property_name的末尾空格,Jolt规则如下:[ { "operation": "modify-overwrite-beta", "spec": { "property_name": "=trim" } } ] - 分组合并JSON:用
MergeContent处理器,设置Merge Strategy为Bin-Packing Algorithm,Correlation Attribute Name为property_name,Delimiter用[和]包裹,将两个JSON合并成一个JSON数组。 - 合并为单个对象:再添加一个
JoltTransformJSON处理器,用Shift规范将数组合并为目标结构,Jolt规则如下:
这个规则会保留数组第一个元素(输入2)的所有字段,同时把第二个元素(预处理后的输入1)的[ { "operation": "shift", "spec": { "0": { "*": "&" }, "1": { "x_read": "&", "y_read": "&" } } } ]x_read、y_read字段添加进去。
方案二:使用ExecuteScript处理器(灵活定制)
如果需要更灵活的逻辑,可以用ExecuteScript编写Groovy脚本实现合并,示例脚本如下:
import groovy.json.JsonSlurper import groovy.json.JsonBuilder def flowFile = session.get() if (!flowFile) return // 假设flowFile内容是包含两个JSON对象的数组 def jsonSlurper = new JsonSlurper() def jsonArray = jsonSlurper.parseText(flowFile.read().getText("UTF-8")) def mergedJson = [:] // 合并输入2的所有字段 mergedJson.putAll(jsonArray[0]) // 添加输入1的指定字段 mergedJson.put("x_read", jsonArray[1].x_read) mergedJson.put("y_read", jsonArray[1].y_read) // 统一property_name格式 mergedJson.property_name = mergedJson.property_name.trim() def jsonBuilder = new JsonBuilder(mergedJson) flowFile.write(jsonBuilder.toPrettyString()) session.transfer(flowFile, REL_SUCCESS)
使用这个方案需要先确保两个JSON已通过MergeContent合并为数组,或者通过属性传递其中一个JSON的内容。
关键注意事项
- 字段冲突处理:两个输入都有
property_name字段,要明确以哪个版本为准(这里以输入2的无空格版本为准)。 - 关联逻辑:如果有多个文档需要合并,必须设置唯一的关联键(比如这里的
property_name)来分组,避免合并错误的文档。 - 性能优先:处理大量数据时,Jolt的性能优于
ExecuteScript,优先推荐Jolt方案。
备注:内容来源于stack exchange,提问作者BHADRAKA HERATH
相关产品推荐
相关产品推荐

