如何在NiFi中解析JSON的message字段并将结果整合回内容?
可行解决方案
方案1:纯处理器组合(无需代码,适合新手)
步骤分解:
保留原始JSON+提取syslog字段到属性
用EvaluateJSONPath处理器,做两项配置:- 规则1:属性名设为
syslog_raw,值为$.message(把JSON里的message字段提取为流文件属性) - 规则2:属性名设为
original_json,值为$(把整个JOLT处理后的JSON提取为属性,用于后续合并) - 将“Destination”设置为
flowfile-attribute
- 规则1:属性名设为
解析syslog属性(二选一)
- 用
ParseSyslog:
配置处理器,将“Source”设为FlowFile Attribute,指定“Attribute Name”为syslog_raw,解析后会自动生成syslog.timestamp、syslog.host、syslog.message等属性。 - 用
ExtractGrok:
配置处理器,“Source”选FlowFile Attribute,“Attribute Name”填syslog_raw,在“Grok Pattern”中填入对应syslog模式(比如%{SYSLOGTIMESTAMP:timestamp} %{SYSLOGHOST:host} %{DATA:process}%{DATA:pid}: %{GREEDYDATA:content}),解析后会生成你定义的字段属性(如timestamp、host)。
- 用
把解析后的属性转成JSON
用AttributesToJSON处理器:- “Attributes to Include”:如果用ParseSyslog填
syslog.*,如果用ExtractGrok填你定义的字段名(比如timestamp,host,content) - “Destination”选
flowfile-attribute,属性名设为parsed_syslog
- “Attributes to Include”:如果用ParseSyslog填
合并原始JSON和解析后的JSON
用JOLTTransformJSON处理器:- “JOLT Specification”使用合并规则:
[ { "operation": "shift", "spec": { "original_json": { "*": "&" }, "parsed_syslog": { "*": "&" } } } ] - “Input JSON Path”设为
$,处理器会读取流文件属性中的original_json和parsed_syslog,合并成新的JSON对象输出到流文件内容。
- “JOLT Specification”使用合并规则:
方案2:用ExecuteScript脚本灵活处理(适合自定义逻辑场景)
如果觉得处理器组合繁琐,可直接用ExecuteScript(选Groovy语言)完成解析+合并:
- 先通过
EvaluateJSONPath把message字段提取为属性syslog_raw。 - 配置
ExecuteScript的脚本示例(基于Syslog解析):import org.apache.nifi.processor.io.StreamCallback import groovy.json.JsonSlurper import groovy.json.JsonBuilder import org.apache.commons.net.syslog.SyslogParser import org.apache.commons.net.syslog.SyslogMessage def flowFile = session.get() if (!flowFile) return flowFile = session.write(flowFile, { inputStream, outputStream -> // 读取原始JSON内容 def originalJson = new JsonSlurper().parse(inputStream) // 获取syslog原始内容 def syslogRaw = flowFile.getAttribute('syslog_raw') if (syslogRaw) { // 解析syslog SyslogMessage syslogMsg = SyslogParser.createParser().parse(syslogRaw) // 将解析结果合并到原始JSON originalJson['syslog_timestamp'] = syslogMsg.timestamp originalJson['syslog_host'] = syslogMsg.hostname originalJson['syslog_content'] = syslogMsg.message } // 写入新JSON到流文件 outputStream.write(new JsonBuilder(originalJson).toByteArray()) } as StreamCallback) session.transfer(flowFile, REL_SUCCESS) - 注意:NiFi通常自带
commons-net依赖,若缺失需手动添加。
关键提示
ExtractGrok和ParseSyslog都支持从流文件属性读取内容,只需在处理器的“Source”选项中选择FlowFile Attribute并指定属性名,无需替换流文件内容。- 合并JSON时,可根据需求调整JOLT规则,比如要删除原
message字段,可在JOLT中添加删除操作。
内容的提问来源于stack exchange,提问作者Mogget
相关产品推荐
相关产品推荐

