如何在Azure Data Factory中同时用数千个不同输入文件运行同一Python脚本
调整Azure Data Factory + Batch方案遍历输入Blob容器txt文件的实现方法
核心逻辑是通过ADF自带的元数据遍历能力批量触发Batch任务,每个任务独立处理单个txt文件,无需在单脚本中遍历所有文件,天然支持数千个文件的并行处理,效率更高。
步骤1:配置Get Metadata活动获取所有txt文件列表
在你现有ADF管道的最前端新增Get Metadata活动:
- 关联你输入Blob容器对应的存储数据集
- 活动配置的字段列表选择
Child Items,过滤器规则设置为*.txt,直接过滤出容器内所有后缀为txt的文件,无需后续额外筛选
步骤2:添加ForEach循环实现批量任务调度
- 将
Get Metadata活动的输出连接到新增的ForEach活动 - ForEach活动的
Items属性填写表达式:@activity('Get Metadata').output.childItems,自动遍历所有获取到的txt文件条目 - 如需并行处理多个文件,关闭ForEach的
Sequential开关,并行度最高可设置为50,数值可以和你的Batch池可用节点数匹配调整
步骤3:修改Batch自定义活动的启动参数
在ForEach活动内部放入你当前在用的Batch自定义活动,调整启动命令:
- 将当前遍历到的文件信息作为参数传给Batch任务,启动命令修改为:
python main.py @item().name - 如果你的输入文件存放在容器的子目录下,替换为
@item().fullPath传递完整路径即可
步骤4:适配Python脚本的输入逻辑
调整你的main.py代码,接收传入的文件参数,完成单个文件的下载、处理、上传逻辑,参考代码如下:
import sys from azure.storage.blob import BlobServiceClient # 你原本的文件处理逻辑不需要修改 def your_original_process_func(raw_content): # 原有处理代码 return processed_content if __name__ == "__main__": # 接收ADF传入的待处理文件名/路径 target_file = sys.argv[1] # 初始化存储客户端,替换为你自己的存储账户连接串、容器名 blob_client = BlobServiceClient.from_connection_string("YOUR_STORAGE_CONN_STR") # 从输入容器下载目标文件 input_blob = blob_client.get_blob_client(container="YOUR_INPUT_CONTAINER", blob=target_file) raw_txt = input_blob.download_blob().readall().decode("utf-8") # 调用原有逻辑处理 processed_txt = your_original_process_func(raw_txt) # 上传处理结果到输出容器 output_blob = blob_client.get_blob_client(container="YOUR_OUTPUT_CONTAINER", blob=f"processed_{target_file}") output_blob.upload_blob(processed_txt, overwrite=True)
可选优化项
- 如果txt文件体积大、总数量远超过50,可给Batch池开启自动缩放,根据待处理任务数自动增减计算节点,降低空闲成本
- 如需避免重复处理已完成的文件,可在Get Metadata活动之后新增过滤逻辑,对比输出容器已有的文件,筛掉已经处理过的条目再进入ForEach循环
内容的提问来源于stack exchange,提问作者James Mason
相关产品推荐
相关产品推荐

