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

如何在NiFi中解析JSON的message字段并将结果整合回内容?

可行解决方案

方案1:纯处理器组合(无需代码,适合新手)

步骤分解:

  1. 保留原始JSON+提取syslog字段到属性
    用EvaluateJSONPath处理器,做两项配置:

    • 规则1:属性名设为syslog_raw,值为$.message(把JSON里的message字段提取为流文件属性)
    • 规则2:属性名设为original_json,值为$(把整个JOLT处理后的JSON提取为属性,用于后续合并)
    • 将“Destination”设置为flowfile-attribute
  2. 解析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)。
  3. 把解析后的属性转成JSON
    用AttributesToJSON处理器:

    • “Attributes to Include”:如果用ParseSyslog填syslog.*,如果用ExtractGrok填你定义的字段名(比如timestamp,host,content)
    • “Destination”选flowfile-attribute,属性名设为parsed_syslog
  4. 合并原始JSON和解析后的JSON
    用JOLTTransformJSON处理器:

    • “JOLT Specification”使用合并规则:
      [
        {
          "operation": "shift",
          "spec": {
            "original_json": {
              "*": "&"
            },
            "parsed_syslog": {
              "*": "&"
            }
          }
        }
      ]
      
    • “Input JSON Path”设为$,处理器会读取流文件属性中的original_json和parsed_syslog,合并成新的JSON对象输出到流文件内容。

方案2:用ExecuteScript脚本灵活处理(适合自定义逻辑场景)

如果觉得处理器组合繁琐,可直接用ExecuteScript(选Groovy语言)完成解析+合并:

  1. 先通过EvaluateJSONPath把message字段提取为属性syslog_raw。
  2. 配置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)
    
  3. 注意:NiFi通常自带commons-net依赖,若缺失需手动添加。

关键提示

  • ExtractGrok和ParseSyslog都支持从流文件属性读取内容,只需在处理器的“Source”选项中选择FlowFile Attribute并指定属性名,无需替换流文件内容。
  • 合并JSON时,可根据需求调整JOLT规则,比如要删除原message字段,可在JOLT中添加删除操作。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 02:01:28