NiFi使用UpdateRecord将JSON字符串转为合法JSON异常的解决咨询
NiFi UpdateRecord 嵌套JSON字符串转对象解决方案
你当前遇到的问题是没有对user字段的JSON字符串做解析转换,且Schema定义为字符串类型,导致处理器直接输出未解析的字符串格式。以下是3种可行解决方案:
方案1:调整UpdateRecord原有配置(无需更换处理器)
这是改动最小的方案,只需调整3处配置即可:
- 首先修改对应流数据的Schema定义:将
user字段的类型从STRING调整为嵌套Record类型,示例Schema如下:
{ "type": "record", "name": "Root", "fields": [ { "name": "user", "type": { "type": "record", "name": "User", "fields": [ {"name": "id", "type": "string"}, {"name": "name", "type": "string"} ] } } ] }
- 其次修改UpdateRecord的替换规则:
- Record Path填写
/user - 替换值填写
${field.value:fromJson()}
- Record Path填写
- 最后检查JsonRecordSetWriter配置:确保没有开启强制字符串转义的选项,使用默认配置即可。
你当前的配置截图如下:
方案2:更换为JoltTransformJSON处理器(无需预定义Schema)
如果不想维护复杂的Schema定义,直接替换UpdateRecord为JoltTransformJSON处理器,转换规则选择Chain,填写如下配置即可:
[ { "operation": "modify-overwrite-beta", "spec": { "user": "=fromJson(@(1,user))" } } ]
该方案会自动将user字段的JSON字符串解析为对象,一步得到你需要的输出,适配性更高。
方案3:脚本实现(仅适合有额外自定义逻辑的场景)
如果后续还有其他定制化处理逻辑,可以使用ExecuteScript处理器选择Groovy语言,参考代码如下:
import groovy.json.JsonSlurper import groovy.json.JsonOutput def flowFile = session.get() if (!flowFile) return try { flowFile = session.read(flowFile, { inputStream, outputStream -> def data = new JsonSlurper().parse(inputStream) data.user = new JsonSlurper().parseText(data.user as String) outputStream.write(JsonOutput.toJson(data).bytes) } as StreamCallback) session.transfer(flowFile, REL_SUCCESS) } catch (Exception e) { session.transfer(flowFile, REL_FAILURE) }
优先推荐前两种方案,无需编写自定义代码即可实现需求。
内容的提问来源于stack exchange,提问作者Repo Code
相关产品推荐
相关产品推荐


