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
相关产品推荐
相关产品推荐

