ADF中Databricks输出JSON过大致ForEach循环失败的解决方案问询
问题背景
- 现有Azure Data Factory(ADF)管道流程:
- Databricks(DB)Notebook轮询挂载文件系统,基于
added元数据过滤出目标文件列表 - 通过
dbutils.notebook.exit(file_list_dict)将JSON格式的文件列表传递给ADF的ForEach活动 - ForEach遍历文件列表,逐个触发数据清洗插入的子管道
- Databricks(DB)Notebook轮询挂载文件系统,基于
- 核心问题:全量历史数据Ingestion时,文件列表JSON超过ADF限制的20MB导致管道失败;尝试用Lookup读取存储的JSON文件,但Lookup仅支持5000行,无法处理大量文件。
可行解决方案
方案一:Databricks Notebook直接分批调用ADF子管道API
跳过将全量文件列表传回ADF的步骤,在Notebook内将文件列表拆分为小批量,通过ADF的REST API逐个触发子管道,传递当前批次的文件列表。
操作步骤:
- 准备ADF调用所需的认证信息:租户ID、服务主体ID、服务主体密钥、订阅ID、资源组名称、ADF工厂名称、子管道名称
- 在Databricks Notebook中编写Python代码实现分批调用:
import requests import json from azure.identity import ClientSecretCredential # 配置认证与ADF信息 tenant_id = "你的租户ID" client_id = "你的服务主体ID" client_secret = "你的服务主体密钥" subscription_id = "你的订阅ID" resource_group = "资源组名称" factory_name = "ADF工厂名称" pipeline_name = "子管道名称" # 拆分文件列表为每1000个文件一组 file_list = [你的全量文件路径列表] batch_size = 1000 batches = [file_list[i:i+batch_size] for i in range(0, len(file_list), batch_size)] # 获取Azure认证Token credential = ClientSecretCredential(tenant_id, client_id, client_secret) token = credential.get_token("https://management.azure.com/.default").token # 循环触发子管道 for batch in batches: request_body = { "parameters": { "fileList": json.dumps(batch) } } api_url = f"https://management.azure.com/subscriptions/{subscription_id}/resourceGroups/{resource_group}/providers/Microsoft.DataFactory/factories/{factory_name}/pipelines/{pipeline_name}/createRun?api-version=2018-06-01" headers = { "Authorization": f"Bearer {token}", "Content-Type": "application/json" } response = requests.post(api_url, headers=headers, json=request_body) response.raise_for_status()
方案二:拆分文件列表为多个小JSON文件,嵌套ForEach遍历
将全量文件列表拆分为多个小JSON文件(每个文件包含≤5000个路径),通过ADF双层ForEach循环处理:外层遍历小JSON文件,内层遍历单个文件路径。
操作步骤:
- 修改Databricks Notebook逻辑,将过滤后的文件列表拆分为小文件存储:
import json file_list = [你的全量文件路径列表] batch_size = 4000 # 小于Lookup的5000行限制 output_path = "/mnt/your-storage/file-lists/" for idx, i in enumerate(range(0, len(file_list), batch_size)): batch = file_list[i:i+batch_size] with open(f"{output_path}file_batch_{idx}.json", "w") as f: json.dump(batch, f)
- 在ADF中配置活动链:
- Get Metadata活动:指向存储小JSON文件的路径,
Field list选择Child items,获取所有小JSON文件的路径 - 外层ForEach活动:遍历Get Metadata返回的子项列表
- Lookup活动:读取当前小JSON文件内容,关闭
First row only开关 - 内层ForEach活动:遍历Lookup返回的文件路径列表,触发子管道处理单个文件
- Lookup活动:读取当前小JSON文件内容,关闭
- Get Metadata活动:指向存储小JSON文件的路径,
方案三:将文件列表写入SQL表,分页遍历处理
借助SQL数据库存储文件列表,通过ADF分页查询的方式遍历所有记录,避开行数限制。
操作步骤:
- 在Databricks Notebook中将文件列表写入SQL表:
from pyspark.sql import SparkSession spark = SparkSession.builder.getOrCreate() # 构造文件列表DataFrame df = spark.createDataFrame([(path,) for path in file_list], ["file_path"]) # 写入SQL表(需提前创建表结构) df.write.mode("overwrite").jdbc( url="jdbc:sqlserver://your-server.database.windows.net:1433;databaseName=your-db", table="file_list_table", properties={ "user": "数据库用户名", "password": "数据库密码" } )
- 在ADF中配置活动链:
- Lookup活动1:查询表的总记录数,计算分页总数(例如每页5000条)
- ForEach活动:遍历分页序号,在循环内执行:
- Lookup活动2:用分页SQL查询当前页的文件路径,示例SQL:
SELECT file_path FROM file_list_table ORDER BY file_path OFFSET @pageOffset ROWS FETCH NEXT 5000 ROWS ONLY - 内层ForEach活动:遍历当前页的文件路径,触发子管道
- Lookup活动2:用分页SQL查询当前页的文件路径,示例SQL:
内容的提问来源于stack exchange,提问作者DatenBergwerker
相关产品推荐
相关产品推荐

