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

如何优化Airflow DAG减少代码重复?多ADF并行执行场景

优化Airflow DAG:并行执行多个ADF管道并消除代码冗余

问题说明

需要实现的DAG逻辑是:先并行运行多个仅名称不同的Azure Data Factory管道,待所有ADF管道完成后,再执行Databricks Notebook。原实现通过重复定义AzureDataFactoryRunPipelineOperator导致代码冗余;尝试用PythonOperator循环调用算子的方式会导致ADF任务串行执行,且无法被Airflow识别为独立并行任务,不符合需求。

最佳实践方案

核心思路是在DAG的上下文解析阶段,通过循环动态生成多个AzureDataFactoryRunPipelineOperator实例,每个实例对应一个独立的ADF管道任务,这样Airflow会将它们识别为并行节点,同时消除代码重复。

优化后的完整代码

from datetime import datetime, timedelta
from typing import List
from airflow.models import DAG
from airflow.providers.microsoft.azure.operators.data_factory import AzureDataFactoryRunPipelineOperator
from airflow.providers.databricks.operators.databricks import DatabricksSubmitRunOperator

# 定义所有需要执行的ADF管道名称列表
ADF_PIPELINE_NAMES: List[str] = [
    "run_adf_pipeline_1",
    "run_adf_pipeline_2",
    "run_adf_pipeline_3",
    "run_adf_pipeline_4"
]

# 假设default_args、DAG_ID、DATABRICKS_NOTEBOOK_PATH已提前定义
with DAG(
    dag_id=DAG_ID,
    start_date=datetime(2023, 3, 15),
    schedule="@daily",
    catchup=False,
    default_args=default_args,
    default_view="graph",
    tags=["development", "azure", "databricks"]
) as dag:
    # 循环生成ADF管道任务
    adf_tasks = []
    for pipeline_name in ADF_PIPELINE_NAMES:
        task = AzureDataFactoryRunPipelineOperator(
            task_id=pipeline_name,
            pipeline_name=pipeline_name,
            # 可添加其他通用参数,如azure_data_factory_conn_id等
        )
        adf_tasks.append(task)

    # 定义Databricks Notebook任务
    new_cluster = {
        "spark_version": "10.4.x-scala2.12",
        "node_type_id": "Standard_DS3_v2",
        "num_workers": 1,
    }

    notebook_task_params = {
        "new_cluster": new_cluster,
        "notebook_task": {
            "notebook_path": DATABRICKS_NOTEBOOK_PATH,
        },
    }

    run_databricks_notebook = DatabricksSubmitRunOperator(
        task_id="run_databricks_notebook", json=notebook_task_params
    )

    # 设置依赖:所有ADF任务完成后执行Databricks任务
    adf_tasks >> run_databricks_notebook

关键说明

  • 动态生成算子:在DAG上下文内循环创建AzureDataFactoryRunPipelineOperator,每个实例都是Airflow的独立任务节点,天然支持并行执行。
  • 消除冗余:通过管道名称列表统一管理所有需要执行的ADF管道,新增/删除管道只需修改列表,无需重复编写算子定义代码。
  • 依赖设置简化:将所有ADF任务存入列表后,直接通过列表 >> 下游任务的方式设置依赖,代码更简洁。

为什么之前的PythonOperator方式不可行

  • PythonOperator的python_callable是在任务运行阶段执行的,而Airflow识别任务节点是在DAG解析阶段(即DAG上下文初始化时),因此在callable里创建的算子不会被Airflow注册为独立任务。
  • 即使忽略注册问题,callable里的循环执行是在单个Airflow任务实例内完成的,必然是串行执行,无法实现并行需求。

内容的提问来源于stack exchange,提问作者Luiz

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 16:07:55