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

Apache Airflow中如何按条件推迟/延迟手动触发的DAG运行?

在Apache Airflow中延迟/推迟外部触发DAG运行的标准方案

我有一个通过REST API由外部组织手动触发的Apache Airflow DAG,想知道有没有标准机制可以根据特定条件(比如检查该组织是否还有其他DAG在运行)来推迟或延迟DAG运行?

我期望实现的逻辑示例

from airflow import DAG
from airflow.operators.python_operator import PythonOperator
from airflow.utils.dates import days_ago

default_args = {
    'owner': 'airflow',
    'start_date': days_ago(1),
}


def main_task(**kwargs):
   if organization_has_running_dag() :
      postpone_dag(minutes=10)
   do_main_job() 


dag = DAG(
    'example_dag',
    default_args=default_args,
    description='An example DAG to postpone runs',
)

main_task = PythonOperator(
    task_id='main_task',
    provide_context=True,
    python_callable=check_org_run,
    dag=dag,
)

main_task

我曾考虑的方案(不推荐)

用while循环加time.sleep()等待:

def main_task(**kwargs):
   while organization_has_running_dag() :
      time.sleep(5*60)
   do_main_job()

但这种方式会持续占用worker资源,且中断后无法恢复等待进度,我希望找到更标准的、能重新调度DAG而非在任务中等待的方案。


标准解决方案

1. 使用PythonSensor(官方推荐)

Airflow的PythonSensor是专门用于等待条件满足的组件,它会定期检查条件,直到满足或超时,不会持续占用worker资源:

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.sensors.python import PythonSensor
from airflow.utils.dates import days_ago
from airflow.utils.state import State
from airflow.models import DagRun
import pendulum

default_args = {
    'owner': 'airflow',
    'start_date': days_ago(1),
}

def check_organization_idle(**kwargs):
    # 从DAG运行配置中获取组织ID
    organization_id = kwargs['dag_run'].conf.get('organization_id')
    # 查询该组织的其他DAG是否有运行中的实例
    running_runs = DagRun.find(
        dag_id=['org_dag_a', 'org_dag_b'],  # 替换为该组织的其他DAG ID
        state=State.RUNNING,
        conf_contains={'organization_id': organization_id}
    )
    # 没有运行中的DAG则返回True,传感器结束等待
    return len(running_runs) == 0

def main_task(**kwargs):
    # 执行你的核心业务逻辑
    print(f"为组织{kwargs['dag_run'].conf.get('organization_id')}执行主任务")

with DAG(
    'sensor_based_delayed_dag',
    default_args=default_args,
    description='基于传感器延迟执行的组织DAG',
    catchup=False
) as dag:
    wait_for_org_idle = PythonSensor(
        task_id='wait_for_organization_idle',
        python_callable=check_organization_idle,
        poke_interval=600,  # 每10分钟检查一次条件
        timeout=7200,  # 最长等待2小时,超时则任务失败
        provide_context=True
    )

    main_task = PythonOperator(
        task_id='main_task',
        python_callable=main_task,
        provide_context=True
    )

    wait_for_org_idle >> main_task
  • 优势:完全符合Airflow调度模型,等待期间任务处于UP_FOR_RETRY状态,不占用worker;支持超时控制,避免无限等待。

2. 结合ShortCircuitOperator与重试机制

通过前置检查任务,条件不满足时触发任务重试,利用Airflow的重试延迟实现调度:

from airflow import DAG
from airflow.operators.python import PythonOperator, ShortCircuitOperator
from airflow.utils.dates import days_ago
from airflow.utils.state import State
from airflow.models import DagRun
import pendulum

default_args = {
    'owner': 'airflow',
    'start_date': days_ago(1),
    'retries': 12,  # 最多重试12次
    'retry_delay': pendulum.duration(minutes=10)  # 每次重试间隔10分钟
}

def check_organization_runs(**kwargs):
    organization_id = kwargs['dag_run'].conf.get('organization_id')
    running_runs = DagRun.find(
        dag_id=['org_dag_a', 'org_dag_b'],
        state=State.RUNNING,
        conf_contains={'organization_id': organization_id}
    )
    # 条件满足则返回True,执行后续任务;否则返回False,触发重试
    return len(running_runs) == 0

def main_task(**kwargs):
    print(f"为组织{kwargs['dag_run'].conf.get('organization_id')}执行主任务")

with DAG(
    'retry_based_delayed_dag',
    default_args=default_args,
    description='基于重试机制延迟的组织DAG',
    catchup=False
) as dag:
    check_task = ShortCircuitOperator(
        task_id='check_organization_running_dags',
        python_callable=check_organization_runs,
        provide_context=True
    )

    main_task = PythonOperator(
        task_id='main_task',
        python_callable=main_task,
        provide_context=True
    )

    check_task >> main_task
  • 原理:ShortCircuitOperator返回False时会跳过后续任务,但重试配置会让任务在指定延迟后重新执行,直到条件满足或达到最大重试次数。

3. 触发延迟的新DAG运行

当条件不满足时,触发一个延迟的新DAG实例,然后终止当前运行:

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.utils.dates import days_ago
from airflow.utils.state import State
from airflow.models import DagRun
from airflow.operators.trigger_dagrun import TriggerDagRunOperator
import pendulum

default_args = {
    'owner': 'airflow',
    'start_date': days_ago(1),
}

def check_and_reschedule(**kwargs):
    organization_id = kwargs['dag_run'].conf.get('organization_id')
    running_runs = DagRun.find(
        dag_id=['org_dag_a', 'org_dag_b'],
        state=State.RUNNING,
        conf_contains={'organization_id': organization_id}
    )
    if len(running_runs) > 0:
        # 触发10分钟后的新DAG运行,传递原配置
        trigger = TriggerDagRunOperator(
            task_id='trigger_delayed_run',
            trigger_dag_id=kwargs['dag'].dag_id,
            execution_date=pendulum.now().add(minutes=10),
            conf=kwargs['dag_run'].conf,
            wait_for_completion=False,
            dag=kwargs['dag']
        )
        trigger.execute(context=kwargs)
        # 标记当前DAG运行为成功,避免占用资源
        kwargs['dag_run'].set_state(State.SUCCESS)
        return
    # 条件满足,执行主任务
    print(f"为组织{organization_id}执行主任务")

with DAG(
    'rescheduled_org_dag',
    default_args=default_args,
    description='支持重新调度的组织DAG',
    catchup=False
) as dag:
    check_and_run_task = PythonOperator(
        task_id='check_and_reschedule',
        python_callable=check_and_reschedule,
        provide_context=True
    )

    check_and_run_task
  • 注意:需确保DAG允许手动触发,且catchup=False避免历史任务干扰。

方案对比

你之前考虑的time.sleep()方案会持续占用worker资源,等待期间worker无法处理其他任务,且一旦worker重启或中断,等待进度会丢失,不适合生产环境。上述三种标准方案均基于Airflow原生机制,资源利用率更高,稳定性更强。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 10:20:08