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

使用Azure Data Factory按属性拆分JSON文件并生成层级存储结构

问题

我会定期(每5分钟或更频繁)从API拉取数据,返回的JSON格式示例如下:

{"timestamp":"2022-09-28T00:33:53Z",
  "data":[
    {
      "id":"bdcb2ad8-9e19-4468-a4f3-b440de3a7b40",
      "value":3,
      "created": "2020-09-28T00:00:00Z"
    },
    {
      "id":"7f8d07eb-433b-404c-a9b3-f1832bdd780f",
      "value":4,
      "created": "2020-09-28T00:00:00Z"
    },
    {
      "id":"7f8d07eb-433b-404c-a9b3-f1832bdd780f",
      "value":6,
      "created": "2020-09-28T00:05:00Z"
    }
  ]
}

着陆区会生成一系列以拉取时间命名的JSON文件(如2022-09-28T00:33:53Z.json、2022-09-28T00:43:44Z.json)。

我需要用Azure Data Factory(ADF)实现以下需求:

  1. 按id属性拆分文件,生成如下层级存储结构:
/bdcb2ad8-9e19-4468-a4f3-b440de3a7b40/2022-09-28T00:33:53Z.json
/bdcb2ad8-9e19-4468-a4f3-b440de3a7b40/2022-09-28T00:43:44Z.json
/7f8d07eb-433b-404c-a9b3-f1832bdd780f/2022-09-28T00:33:53Z.json
/7f8d07eb-433b-404c-a9b3-f1832bdd780f/2022-09-28T00:43:44Z.json

每个目标文件仅包含对应id的相关数据。
2. 额外需求(非必需):每次ADF运行时,将同一id的所有数据合并为单个文件。

解决方案

一、配置基础数据集

  1. 着陆区数据集:创建指向存储着陆区JSON文件的数据集(如Blob存储/ADLS Gen2的JSON数据集),开启递归遍历,并设置文件筛选为*.json,确保只处理目标文件。
  2. 目标存储数据集:创建指向最终存储位置的数据集,选择JSON格式,暂不设置具体路径,后续通过数据流动态指定。

二、构建数据流实现拆分/合并逻辑

数据流是实现该需求的核心组件,按以下步骤配置:

1. 源节点:读取并展开嵌套数据

  • 连接着陆区数据集,在JSON设置中,将"数组根"指定为$data,直接把每个JSON文件里的data数组展开为一行一行的明细数据。
  • 添加派生列,命名为source_filename,用表达式split(name(), '.')[0]提取源文件的时间戳(比如从2022-09-28T00:33:53Z.json中取出2022-09-28T00:33:53Z),用于后续命名目标文件。

2. 可选:聚合节点实现同ID数据合并

如果需要将同一id的所有数据合并为单个文件,添加聚合节点:

  • 分组键选择id字段。
  • 添加聚合列,命名为merged_data,用表达式collect(@{})将当前id的所有明细数据收集为一个数组。

3. Sink节点:输出到目标层级结构

  • 连接目标存储数据集,在设置标签页:
    • 如果不做合并:选择按分区存储,分区列指定为id,ADF会自动创建以id命名的文件夹;文件命名规则设置为@concat(source_filename, '.json'),确保每个源文件对应目标文件夹下的同名文件,内容为该id在源文件中的所有数据。
    • 如果做合并:同样按id分区,文件命名规则设置为@concat(id, '.json')(或添加时间后缀如@concat(id, '_', pipeline().RunId, '.json')避免重复),JSON格式选择"数组",确保合并后的数据以数组形式存储。

三、创建管道与调度

  1. 构建管道:新建管道,添加数据流活动,关联上面创建的数据流。
  2. 设置触发器:
    • 若和API拉取频率匹配,选择定时触发器(比如每5分钟运行一次)。
    • 若要实时处理新文件,选择事件触发器(监听着陆区的文件创建事件)。
  3. 增量处理优化:如果用定时触发器,可在着陆区数据集的筛选条件中添加lastModified > @pipeline().parameters.lastRunTime,结合管道参数记录上次运行时间,避免重复处理文件。

内容的提问来源于stack exchange,提问作者Wouter De Raeve

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 06:25:22