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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 21:54:04