如何在Airflow中定义DAG依赖?多DAG完成后触发BigQuery存储过程
实现思路:基于Airflow的多DAG依赖触发方案
针对你的需求,这里提供两种纯Airflow实现的方案,均无需使用Sensor,满足三个上游DAG全部成功后再触发精修层存储过程执行的要求:
方案一:利用Airflow Dataset触发(推荐,Airflow 2.4+)
Airflow 2.4引入的Dataset功能可以实现数据驱动的DAG调度,当所有依赖的Dataset都被上游任务更新后,下游DAG才会触发执行,完美匹配你的需求。
实现步骤
- 定义Dataset:为三个上游DAG的临时表分别创建Dataset,作为任务完成的标记
- 配置上游DAG:在每个上游DAG的最终数据加载任务中,通过
outlets关联对应的Dataset - 配置下游DAG:将精修层DAG的调度规则设置为依赖这三个Dataset,确保只有全部Dataset更新后才触发
代码示例
1. 定义全局Dataset
from airflow import Dataset # 对应三个临时表的Dataset,格式为BigQuery资源路径 dataset_staging_a = Dataset("bigquery://your-project/staging_dataset/staging_a") dataset_staging_b = Dataset("bigquery://your-project/staging_dataset/staging_b") dataset_staging_c = Dataset("bigquery://your-project/staging_dataset/staging_c")
2. 上游DAG(以DAG A为例)
from airflow.decorators import dag from airflow.providers.google.cloud.operators.bigquery import BigQueryInsertJobOperator from datetime import datetime @dag( schedule="0 9 * * *", # 每日上午9点调度 start_date=datetime(2024, 1, 1), catchup=False, tags=["staging_layer"] ) def dag_a(): # 数据加载到临时表任务 load_to_staging_a = BigQueryInsertJobOperator( task_id="load_staging_a", configuration={ "load": { "sourceUris": ["gs://your-bucket/path/to/a_data/*.parquet"], "destinationTable": { "projectId": "your-project", "datasetId": "staging_dataset", "tableId": "staging_a" }, "writeDisposition": "WRITE_TRUNCATE" } }, outlets=[dataset_staging_a] # 任务完成后标记Dataset更新 ) dag_a()
DAG B和DAG C的配置与DAG A一致,仅需替换对应的数据源路径、临时表名和关联的Dataset。
3. 精修层DAG(curated_sp_dag)
from airflow.decorators import dag from airflow.providers.google.cloud.operators.bigquery import BigQueryExecuteQueryOperator from datetime import datetime @dag( # 依赖三个Dataset,全部更新时触发 schedule=[dataset_staging_a, dataset_staging_b, dataset_staging_c], start_date=datetime(2024, 1, 1), catchup=False, tags=["curated_layer"] ) def curated_sp_dag(): # 执行Union All存储过程 execute_union_sp = BigQueryExecuteQueryOperator( task_id="execute_union_staging_sp", sql="CALL `your-project.curated_dataset.union_staging_tables`();", use_legacy_sql=False ) curated_sp_dag()
方案二:TriggerDagRunOperator + 状态检查(兼容旧版本Airflow)
如果你的Airflow版本低于2.4,可以采用上游触发+下游状态校验的方式:每个上游DAG成功后触发下游DAG,下游DAG先检查三个上游DAG的当日执行状态,全部成功再执行存储过程。
实现步骤
- 上游DAG添加触发任务:每个上游DAG的最后添加
TriggerDagRunOperator,触发精修层DAG - 下游DAG添加状态检查:在精修层DAG中,用Python任务检查三个上游DAG的当日执行状态是否为成功
- 执行存储过程:状态检查通过后,执行BigQuery存储过程
代码示例
1. 上游DAG添加触发任务(以DAG A为例)
from airflow.operators.trigger_dagrun import TriggerDagRunOperator # 数据加载任务(同方案一的load_to_staging_a) load_to_staging_a = BigQueryInsertJobOperator(...) # 触发精修层DAG trigger_curated_dag = TriggerDagRunOperator( task_id="trigger_curated_sp_dag", trigger_dag_id="curated_sp_dag", execution_date="{{ ds }}", # 传递当日执行日期 wait_for_completion=False, reset_dag_run=True ) load_to_staging_a >> trigger_curated_dag
DAG B和DAG C需添加相同的触发任务。
2. 精修层DAG(curated_sp_dag)
from airflow.decorators import dag, task from airflow.providers.google.cloud.operators.bigquery import BigQueryExecuteQueryOperator from airflow.models import DagRun from airflow.utils.state import State from datetime import datetime @dag( schedule=None, # 仅通过外部触发 start_date=datetime(2024, 1, 1), catchup=False, tags=["curated_layer"] ) def curated_sp_dag(): @task(retries=3, retry_delay=300) # 增加重试,避免上游DAG未完成时误判 def check_upstream_success(**context): execution_date = context["ds"] upstream_dag_ids = ["dag_a", "dag_b", "dag_c"] for dag_id in upstream_dag_ids: # 查询当日该DAG的执行记录 dag_runs = DagRun.find(dag_id=dag_id, execution_date=execution_date) if not dag_runs or dag_runs[0].state != State.SUCCESS: raise ValueError(f"上游DAG {dag_id} 在 {execution_date} 未成功执行") return True # 状态检查任务 upstream_check = check_upstream_success() # 执行存储过程任务 execute_union_sp = BigQueryExecuteQueryOperator( task_id="execute_union_staging_sp", sql="CALL `your-project.curated_dataset.union_staging_tables`();", use_legacy_sql=False ) upstream_check >> execute_union_sp curated_sp_dag()
方案对比
| 方案 | 适用版本 | 优点 | 注意事项 |
|---|---|---|---|
| Dataset触发 | Airflow 2.4+ | 原生支持,代码简洁,无需手动状态检查,触发逻辑精准 | 需升级到对应版本 |
| Trigger+状态检查 | 全版本 | 兼容旧版本,无需升级Airflow | 下游DAG会被触发三次,需配置重试机制避免误判 |
内容的提问来源于stack exchange,提问作者Sandeep Mohanty
相关产品推荐
相关产品推荐

