Airflow中基于其他DAG执行结果调度DAG运行的可行性咨询
当然可以实现!Airflow原生就支持这种基于其他DAG执行结果的触发逻辑,完全能满足你说的「dag1执行成功才触发dag2,失败则不触发」的需求。下面给你两种最常用的实现方案:
方案一:用TriggerDagRunOperator(兼容所有Airflow版本)
这是最直接的方式——在dag1的任务流末尾添加一个专门触发dag2的任务,并且设置只有当dag1的所有前置任务都成功时,这个触发任务才会执行。
举个代码例子:
from airflow import DAG from airflow.operators.dummy import DummyOperator from airflow.operators.trigger_dagrun import TriggerDagRunOperator from datetime import datetime with DAG( dag_id="dag1", start_date=datetime(2024, 1, 1), schedule_interval="@daily", catchup=False ) as dag1: task1 = DummyOperator(task_id="task1") task2 = DummyOperator(task_id="task2") # 触发dag2的任务,只有当前置任务全部成功才会执行 trigger_dag2 = TriggerDagRunOperator( task_id="trigger_dag2", trigger_dag_id="dag2", trigger_rule="all_success", # 明确指定只有所有前置成功才触发 wait_for_completion=False, # 不需要等待dag2执行完成,按需设置 dag=dag1 ) task1 >> task2 >> trigger_dag2
这样一来,只要dag1里的task1、task2有任何一个失败,trigger_dag2任务都不会运行,自然也就不会触发dag2。
方案二:用Airflow 2.0+的Dataset功能(更优雅的依赖管理)
如果你用的是Airflow 2.0及以上版本,推荐用Dataset来定义DAG之间的依赖关系——它是一种基于「数据产出」的调度方式,dag1成功产出某个Dataset后,自动触发依赖这个Dataset的dag2。
代码示例如下:
首先在dag1中定义产出的Dataset:
from airflow import DAG from airflow.operators.dummy import DummyOperator from airflow.models.dataset import Dataset from datetime import datetime # 定义一个Dataset,URI可以自定义,比如用文件路径或逻辑标识 my_dataset = Dataset("s3://my-bucket/dag1-output") with DAG( dag_id="dag1", start_date=datetime(2024, 1, 1), schedule_interval="@daily", catchup=False ) as dag1: task1 = DummyOperator(task_id="task1") # 标记这个任务产出my_dataset task2 = DummyOperator(task_id="task2", outlets=[my_dataset]) task1 >> task2
然后在dag2中设置依赖这个Dataset:
from airflow import DAG from airflow.operators.dummy import DummyOperator from airflow.models.dataset import Dataset from datetime import datetime my_dataset = Dataset("s3://my-bucket/dag1-output") with DAG( dag_id="dag2", start_date=datetime(2024, 1, 1), # 调度规则设置为依赖my_dataset,当Dataset被更新时自动触发 schedule=[my_dataset], catchup=False ) as dag2: task_a = DummyOperator(task_id="task_a") task_b = DummyOperator(task_id="task_b") task_a >> task_b
这种方式的好处是不需要在dag1里硬编码触发dag2的逻辑,而是通过Dataset来解耦两个DAG的依赖,更符合Airflow 2.x的设计理念。
额外注意事项
- 不管用哪种方案,都要确保执行dag1的Airflow角色拥有触发dag2的权限(可以在Airflow UI的「Security」-「Roles」里配置)
- 如果dag2本身有固定的调度周期,加上这种触发逻辑后,它会同时响应调度周期和外部触发,你可以根据需求调整
内容的提问来源于stack exchange,提问作者Kevin Nash
相关产品推荐
相关产品推荐

