如何在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处理时快速查询:
GetFile:读取CSV文件,配置输入目录指向CSV所在路径。ConvertCSVToRecord:将CSV转换为Record格式:- 设置
Delimiter为\s+(匹配多个空格分隔符) - 指定Record Schema:
{ "type": "record", "name": "KeyValue", "fields": [ {"name": "Key", "type": "string"}, {"name": "Value", "type": "string"} ] }
- 设置
PutDistributedMapCache:将每条Record存入缓存:Cache Key设为${record:value('/Key')}Cache Value设为${record:value('/Value')}- 提前确保DistributedMapCacheServer服务已启动。
二、处理JSON并替换字段
对JSON数组逐个元素处理,替换name字段:
GetFile:读取目标JSON文件,配置输入目录指向JSON所在路径(可通过Wait处理器确保缓存已加载完成后再执行此步骤)。SplitJson:将JSON数组拆分为单个对象:JsonPath Expression设为$[*],每个数组元素会生成独立的FlowFile。
EvaluateJsonPath:提取当前对象的key值作为FlowFile属性:Destination选择flowfile-attribute- 添加属性
current_key,值为$.key
FetchDistributedMapCache:根据current_key查询缓存中的Value:Cache Key设为${current_key}- 查询到的Value会存入
cached_value属性。
JoltTransformJSON:替换name字段值:- 设置Jolt Spec为:
[ { "operation": "modify-overwrite-beta", "spec": { "name": "${cached_value}" } } ]
- 设置Jolt Spec为:
MergeContent:将单个JSON对象合并回数组:Merge Strategy选择Bin-Packing AlgorithmDelimiter设为,\nHeader设为[,Footer设为]
PutFile:将处理后的JSON写入目标输出目录。
注意事项
- 确保DistributedMapCacheServer服务正常运行,缓存配置的过期时间足够覆盖JSON处理时长。
- 若需批量处理,可通过
Wait处理器监听缓存加载完成的信号,避免JSON处理时缓存未就绪。 - 若CSV或JSON格式有变动,需同步调整
ConvertCSVToRecord的Schema或JoltTransformJSON的Spec。
内容的提问来源于stack exchange,提问作者mycuser
相关产品推荐
相关产品推荐

