如何在Apache NiFi中基于ID合并两个处理器的JSON数据?
下面是几种实用的实现方案,可根据你的数据规模和结构选择:
方案一:使用 MergeRecord(推荐,适合结构化JSON)
该方案基于NiFi的Record API,处理结构化JSON更高效规范:
步骤1:提取共同ID到FlowFile属性
给两个数据流分别添加ExtractText或JoltTransformJSON处理器,把JSON里的共同ID字段(比如id)提取为FlowFile属性,例如命名为target.id。
示例ExtractText配置:- 正则表达式:
"id":\s*"([^"]+)"(适配字符串ID,数字ID可调整为"id":\s*(\d+)) - 属性名称:
target.id
- 正则表达式:
步骤2:合并数据流到同一队列
用Funnel处理器将Jolt Transform和输入端口的输出流合并到同一个队列,确保两个流都能进入后续的MergeRecord处理器。步骤3:配置MergeRecord处理器
- Record Reader:选
JsonTreeReader(支持任意JSON结构)或JsonReader(需指定Schema),根据你的JSON格式调整。 - Record Writer:选
JsonRecordSetWriter,可配置输出为单个合并后的JSON对象或数组。 - 合并策略:设置
Merge Strategy为Merge by Attribute,Attribute Name填之前提取的target.id,这样相同ID的FlowFile会被分到同一组。 - 记录限制:设置
Max Records per Merge为2,保证每组只合并对应ID的两条数据。
- Record Reader:选
步骤4:合并字段(可选)
如果MergeRecord输出的是包含两条记录的JSON数组,要合并为单个对象,可添加JoltTransformJSON处理器,使用如下Jolt规范:[ { "operation": "shift", "spec": { "*": { "*": "&" } } } ]该规则会把数组中所有对象的字段合并到一个对象里,若存在相同字段,后出现的记录字段会覆盖前者。
方案二:使用 MergeContent + Jolt(通用文本合并)
如果不想依赖Record API,可采用文本合并后再做结构化处理:
步骤1:提取共同ID属性
同方案一,把共同ID提取为FlowFile属性target.id。步骤2:配置MergeContent处理器
- 合并策略:选
Bin by Attribute,Attribute Name为target.id。 - 分隔符设置:
Delimiter Strategy选Text,分隔符填",",确保两个JSON对象用逗号分隔。 - 分组限制:设置
Max Bin Items为2,保证每个分组只有对应ID的两条数据。
合并后输出内容示例:{"id":1,"fieldA":"valA"},{"id":1,"fieldB":"valB"}
- 合并策略:选
步骤3:转换为合并对象
添加ReplaceText处理器,在合并内容前后加上[和],转成JSON数组;再用JoltTransformJSON处理器,用方案一中的Jolt规则合并为单个对象。
方案三:使用分布式缓存 + LookupRecord(适合大数据量场景)
如果其中一个数据流的数据量较小,可将其存入缓存,再从另一个流中查询合并:
步骤1:加载小数据集到缓存
把输入端口的数据流接入PutDistributedMapCache处理器,配置:- Key Attribute:设为提取的ID属性
target.id - Value Serializer:选
JSON Serializer,把整个JSON内容作为缓存值
- Key Attribute:设为提取的ID属性
步骤2:查询缓存并合并
把Jolt Transform的数据流接入LookupRecord处理器,配置:- Record Reader:
JsonTreeReader - Lookup Service:选
DistributedMapCacheLookup服务,关联之前的缓存 - Lookup Key:设为JSON中的ID字段(比如
/id) - Record Writer:
JsonRecordSetWriter,配置将原记录和查询到的缓存记录合并为单个对象
- Record Reader:
内容的提问来源于stack exchange,提问作者Sathish

