NiFi中通过Jolt实现数组匹配字段插入的技术问询
解决NiFi中MergeRecord后关联数据合并的Jolt方案及备选实现
可行的Jolt Transform DSL规范
以下Jolt规范可实现按erp_projectid关联,将含project_id的条目值合并到同项目下所有含subject的主条目:
[ // 按erp_projectid分组,区分主记录与项目ID记录 { "operation": "shift", "spec": { "*": { "erp_projectid": { "*": { "@(2)": { "subject": { "@(3)": "groups.&1.main[]" }, "project_id": { "@(3)": "groups.&1.project" } } } } } } }, // 将project_id合并到每个主记录 { "operation": "shift", "spec": { "groups": { "*": { "main": { "*": { "*": "&4_&3_&2.&", "@(2,project.project_id)": "&4_&3_&2.project_id" } } } } } }, // 展开为顶层数组 { "operation": "shift", "spec": { "*": "[]" } } ]
逻辑说明
- 分组阶段:将输入数组按
erp_projectid拆分,把含subject的主记录归入main数组,含project_id的记录归入project节点。 - 合并阶段:将同项目下的
project_id值注入到每个主记录中。 - 展开阶段:把分组后的合并结果转为顶层数组,匹配期望输出格式。
备选实现方案(ExecuteScript处理器)
若Jolt无法满足更复杂的关联需求,可使用NiFi的ExecuteScript处理器,用Groovy代码实现灵活的数据合并:
import groovy.json.JsonSlurper import groovy.json.JsonBuilder def flowFile = session.get() if (!flowFile) return // 读取并解析输入JSON def slurper = new JsonSlurper() def inputData = slurper.parseText(flowFile.read().text) // 构建项目ID映射表:key为erp_projectid,value为对应的project_id def projectIdMap = inputData.findAll { it.containsKey('project_id') } .collectEntries { [(it.erp_projectid): it.project_id] } // 过滤主记录并添加project_id def outputData = inputData.findAll { it.containsKey('subject') } .collect { entry -> entry + [project_id: projectIdMap[entry.erp_projectid]] } // 写入输出JSON flowFile.write(new ByteArrayInputStream(new JsonBuilder(outputData).toByteArray())) session.transfer(flowFile, REL_SUCCESS)
优势
- 支持多字段关联、自定义匹配逻辑等复杂场景。
- 代码可读性高,便于调试和扩展。
内容的提问来源于stack exchange,提问作者Sanidhya Sharma
相关产品推荐
相关产品推荐

