如何优化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
相关产品推荐
相关产品推荐

