如何在NiFi ExecuteGroovy处理器中实现FlowFile键值对移位
在NiFi的ExecuteGroovy处理器中实现FlowFile数据移位逻辑
需求说明
场景1:当FlowFile中存在值为
"M"的项时输入示例:
{ "cv1": "A", "cv2": "B", "cv3": "M", "cv4": "D", "cv5": "C" }输出要求:移除值为
"M"的项,后续值依次向上移位,最后一个位置设为空字符串{ "cv1": "A", "cv2": "B", "cv3": "D", "cv4": "C", "cv5": "" }场景2:当FlowFile中不存在值为
"M"的项时输入示例:
{ "cv1": "A", "cv2": "B", "cv3": "C", "cv4": "D", "cv5": "E" }输出要求:所有值整体向上移位一位,最后一个位置设为空字符串
{ "cv1": "B", "cv2": "C", "cv3": "D", "cv4": "E", "cv5": "" }
现有代码片段
import groovy.json.JsonOutput import groovy.json.JsonSlurper import org.apache.nifi.processor.io.StreamCallback import java.nio.charset.StandardCharsets def flowFile = session.get() if(!flowFile) return try{ def inputStream1 = session.read(flowFile) def json = new groovy.json.JsonSlurper().parse(inputStream1)
完整实现代码
import groovy.json.JsonOutput import groovy.json.JsonSlurper import org.apache.nifi.processor.io.StreamCallback import java.nio.charset.StandardCharsets def flowFile = session.get() if (!flowFile) return try { flowFile = session.write(flowFile, { inputStream, outputStream -> // 解析输入的JSON数据 def json = new JsonSlurper().parse(inputStream) // 按cv1到cv5的顺序提取所有值,缺失项默认空字符串 def values = (1..5).collect { idx -> json["cv$idx"] ?: "" } def hasM = values.contains("M") def newValues = [] if (hasM) { // 找到第一个"M"的位置,移除后后续元素前移 def mIndex = values.indexOf("M") newValues = values[0..<mIndex] + values[mIndex+1..-1] // 补空字符串保证总长度为5 if (newValues.size() < 5) { newValues.add("") } } else { // 整体上移一位,最后补空字符串 newValues = values[1..-1] + [""] } // 将新值映射回cv1到cv5的结构 def result = (1..5).collectEntries { idx -> ["cv$idx", newValues[idx-1]] } // 输出JSON结果 outputStream.write(JsonOutput.toJson(result).getBytes(StandardCharsets.UTF_8)) } as StreamCallback) session.transfer(flowFile, REL_SUCCESS) } catch (Exception e) { log.error("处理FlowFile失败: ${e.getMessage()}", e) session.transfer(flowFile, REL_FAILURE) }
代码说明
- 数据解析与提取:通过
JsonSlurper读取输入流,按cv1到cv5的固定顺序提取值,确保处理逻辑的顺序一致性。 - 核心逻辑处理:
- 检测到
"M"时,移除第一个匹配项,后续值自动前移,最后补空保证结构长度不变。 - 未检测到
"M"时,直接截取从第二个值开始的所有元素,末尾补空实现整体上移。
- 检测到
- 结果构建与输出:将处理后的新值重新映射为
cv1到cv5的JSON结构,写入输出流。 - 异常处理:捕获处理过程中的异常,将FlowFile转发到失败关系并记录错误日志。
内容的提问来源于stack exchange,提问作者Yash Anand
相关产品推荐
相关产品推荐

