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

如何在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)
}

代码说明

  1. 数据解析与提取:通过JsonSlurper读取输入流,按cv1到cv5的固定顺序提取值,确保处理逻辑的顺序一致性。
  2. 核心逻辑处理:
    • 检测到"M"时,移除第一个匹配项,后续值自动前移,最后补空保证结构长度不变。
    • 未检测到"M"时,直接截取从第二个值开始的所有元素,末尾补空实现整体上移。
  3. 结果构建与输出:将处理后的新值重新映射为cv1到cv5的JSON结构,写入输出流。
  4. 异常处理:捕获处理过程中的异常,将FlowFile转发到失败关系并记录错误日志。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 19:40:33