如何使用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读取数组下的字段会返回多值集合,无法直接做等值匹配,需要先将数组展开为单行记录再关联。
配置步骤
- 配置QueryRecord的Record Reader为
JsonTreeReader,Schema访问策略选择自动推断即可。 - 配置Record Writer为
JsonRecordSetWriter,同样配置自动推断Schema输出JSON。 - 写入如下查询语句:
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
相关产品推荐
相关产品推荐

