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

Apache NiFi中JOLT合并无共同键数组及替代方案问询

实现JSON数组笛卡尔积合并的JOLT规则及Apache NiFi替代方案

问题概述

需要将输入JSON中questionResponses节点下的responses数组与answers数组执行笛卡尔积合并,生成包含所有组合的输出数组。当前JOLT规则仅能提取两个数组,无法完成合并,需修正JOLT规则,同时提供Apache NiFi中JOLT方案无效时的替代处理器方法。


输入输出示例

输入JSON

{
  "questionResponses": {
    "responses": ["A", "B"],
    "answers": ["X", "Y", "Z"]
  }
}

目标输出JSON

[
  {"response": "A", "answer": "X"},
  {"response": "A", "answer": "Y"},
  {"response": "A", "answer": "Z"},
  {"response": "B", "answer": "X"},
  {"response": "B", "answer": "Y"},
  {"response": "B", "answer": "Z"}
]

当前JOLT规则与不理想输出

当前JOLT规则

[
  {
    "operation": "shift",
    "spec": {
      "questionResponses": {
        "responses": "responses",
        "answers": "answers"
      }
    }
  }
]

不理想输出

{
  "responses": ["A", "B"],
  "answers": ["X", "Y", "Z"]
}

仅提取了两个数组,未完成笛卡尔积合并。


正确JOLT转换规则

以下JOLT Spec通过两步shift操作实现笛卡尔积合并:

[
  {
    "operation": "shift",
    "spec": {
      "questionResponses": {
        "responses": {
          "*": {
            "@": "&3.&2[#2].response",
            "@(2,answers)": "&3.&2[#2].answers[]"
          }
        }
      }
    }
  },
  {
    "operation": "shift",
    "spec": {
      "questionResponses": {
        "responses": {
          "*": {
            "answers": {
              "*": {
                "@(2,response)": "[#4].response",
                "@": "[#4].answer"
              }
            }
          }
        }
      }
    }
  }
]

规则说明

  1. 第一步shift:将每个response与完整的answers数组绑定,形成responses[N].response和responses[N].answers的结构。
  2. 第二步shift:遍历每个responses条目下的answers数组,将每个answer与对应的response拆分为独立的键值对,最终生成笛卡尔积结果数组。

Apache NiFi替代处理器方案

若JOLT规则无法满足需求(如复杂嵌套结构或特殊数据类型),可使用ExecuteScript处理器编写脚本实现笛卡尔积:

Groovy脚本示例

import groovy.json.JsonSlurper
import groovy.json.JsonBuilder

def flowFile = session.get()
if (!flowFile) return

try {
    // 读取输入JSON
    def inputJson = new JsonSlurper().parseText(flowFile.read().text)
    def responses = inputJson.questionResponses.responses
    def answers = inputJson.questionResponses.answers
    
    // 执行笛卡尔积合并
    def resultList = []
    responses.each { resp ->
        answers.each { ans ->
            resultList.add([response: resp, answer: ans])
        }
    }
    
    // 写入结果到FlowFile
    flowFile.write(new JsonBuilder(resultList).toString())
    session.transfer(flowFile, REL_SUCCESS)
} catch (Exception e) {
    log.error("Failed to process Cartesian product", e)
    session.transfer(flowFile, REL_FAILURE)
}

使用说明

  1. 在NiFi中添加ExecuteScript处理器,选择Groovy作为脚本语言。
  2. 将上述脚本粘贴到处理器的Script属性中。
  3. 配置处理器的输入输出关系,确保成功/失败分支正确路由。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 00:06:24