如何在Azure Data Factory中运行Python ETL脚本并选择最优方案?
在Azure Data Factory中运行Python ETL脚本的方法及选型考量
一、ADF中运行Python ETL脚本的核心方法
1. Azure Function活动
把Python ETL脚本封装成Azure Function(推荐HTTP触发类型),通过ADF的Azure Function活动调用执行。
- 适用场景:中小数据量ETL、需要事件触发(如Blob上传触发)、希望按需缩放的场景。
- 实现要点:
- 用Python函数模板编写ETL逻辑,通过环境变量配置存储/数据库连接信息。
- ADF活动中配置函数APP、函数名,传递数据源路径、目标位置等参数。
- 代码示例:
import azure.functions as func import pandas as pd import os from azure.storage.blob import BlobServiceClient def main(req: func.HttpRequest) -> func.HttpResponse: # 接收ADF传递的参数 source_blob = req.params.get('source_blob') target_container = req.params.get('target_container') # 初始化存储客户端 conn_str = os.getenv("AZURE_STORAGE_CONNECTION_STRING") blob_service_client = BlobServiceClient.from_connection_string(conn_str) # 抽取:读取原始数据 raw_blob_client = blob_service_client.get_blob_client(container="raw", blob=source_blob) df = pd.read_csv(raw_blob_client.download_blob()) # 转换:清洗数据 df_clean = df.dropna(subset=["user_id"]).assign(processed_at=pd.Timestamp.now()) # 加载:写入处理后数据 processed_blob_client = blob_service_client.get_blob_client(container=target_container, blob=f"processed_{source_blob}") processed_blob_client.upload_blob(df_clean.to_csv(index=False), overwrite=True) return func.HttpResponse("ETL任务执行完成", status_code=200)
2. 自定义活动(HDInsight/Azure Batch)
(1)HDInsight自定义活动
将Python脚本(支持PySpark分布式逻辑)上传到ADLS/Blob存储,通过ADF调用HDInsight集群(按需或持久)运行脚本。
- 适用场景:TB级以上大数据量ETL、需要分布式处理的场景。
- 实现要点:指定脚本的ADLS路径(如
abfss://container@storageaccount.dfs.core.windows.net/scripts/etl_spark.py),传递运行参数。
(2)Azure Batch自定义活动
创建Batch池,上传脚本和依赖包,ADF调用Batch任务执行Python脚本。
- 适用场景:需要特定计算资源(如GPU、自定义Python版本)、批量处理多任务的场景。
3. Databricks Notebook活动
将PySpark或Python ETL逻辑写入Databricks Notebook,通过ADF的Databricks Notebook活动调用执行。
- 适用场景:大数据量分布式ETL、希望托管式集群管理的场景。
- 实现要点:
- Notebook中通过
dbutils.widgets.get()获取ADF传递的参数。 - ADF活动中配置Databricks工作区、集群(按需或持久)及参数。
- Notebook中通过
- 代码示例:
from pyspark.sql import SparkSession import pyspark.sql.functions as F # 获取ADF参数 source_path = dbutils.widgets.get("source_path") target_path = dbutils.widgets.get("target_path") # 初始化Spark会话 spark = SparkSession.builder.appName("Python ETL").getOrCreate() # 抽取数据 df = spark.read.csv(source_path, header=True, inferSchema=True) # 转换数据 df_clean = df.filter(F.col("user_id").isNotNull()).withColumn("processed_at", F.current_timestamp()) # 加载数据 df_clean.write.mode("overwrite").csv(target_path, header=True)
4. Shell活动(配合VM/ACI)
利用ADF的Shell活动,通过SSH连接Azure VM,或调用Azure Container Instance(ACI)运行Python脚本。
- 适用场景:已有现成Python环境、需要自定义环境配置的场景。
- 实现要点:
- VM场景:开放SSH端口,ADF配置SSH连接信息,执行命令如
python /opt/scripts/etl.py --source /data/raw.csv。 - ACI场景:打包Python环境为Docker镜像,ADF调用ACI执行脚本。
- VM场景:开放SSH端口,ADF配置SSH连接信息,执行命令如
二、选型核心考量因素
- 数据量级:大数据量(TB+)优先选Databricks/HDInsight;中小数据量选Azure Function/Shell活动。
- 资源需求:需要特定库、GPU或自定义环境,选Azure Batch/ACI/VM;通用场景选Azure Function/Databricks。
- 成本控制:按需付费优先选Azure Function(按调用次数)、Databricks按需集群;长期稳定任务可选持久集群。
- 集成复杂度:需与ADF其他活动(复制、数据流)紧密集成,选Databricks Notebook/Azure Function活动,参数传递和状态监控更顺畅。
- 维护成本:托管式服务(Azure Function/Databricks)维护量低;HDInsight/Batch/VM需自行维护集群/环境,成本较高。
- 触发需求:事件触发(如Blob上传)选Azure Function;定时触发可搭配ADF管道触发器使用任意活动。
内容的提问来源于stack exchange,提问作者Aadhil Imam
相关产品推荐
相关产品推荐

