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

如何在Apache NiFi中用CSV替换JSON中的值

使用Apache NiFi实现CSV值替换JSON对应字段

需求说明

需要将CSV中的Key-Value映射,替换JSON数组中每个对象的name字段(原name与key字段值一致,替换为CSV中对应Key的Value)。

原始数据

原JSON数组

[
  {
    "key": "key_name1",
    "name": "key_name1"
  },
  {
    "key": "key_name2",
    "name": "key_name2"
  },
  {
    "key": "key_name3",
    "name": "key_name3"
  },
  {
    "key": "key_name4",
    "name": "key_name4"
  }
]

原CSV文件(空格分隔)

Key            Value
key_name1      key_value1
key_name2      key_value2
key_name3      key_value3
key_name4      key_value4

期望输出JSON

[
  {
    "key": "key_name1",
    "name": "key_value1"
  },
  {
    "key": "key_name2",
    "name": "key_value2"
  },
  {
    "key": "key_name3",
    "name": "key_value3"
  },
  {
    "key": "key_name4",
    "name": "key_value4"
  }
]

实现步骤

一、构建CSV键值映射缓存

先将CSV中的数据存入分布式缓存,方便后续JSON处理时快速查询:

  1. GetFile:读取CSV文件,配置输入目录指向CSV所在路径。
  2. ConvertCSVToRecord:将CSV转换为Record格式:
    • 设置Delimiter为\s+(匹配多个空格分隔符)
    • 指定Record Schema:
      {
        "type": "record",
        "name": "KeyValue",
        "fields": [
          {"name": "Key", "type": "string"},
          {"name": "Value", "type": "string"}
        ]
      }
      
  3. PutDistributedMapCache:将每条Record存入缓存:
    • Cache Key设为${record:value('/Key')}
    • Cache Value设为${record:value('/Value')}
    • 提前确保DistributedMapCacheServer服务已启动。

二、处理JSON并替换字段

对JSON数组逐个元素处理,替换name字段:

  1. GetFile:读取目标JSON文件,配置输入目录指向JSON所在路径(可通过Wait处理器确保缓存已加载完成后再执行此步骤)。
  2. SplitJson:将JSON数组拆分为单个对象:
    • JsonPath Expression设为$[*],每个数组元素会生成独立的FlowFile。
  3. EvaluateJsonPath:提取当前对象的key值作为FlowFile属性:
    • Destination选择flowfile-attribute
    • 添加属性current_key,值为$.key
  4. FetchDistributedMapCache:根据current_key查询缓存中的Value:
    • Cache Key设为${current_key}
    • 查询到的Value会存入cached_value属性。
  5. JoltTransformJSON:替换name字段值:
    • 设置Jolt Spec为:
      [
        {
          "operation": "modify-overwrite-beta",
          "spec": {
            "name": "${cached_value}"
          }
        }
      ]
      
  6. MergeContent:将单个JSON对象合并回数组:
    • Merge Strategy选择Bin-Packing Algorithm
    • Delimiter设为,\n
    • Header设为[,Footer设为]
  7. PutFile:将处理后的JSON写入目标输出目录。

注意事项

  • 确保DistributedMapCacheServer服务正常运行,缓存配置的过期时间足够覆盖JSON处理时长。
  • 若需批量处理,可通过Wait处理器监听缓存加载完成的信号,避免JSON处理时缓存未就绪。
  • 若CSV或JSON格式有变动,需同步调整ConvertCSVToRecord的Schema或JoltTransformJSON的Spec。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 17:10:11