如何批量将多字段JSON转换为NiFi属性?
解决方案
一、批量将JSON字段转为NiFi属性
无需逐个配置EvaluateJsonPath,直接使用ConvertJSONtoAttributes组件即可实现批量转换:
- 该组件默认会把JSON顶层的键值对直接转为FlowFile属性(键为JSON字段名,值为对应字段值)
- 若JSON存在嵌套结构,可先通过
JoltTransformJSON将嵌套字段扁平化为顶层结构,示例Jolt规范:
[ { "operation": "shift", "spec": { "*": "&", "nestedObj": { "*": "nestedObj_&" } } } ]
- 如需过滤特定字段,可在
ConvertJSONtoAttributes的Include Attributes中指定字段前缀或正则,例如只保留大写字母开头的字段:^[A-Z]+$
二、解决ConvertJsontoSQL适配Phoenix的问题
ConvertJsontoSQL生成的SQL通常不符合Phoenix语法要求,可通过以下两种方式处理:
方式1:Jolt+ReplaceText生成合规SQL
- 使用
JoltTransformJSON将JSON转为SQL模板所需结构,示例Jolt规范:
[ { "operation": "shift", "spec": { "*": { "$": "columns[]", "@": "values[]" } } } ]
转换后输出结构:
{ "columns": ["AAAA", "BBBB", "CCCC", ...], "values": ["AAAA", "BBBB", "CCCC", ...] }
- 用
EvaluateJsonPath提取columns和values为FlowFile属性,配置路径分别为$.columns和$.values - 通过
ReplaceText组件生成Phoenix插入SQL,模板如下:
UPSERT INTO YOUR_TABLE_NAME ("${columns:replace(',', '","')}") VALUES ('${values:replace(',', "','")}')
注:Phoenix要求字段名加双引号(区分大小写),字符串值加单引号,且UPSERT是Phoenix标准的插入/更新语法
方式2:ExecuteScript自定义生成SQL
用Groovy脚本读取FlowFile内容和属性,直接生成符合Phoenix要求的SQL:
import org.apache.nifi.processor.io.StreamCallback import groovy.json.JsonSlurper def flowFile = session.get() if (!flowFile) return flowFile = session.write(flowFile, { inputStream, outputStream -> def json = new JsonSlurper().parse(inputStream) def columns = json.keySet().collect { "\"$it\"" }.join(',') def values = json.values().collect { "'${it.replace("'", "''")}'" }.join(',') def sql = "UPSERT INTO YOUR_TABLE_NAME ($columns) VALUES ($values)" outputStream.write(sql.getBytes('UTF-8')) } as StreamCallback) session.transfer(flowFile, REL_SUCCESS)
- 替换
YOUR_TABLE_NAME为实际Phoenix表名 - 脚本中已包含单引号转义逻辑,避免SQL语法错误
三、Phoenix数据写入
生成合规SQL后,使用PutSQL组件连接Phoenix:
- 提前将Phoenix JDBC驱动包放入NiFi的依赖库目录
- 配置JDBC URL为Phoenix连接地址(例如
jdbc:phoenix:zk-host:2181) - 直接执行生成的SQL即可完成数据写入
内容的提问来源于stack exchange,提问作者veganzombie
相关产品推荐
相关产品推荐

