Airflow跨Super DAG依赖检查:如何实现10个任务全成功后触发下游DAG
方案1:在Super DAG中添加汇总任务(兼容所有Airflow版本)
在Super DAG里新增一个任务,依赖所有10个业务任务,并且设置trigger_rule='all_success'——只有当所有上游任务都成功时,这个汇总任务才会执行成功。然后第二个DAG的External Sensor只需要检查这个汇总任务的成功状态即可。
Super DAG 代码示例:
from airflow.models import DAG from airflow.operators.dummy import DummyOperator from datetime import datetime with DAG('super_dag', start_date=datetime(2024, 1, 1), schedule_interval='@daily') as dag: # 定义10个业务任务 task1 = DummyOperator(task_id='task1') task2 = DummyOperator(task_id='task2') # ... 依次定义task3到task9 task10 = DummyOperator(task_id='task10') # 汇总任务:仅当所有10个任务成功时才会成功 all_tasks_success = DummyOperator( task_id='all_tasks_success', trigger_rule='all_success' ) # 让所有业务任务都指向汇总任务 [task1, task2, task3, task4, task5, task6, task7, task8, task9, task10] >> all_tasks_success
第二个DAG 代码示例:
from airflow.models import DAG from airflow.sensors.external_task import ExternalTaskSensor from datetime import datetime with DAG('second_dag', start_date=datetime(2024, 1, 1), schedule_interval='@daily') as dag: # 等待Super DAG的汇总任务成功 wait_for_super_dag = ExternalTaskSensor( task_id='wait_for_super_dag', external_dag_id='super_dag', external_task_id='all_tasks_success', allowed_states=['success'], failed_states=['failed', 'skipped'], mode='poke' ) # 第二个DAG的后续任务 downstream_task = DummyOperator(task_id='downstream_task') wait_for_super_dag >> downstream_task
优势:逻辑直观,所有Airflow版本都支持,不需要修改现有业务任务的依赖关系;劣势:需要对Super DAG做少量修改。
方案2:在第二个DAG中检查所有10个任务的状态(无需修改Super DAG)
如果不想改动Super DAG,可以用Airflow 2.0+支持的external_task_ids参数,让External Sensor同时检查Super DAG的10个任务,只有当所有任务都成功时才触发后续流程。
第二个DAG 代码示例:
from airflow.models import DAG from airflow.sensors.external_task import ExternalTaskSensor from datetime import datetime with DAG('second_dag', start_date=datetime(2024, 1, 1), schedule_interval='@daily') as dag: # 等待Super DAG的所有10个任务成功 wait_for_all_super_tasks = ExternalTaskSensor( task_id='wait_for_all_super_tasks', external_dag_id='super_dag', external_task_ids=['task1', 'task2', 'task3', 'task4', 'task5', 'task6', 'task7', 'task8', 'task9', 'task10'], allowed_states=['success'], failed_states=['failed', 'skipped'], mode='poke', trigger_rule='all_success' ) downstream_task = DummyOperator(task_id='downstream_task') wait_for_all_super_tasks >> downstream_task
如果你的Airflow版本低于2.0,可以创建10个独立的ExternalTaskSensor(每个对应一个任务),然后用一个DummyOperator汇总,设置trigger_rule='all_success':
# 示例:创建10个sensor task1_sensor = ExternalTaskSensor(task_id='wait_task1', external_dag_id='super_dag', external_task_id='task1', ...) task2_sensor = ExternalTaskSensor(task_id='wait_task2', external_dag_id='super_dag', external_task_id='task2', ...) # ... 直到task10_sensor # 汇总sensor结果 all_sensors_success = DummyOperator(task_id='all_sensors_success', trigger_rule='all_success') [task1_sensor, task2_sensor, ..., task10_sensor] >> all_sensors_success all_sensors_success >> downstream_task
优势:无需修改Super DAG;劣势:Airflow版本受限(低版本需要写重复代码),任务列表较长时配置繁琐。
方案3:使用Airflow Dataset(Airflow 2.4+)
Airflow 2.4+引入的Dataset功能可以基于数据事件触发DAG。我们可以让Super DAG的汇总任务(带all_success触发规则)成功后生成一个Dataset事件,第二个DAG直接依赖这个Dataset即可。
Super DAG 代码示例:
from airflow.models import DAG, Dataset from airflow.operators.dummy import DummyOperator from datetime import datetime # 定义Dataset(标识可以是任意字符串,不一定对应实际文件) all_tasks_success_dataset = Dataset('super_dag_all_tasks_success') with DAG('super_dag', start_date=datetime(2024, 1, 1), schedule_interval='@daily') as dag: # 10个业务任务定义... task1 = DummyOperator(task_id='task1') # ... task2到task10 all_tasks_success = DummyOperator( task_id='all_tasks_success', trigger_rule='all_success', outlets=[all_tasks_success_dataset] # 任务成功时触发Dataset事件 ) [task1, ..., task10] >> all_tasks_success
第二个DAG 代码示例:
from airflow.models import DAG, Dataset from airflow.operators.dummy import DummyOperator from datetime import datetime all_tasks_success_dataset = Dataset('super_dag_all_tasks_success') # 直接依赖Dataset,只有当Dataset事件触发时才运行 with DAG('second_dag', start_date=datetime(2024, 1, 1), schedule=[all_tasks_success_dataset]) as dag: downstream_task = DummyOperator(task_id='downstream_task')
优势:无需使用Sensor,代码更简洁,符合Airflow现代化的事件驱动理念;劣势:需要升级到Airflow 2.4及以上版本。
关键注意点
- 无论用哪种方案,核心都是确保只有10个任务全部成功才会触发第二个DAG,所以必须严格使用
trigger_rule='all_success'(或等价的检查逻辑),避免任务跳过导致的误触发。 - 如果Super DAG中存在任务可能被主动跳过(比如用
ShortCircuitOperator或手动跳过),all_success规则会让汇总任务(或Sensor)无法成功,正好符合需求。
内容的提问来源于stack exchange,提问作者hellopa1

