如何实现TriggerDagRunOperator的conditionally_trigger函数,dag_b暂停时跳过触发?
解决TriggerDagRunOperator在目标DAG暂停时跳过触发的问题
核心思路是借助Airflow内置的DagModel类获取目标DAG的实时状态,在自定义的conditionally_trigger函数中判断是否执行触发操作。
步骤1:实现自定义触发判断函数
这个函数会检查dag_b的暂停状态,仅当它处于开启状态时才触发:
from airflow.models import DagModel def conditionally_trigger(context, dag_run_obj): # 从元数据库中获取目标DAG的状态模型 target_dag = DagModel.get_dagmodel(dag_id='dag_b') # 检查目标DAG是否未暂停(处于开启状态) if not target_dag.is_paused: return dag_run_obj # 返回触发对象,执行触发逻辑 else: # 可选:通过XCom记录本次跳过触发的状态,方便后续排查 context['ti'].xcom_push(key='trigger_skipped', value=True) return None # 返回None,跳过触发动作
步骤2:在dag_a中配置TriggerDagRunOperator
将自定义函数传入trigger_callable参数,替换默认的触发逻辑:
from airflow import DAG from airflow.operators.trigger_dagrun import TriggerDagRunOperator from datetime import datetime with DAG( dag_id='dag_a', start_date=datetime(2024, 1, 1), schedule_interval='@daily', catchup=False ) as dag: trigger_dag_b_task = TriggerDagRunOperator( task_id='trigger_dag_b', trigger_dag_id='dag_b', trigger_callable=conditionally_trigger, # 如需传递参数给dag_b,可添加conf参数 # conf={'source_dag': 'dag_a', 'run_time': '{{ ds }}'}, wait_for_completion=False ) # 关联dag_a中的前置任务到该触发任务 # your_previous_task >> trigger_dag_b_task
关键说明
DagModel.get_dagmodel()直接读取Airflow元数据库中的DAG状态,确保判断的是实时的暂停/开启状态。- 返回
None或False会让TriggerDagRunOperator跳过触发动作,不会在dag_b中创建待执行任务队列。 - 新增的XCom记录可以在Airflow UI的任务实例详情中查看,方便确认是否因为dag_b暂停而跳过了触发。
- 该方案兼容Airflow 2.x全版本,Airflow 1.x仅需调整少量导入语法即可适配。
内容的提问来源于stack exchange,提问作者Andrew Yar
相关产品推荐
相关产品推荐

