如何让Airflow主DAG跳过调用已禁用的子DAG?
解决Airflow父DAG自动跳过已禁用子DAG的问题
问题背景
禁用dbt_api_dag后,父DAG main 因TriggerDagRunOperator持续等待已禁用的子DAG完成,最终挂起并标记为失败,需手动干预。
解决方案
通过在触发子DAG前检查其启用状态,自动跳过已禁用的子DAG任务,避免父DAG挂起。核心思路是用ShortCircuitOperator控制子DAG触发任务的执行逻辑:
步骤1:导入必要模块
在原有代码基础上添加以下导入:
from airflow.models import DagModel from airflow.operators.python import ShortCircuitOperator import logging
步骤2:添加DAG状态检查函数
定义检查dbt_api_dag是否启用的函数:
def check_dbt_dag_enabled(): logger = logging.getLogger(__name__) # 获取目标DAG的元数据 dbt_dag = DagModel.get_dagmodel(dag_id="dbt_api_dag") if not dbt_dag: logger.info("子DAG dbt_api_dag不存在,跳过执行") return False if not dbt_dag.is_active: logger.info("子DAG dbt_api_dag已被禁用,跳过执行") return False logger.info("子DAG dbt_api_dag已启用,将触发执行") return True
步骤3:修改main DAG任务结构
调整任务依赖,用ShortCircuitOperator控制dbt任务的执行:
@dag( start_date=datetime(2019, 1, 1), max_active_runs=1, schedule_interval=schedule(daily_schedule), # Run at 5:00am MST default_args=default_args, catchup=False, concurrency=4, tags=['extract', 'main'] ) def main(): with TaskGroup(group_id='extract') as extract: TriggerDagRunOperator( task_id='extract', trigger_dag_id='extract', wait_for_completion=True, trigger_rule="all_done" # Run even if previous tasks failed ) # 新增检查子DAG状态的任务 check_dbt_enabled = ShortCircuitOperator( task_id='check_dbt_enabled', python_callable=check_dbt_dag_enabled, trigger_rule="all_done" ) dbt = TriggerDagRunOperator( task_id='dbt', trigger_dag_id='dbt_api_dag', wait_for_completion=True, trigger_rule="all_done" ) # 调整任务依赖 extract >> check_dbt_enabled >> dbt taskflow = main()
逻辑说明
ShortCircuitOperator会根据check_dbt_dag_enabled的返回值决定是否执行后续的dbt任务:- 返回
True:继续执行dbt任务,正常触发子DAG - 返回
False:自动跳过dbt任务,父DAG不会等待未启用的子DAG,直接完成流程
- 返回
- 函数中添加了日志输出,便于排查任务跳过的原因
内容的提问来源于stack exchange,提问作者Jose Pla
相关产品推荐
相关产品推荐

