使用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)实现以下需求:
- 按
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的所有数据合并为单个文件。
解决方案
一、配置基础数据集
- 着陆区数据集:创建指向存储着陆区JSON文件的数据集(如Blob存储/ADLS Gen2的JSON数据集),开启递归遍历,并设置文件筛选为
*.json,确保只处理目标文件。 - 目标存储数据集:创建指向最终存储位置的数据集,选择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格式选择"数组",确保合并后的数据以数组形式存储。
- 如果不做合并:选择按分区存储,分区列指定为
三、创建管道与调度
- 构建管道:新建管道,添加数据流活动,关联上面创建的数据流。
- 设置触发器:
- 若和API拉取频率匹配,选择定时触发器(比如每5分钟运行一次)。
- 若要实时处理新文件,选择事件触发器(监听着陆区的文件创建事件)。
- 增量处理优化:如果用定时触发器,可在着陆区数据集的筛选条件中添加
lastModified > @pipeline().parameters.lastRunTime,结合管道参数记录上次运行时间,避免重复处理文件。
内容的提问来源于stack exchange,提问作者Wouter De Raeve
相关产品推荐
相关产品推荐

