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

NiFi中实现JSON转指定格式Flatfile的方案咨询

问题描述

输入JSON

{
  "BAG": {
    "Ett": "WWWLOG",
    "refD": "DVJOB1415738",
    "TO": "JOB",
    "tst": ""
  },
  "TC": {
    "Command": "DVJOB1415738",
    "dateCmd": "20/11/2026 10:25:16",
    "codeTransporteur": "999",
    "RProjet": "62165SUN_XXXXX_XXXXXX"
  },
  "buyer": {
    "id": "929292",
    "CMDM": "867220 SAN ANDREAS",
    "Exp": "021",
    "refB": "CIN1887881",
    "lieuL": "867220"
  },
  "CntShip": {
    "TC": "2",
    "cleR": "WWW",
    "FN": "ITASM",
    "LN": "DATA1",
    "Phone": "07********",
    "mail": "data.iteam@xxxx.net"
  },
  "LC": [
    {
      "N": 1,
      "CArt": "3561920920832",
      "CArtDAPO": "",
      "Q": 1
    },
    {
      "N": 2,
      "CArt": "3561920936604",
      "CArtDAPO": "",
      "Q": 1
    },
    {
      "N": 3,
      "CArt": "2619000120519",
      "CArtDAPO": "",
      "Q": 1
    }
  ]
}

目标Flatfile格式

11|SHIP2023112009251776006143|WWWLOG|DVJOB1415738|WWW|OFFER|JOB||20251120102517|
12|SHIP2023112009251776006143|WWWLOG|DVJOB1415738|DVJOB1415738|20231120092516|999||929292|867220|||867220 SAN ANDREAS|021|||CIN1887881|
15|SHIP2023112009251776006143|WWWLOG|DVJOB1415738|2|WWW|ITASM|DATA1|07********||data.iteam@xxxx.net|
8|SHIP2023112009251776006143|WWWLOG|DVJOB1415738|1|3561920956992|||1|||1||||62165SUN_XXXXX_XXXXXX||||||
8|SHIP2023112009251776006143|WWWLOG|DVJOB1415738|2|3561920936604|||1|||2||||62165SUN_XXXXX_XXXXXX||||||
8|SHIP2023112009251776006143|WWWLOG|DVJOB1415738|3|2619000120519|||1|||3||||62165SUN_XXXXX_XXXXXX||||||

转换规则

  • 每行首数字对应JSON层级:11对应BAG,12对应TC+buyer,15对应CntShip,8对应LC下每条商品数据
  • 生成统一的SHIP+yyyyMMddHHMMSS+8位随机数格式MSG_N,所有行保持一致
  • 第2、3、4字段在所有行中相同
  • 补充输入JSON未提及的非空字段(如第一行第6位的"OFFER")
  • 空字段保留竖线占位符,匹配目标格式的所有字段

实现方案

仅用JoltTransformation是否足够?

不行。Jolt仅能处理JSON结构的转换,无法生成全局唯一的随机MSG_N、无法将JSON数组展开为多行文本、也无法完成最终的竖线分隔文本拼接,必须搭配其他NiFi处理器才能完成全流程。

需要搭配的处理器

  • UpdateAttribute:生成全局唯一的MSG_N和共享属性
  • JoltTransformationJSON:将原始JSON转换为包含所有行数据的结构化JSON数组
  • SplitJson:将Jolt输出的数组拆分为单个JSON对象
  • RouteOnAttribute:按行类型分流,适配不同的文本模板
  • ReplaceText:将单个JSON对象转换为竖线分隔的文本行
  • MergeContent:将所有文本行合并为最终的Flatfile

具体步骤

1. 使用UpdateAttribute生成全局变量

为FlowFile添加以下共享属性:

  • msg_n:值设为SHIP${now():format("yyyyMMddHHmmss")}${random():format("00000000")},生成符合要求的统一MSG_N
  • ett:值设为${json-path:read('$.BAG.Ett')},提取BAG.Ett的值
  • ref_d:值设为${json-path:read('$.BAG.refD')},提取BAG.refD的值
  • r_projet:值设为${json-path:read('$.TC.RProjet')},提取TC.RProjet的值

2. JoltTransformationJSON转换结构

使用以下Jolt规范,将原始JSON转换为包含所有行数据的数组,同时填充固定字段和引用全局属性:

[
  {
    "operation": "shift",
    "spec": {
      "BAG": {
        "TO": "rows[0].to"
      },
      "TC": {
        "Command": "rows[1].command",
        "dateCmd": "rows[1].dateCmd",
        "codeTransporteur": "rows[1].codeTrans"
      },
      "buyer": {
        "id": "rows[1].buyerId",
        "lieuL": "rows[1].lieuL",
        "CMDM": "rows[1].cmdm",
        "Exp": "rows[1].exp",
        "refB": "rows[1].refB"
      },
      "CntShip": {
        "TC": "rows[2].tc",
        "cleR": "rows[2].cleR",
        "FN": "rows[2].fn",
        "LN": "rows[2].ln",
        "Phone": "rows[2].phone",
        "mail": "rows[2].mail"
      },
      "LC": {
        "*": {
          "N": "rows[3&].n",
          "CArt": "rows[3&].cart",
          "Q": "rows[3&].q"
        }
      }
    }
  },
  {
    "operation": "default",
    "spec": {
      "rows": {
        "0": {
          "type": "11",
          "cleR": "WWW",
          "fixedVal": "OFFER",
          "date": "20251120102517"
        },
        "1": {
          "type": "12"
        },
        "2": {
          "type": "15"
        },
        "*": {
          "type": "8"
        }
      }
    }
  },
  {
    "operation": "modify-overwrite-beta",
    "spec": {
      "rows": {
        "*": {
          "msg_n": "${msg_n}",
          "ett": "${ett}",
          "refD": "${ref_d}",
          "rProjet": "${r_projet}"
        }
      }
    }
  }
]

转换后会得到一个包含所有行数据的JSON数组,每个元素对应目标Flatfile的一行。

3. SplitJson拆分数组

配置Json Path Expression为$.rows,将数组中的每个元素拆分为单独的FlowFile。

4. RouteOnAttribute按行类型分流

添加路由规则,根据type字段的值(11/12/15/8)将FlowFile分到不同的处理分支,方便后续使用不同的文本模板。

5. ReplaceText生成单行文本

针对每个分支,使用ReplaceText的正则替换模式,将单个JSON对象转换为竖线分隔的文本:

  • 11类型分支替换模板:
    ${json-path:read('$.type')}|${json-path:read('$.msg_n')}|${json-path:read('$.ett')}|${json-path:read('$.refD')}|${json-path:read('$.cleR','')}|${json-path:read('$.fixedVal','')}|${json-path:read('$.to','')}||${json-path:read('$.date','')}|
    
  • 12类型分支替换模板:
    ${json-path:read('$.type')}|${json-path:read('$.msg_n')}|${json-path:read('$.ett')}|${json-path:read('$.refD')}|${json-path:read('$.command','')}|${json-path:read('$.dateCmd':replaceAll('/',''):replaceAll(' ',''))}|${json-path:read('$.codeTrans','')}||${json-path:read('$.buyerId','')}|${json-path:read('$.lieuL','')}|||${json-path:read('$.cmdm','')}|${json-path:read('$.exp','')}|||${json-path:read('$.refB','')}|
    
  • 15类型分支替换模板:
    ${json-path:read('$.type')}|${json-path:read('$.msg_n')}|${json-path:read('$.ett')}|${json-path:read('$.refD')}|${json-path:read('$.tc','')}|${json-path:read('$.cleR','')}|${json-path:read('$.fn','')}|${json-path:read('$.ln','')}|${json-path:read('$.phone','')}||${json-path:read('$.mail','')}|
    
  • 8类型分支替换模板:
    ${json-path:read('$.type')}|${json-path:read('$.msg_n')}|${json-path:read('$.ett')}|${json-path:read('$.refD')}|${json-path:read('$.n','')}|${json-path:read('$.cart','')}|||${json-path:read('$.q','')}|||${json-path:read('$.n','')}||||${json-path:read('$.rProjet','')}||||||
    

注:日期格式转换可在Jolt阶段提前处理,或在ReplaceText中用表达式完成。

6. MergeContent合并为Flatfile

配置MergeContent:

  • Delimiter Strategy:Text,设置为换行符\n
  • Minimum Number of Entries:1
  • Maximum Number of Entries:根据实际数据量设置,确保所有行都被合并
  • Attribute Strategy:Keep Only Common Attributes

最终得到符合要求的Flatfile。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 21:47:10