You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Apache NiFi中基于ID合并两个处理器的JSON数据?

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的两条数据。
  • 步骤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内容作为缓存值
  • 步骤2:查询缓存并合并
    把Jolt Transform的数据流接入LookupRecord处理器,配置:

    • Record Reader:JsonTreeReader
    • Lookup Service:选DistributedMapCacheLookup服务,关联之前的缓存
    • Lookup Key:设为JSON中的ID字段(比如/id)
    • Record Writer:JsonRecordSetWriter,配置将原记录和查询到的缓存记录合并为单个对象

内容的提问来源于stack exchange,提问作者Sathish

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.21 15:27:21