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

如何在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工作区、集群(按需或持久)及参数。
  • 代码示例:
    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执行脚本。

二、选型核心考量因素

  • 数据量级:大数据量(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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 18:42:29