Airflow新调度周期触发时是否可终止仍在运行的上一周期DAG实例?
Airflow 调度自动终止过期运行实例实现方案
你需求的功能有两种实现路径,优先推荐无代码侵入的内置参数配置方案,无需修改现有Sensor逻辑即可实现:
方案1:官方内置参数配置(推荐,无业务逻辑侵入)
仅需在DAG初始化时配置3个核心参数即可满足需求:
max_active_runs = 1:限制该DAG同时最多只能有1个活跃运行实例,从根源上避免多个实例并行拉取重复数据dagrun_timeout = timedelta(hours=23):设置单DAG实例的最大运行时长,超过该时长的实例会被Airflow自动强制终止。这里设置为23小时,比每日调度的24小时间隔短1小时,保证下一次调度触发前,卡死的旧实例已经被清理catchup = False:关闭历史回溯跑数逻辑,避免调度器积压大量未运行的历史实例占用资源
同时建议给检查报告生成的Sensor单独配置timeout参数,比如timeout=21600(即6小时超时),避免单个Sensor卡死长期占用任务槽位。
配置示例代码:
from datetime import timedelta from airflow import DAG default_args = { 'owner': 'airflow', # 可将任务级全局默认配置放在此处 } dag = DAG( dag_id='your_cooperation_report_dag', default_args=default_args, schedule_interval='0 2 * * *', # 你的每日调度时间 max_active_runs=1, dagrun_timeout=timedelta(hours=23), catchup=False, )
方案2:新实例启动主动终止旧实例(适合需要即时清理的特殊场景)
如果你需要新调度触发后立即终止旧实例,不想等待dagrun_timeout到期,可以在DAG的第一个任务中添加少量逻辑,主动查询并终止旧的活跃实例,同样无需修改现有Sensor逻辑:
- 逻辑步骤:
- 在DAG的首个PythonOperator任务中,通过Airflow内置的
DagRun模型查询当前DAG所有状态为running的实例 - 排除当前正在执行的实例本身
- 把其余运行中的旧实例设置为
failed状态,Airflow会自动终止该实例下的所有运行中任务
- 在DAG的首个PythonOperator任务中,通过Airflow内置的
代码示例:
from airflow.models import DagRun from airflow.operators.python import PythonOperator def kill_old_dag_runs(**context): dag_id = context['dag'].dag_id current_run_id = context['dag_run'].run_id # 查询当前DAG所有活跃实例 active_runs = DagRun.find(dag_id=dag_id, state='running') for run in active_runs: if run.run_id != current_run_id: # 终止旧实例 run.set_state('failed') # 放在DAG的最开头作为第一个执行的任务 kill_old_task = PythonOperator( task_id='kill_old_runs', python_callable=kill_old_dag_runs, provide_context=True, dag=dag, )
内容的提问来源于stack exchange,提问作者Tevett Goad
相关产品推荐
相关产品推荐

