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

如何使用Nifi QueryRecord实现FlowFile内JSON字段关联查询

问题场景

FlowFile中存储了经转换的JSON数据,需要关联JSON内不同节点的字段完成数据整合:按ID匹配从Original.ScoreInfo.Team节点提取R、PA字段值,同时基于Team.ID匹配获取Enriched.Team.Source.HV、Enriched.Team.SourceMapped.HV的对应值,最终期望输出的JSON结构如下:

[{
     "TeamID": 1000,
     "ValueR": 3, 
     "ValuePA": 13},
{
     "TeamID": 2000,
     "ValueR": 1, 
     "ValuePA": 14}
]

示例FlowFile数据

[
  {
    "Mapping": {
      "ScoreInfo": {
        "Team": [
          {
            "ID": 1,
            "Source": {
              "HV": 1,
              "ID": 1
            }
          },
          {
            "ID": 3,
            "Source": {
              "HV": 2,
              "ID": 3
            }
          }
        ]
      }
    },
    "Original": {
      "ScoreInfo": {
        "Team": [
          {
            "HV": 1,
            "ID": 1,
            "R": 1,
            "PA": 8,
            "H": 2,
            "BB": 2,
            "SB": 0,
            "E": 0,
            "Score": [
              {
                "Inn": 1,
                "TB": 2,
                "R": 0,
                "H": 0,
                "BB": 1
              },
              {
                "Inn": 2,
                "TB": 2,
                "R": 1,
                "H": 2,
                "BB": 1
              }
            ]
          },
          {
            "HV": 2,
            "ID": 3,
            "R": 1,
            "PA": 10,
            "H": 3,
            "BB": 0,
            "SB": 0,
            "E": 0,
            "Score": [
              {
                "Inn": 1,
                "TB": 1,
                "R": 1,
                "H": 2,
                "BB": 0
              },
              {
                "Inn": 2,
                "TB": 1,
                "R": 0,
                "H": 0,
                "BB": 0
              }
            ]
          }
        ]
      }
    },
    "MatchValue": "99999999",
    "Enriched": {
      "Team": [
        {
          "ID": 1,
          "Source": {
            "HV": 1,
            "ID": 1
          },
          "SourceMapped": {
            "HV": 1000,
            "ID": 1000
          }
        },
        {
          "ID": 2,
          "Source": {
            "HV": 2,
            "ID": 2
          },
          "SourceMapped": {
            "HV": 2000,
            "ID": 2000
          }
        }
      ],
      "MatchValue": "99999999"
    }
  }
]

问题说明

尝试使用QueryRecord结合RPATH实现需求,执行SELECT * FROM FLOWFILE WHERE RPATH(Original, '/ScoreInfo/Team/ID') = RPATH(Enriched, '/Team/ID')时返回空结果,需要给出正确的QueryRecord实现方案。


解决方案

问题原因

原有写法返回空的核心原因是:Original.ScoreInfo.Team和Enriched.Team均为数组结构,直接用RPATH读取数组下的字段会返回多值集合,无法直接做等值匹配,需要先将数组展开为单行记录再关联。

配置步骤

  1. 配置QueryRecord的Record Reader为JsonTreeReader,Schema访问策略选择自动推断即可。
  2. 配置Record Writer为JsonRecordSetWriter,同样配置自动推断Schema输出JSON。
  3. 写入如下查询语句:
SELECT
  enrich_team.SourceMapped.HV AS TeamID,
  SUM(orig_team.R) AS ValueR,
  SUM(orig_team.PA) AS ValuePA
FROM 
  flowfile,
  UNNEST(Original.ScoreInfo.Team) AS t(orig_team),
  UNNEST(Enriched.Team) AS t(enrich_team)
WHERE
  orig_team.HV = enrich_team.Source.HV
GROUP BY
  enrich_team.SourceMapped.HV

如果需求是按Team的ID字段关联,只需将WHERE条件替换为orig_team.ID = enrich_team.ID即可,执行后输出结果与期望结构完全一致。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 23:57:04