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

如何实现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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 03:16:26