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

如何在Airflow中定义DAG依赖?多DAG完成后触发BigQuery存储过程

实现思路:基于Airflow的多DAG依赖触发方案

针对你的需求,这里提供两种纯Airflow实现的方案,均无需使用Sensor,满足三个上游DAG全部成功后再触发精修层存储过程执行的要求:

方案一:利用Airflow Dataset触发(推荐,Airflow 2.4+)

Airflow 2.4引入的Dataset功能可以实现数据驱动的DAG调度,当所有依赖的Dataset都被上游任务更新后,下游DAG才会触发执行,完美匹配你的需求。

实现步骤

  1. 定义Dataset:为三个上游DAG的临时表分别创建Dataset,作为任务完成的标记
  2. 配置上游DAG:在每个上游DAG的最终数据加载任务中,通过outlets关联对应的Dataset
  3. 配置下游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的当日执行状态,全部成功再执行存储过程。

实现步骤

  1. 上游DAG添加触发任务:每个上游DAG的最后添加TriggerDagRunOperator,触发精修层DAG
  2. 下游DAG添加状态检查:在精修层DAG中,用Python任务检查三个上游DAG的当日执行状态是否为成功
  3. 执行存储过程:状态检查通过后,执行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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 12:16:04