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

如何为依赖TriggerDagRunOperator的DAG实现截止时间调度与触发规避?

实现非严格依赖的DAG调度方案

针对你提出的需求,有两种可行的实现思路,以下是具体操作步骤:

一、核心需求拆解

要实现的逻辑是:

  1. 第二个DAG支持两种启动方式:被第一个DAG触发,或到达指定截止时间自动启动
  2. 第一个DAG触发前需检查第二个DAG是否已启动,避免重复触发

二、方案一:双触发机制+前置状态检查(推荐)

这种方案更轻量,无需让DAG长期处于等待状态,适合大多数场景。

1. 配置第二个DAG(自动截止启动+防重复)

给第二个DAG设置定时调度(截止时间),同时允许外部触发,并限制同一时间仅运行一个实例:

from airflow import DAG
from airflow.operators.dummy import DummyOperator
from datetime import datetime

default_args = {
    'owner': 'airflow',
    'start_date': datetime(2024, 1, 1),
    'retries': 0
}

# 假设截止时间为每天10:00
with DAG(
    'second_dag',
    default_args=default_args,
    schedule_interval='0 10 * * *',  # 每天10点自动触发
    max_active_runs=1,  # 禁止同一DAG同时运行多个实例
    catchup=False
) as dag:
    start = DummyOperator(task_id='start')
    # 此处添加你的业务任务
    end = DummyOperator(task_id='end')
    
    start >> end

2. 第一个DAG添加前置检查逻辑

在触发第二个DAG之前,先检查目标DAG是否已有运行/待运行实例,再决定是否触发:

from airflow import DAG
from airflow.operators.python import BranchPythonOperator
from airflow.operators.trigger_dagrun import TriggerDagRunOperator
from airflow.operators.dummy import DummyOperator
from airflow.models import DagRun
from datetime import datetime

default_args = {
    'owner': 'airflow',
    'start_date': datetime(2024, 1, 1),
    'retries': 1
}

def check_second_dag_status(**context):
    # 获取当前DAG的执行日期,确保检查同周期的目标DAG实例
    exec_date = context['execution_date']
    # 查找目标DAG已启动/运行/成功的实例
    existing_runs = DagRun.find(
        dag_id='second_dag',
        execution_date=exec_date,
        state=['running', 'queued', 'scheduled', 'success']
    )
    # 已有实例则跳过触发,否则执行触发
    return 'skip_trigger' if existing_runs else 'trigger_second_dag'

with DAG(
    'first_dag',
    default_args=default_args,
    schedule_interval='0 8 * * *',  # 示例:每天8点启动
    catchup=False
) as dag:
    # 检查目标DAG状态的分支任务
    check_status = BranchPythonOperator(
        task_id='check_second_dag_status',
        python_callable=check_second_dag_status,
        provide_context=True
    )
    
    # 触发第二个DAG的任务
    trigger_task = TriggerDagRunOperator(
        task_id='trigger_second_dag',
        trigger_dag_id='second_dag',
        execution_date='{{ execution_date }}',  # 传递相同执行日期,便于匹配
        wait_for_completion=False
    )
    
    # 跳过触发的占位任务
    skip_trigger = DummyOperator(task_id='skip_trigger')
    
    # 后续业务任务(无论触发与否都继续执行)
    continue_process = DummyOperator(
        task_id='continue_processing',
        trigger_rule='none_failed_min_one_success'
    )
    
    # 任务流配置
    check_status >> [trigger_task, skip_trigger] >> continue_process

三、方案二:双传感器等待(适合需明确触发来源的场景)

让第二个DAG提前启动,同时监听两个条件:第一个DAG完成,或到达截止时间,满足任一条件则执行后续任务。

from airflow import DAG
from airflow.sensors.time_sensor import TimeSensor
from airflow.sensors.external_task import ExternalTaskSensor
from airflow.operators.dummy import DummyOperator
from datetime import datetime

default_args = {
    'owner': 'airflow',
    'start_date': datetime(2024, 1, 1),
    'retries': 0
}

with DAG(
    'second_dag',
    default_args=default_args,
    schedule_interval='0 8 * * *',  # 和第一个DAG同一时间启动
    max_active_runs=1,
    catchup=False
) as dag:
    # 等待截止时间(每天10:00)
    wait_deadline = TimeSensor(
        task_id='wait_until_deadline',
        target_time=datetime.strptime('10:00', '%H:%M').time(),
        poke_interval=60  # 每分钟检查一次
    )
    
    # 等待第一个DAG的最后一个任务完成
    wait_first_dag = ExternalTaskSensor(
        task_id='wait_for_first_dag',
        external_dag_id='first_dag',
        external_task_id='end_of_first_dag',  # 替换为第一个DAG的最后一个任务ID
        execution_date='{{ execution_date }}',
        poke_interval=60,
        timeout=7200  # 超时时间:2小时后不再等待第一个DAG
    )
    
    # 后续业务任务(任一条件满足即执行)
    start_task = DummyOperator(
        task_id='start_task',
        trigger_rule='none_failed_min_one_success'
    )
    
    # 任务流配置
    [wait_deadline, wait_first_dag] >> start_task

四、关键注意点

  • 两种方案都需要给第二个DAG设置max_active_runs=1,防止重复执行
  • 方案一中的日期匹配逻辑可根据实际业务调整,比如若执行周期不同,可放宽日期匹配条件
  • 方案二中的timeout参数需根据第一个DAG的最长执行时长合理设置

内容的提问来源于stack exchange,提问作者Azaleum

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 04:35:13