Airflow 1.10.2中dagrun_timeout设置无效,如何解决?
为什么Airflow 1.10.2的dagrun_timeout没阻止长任务运行?
你没选错参数,但dagrun_timeout的作用和你预期的完全不同!它不是用来限制单个任务的运行时长,而是管控整个DAG Run从启动到结束的总耗时。
具体来说,Airflow只会在整个DAG Run完成后,才会检查总时长是否超过dagrun_timeout设定的值。如果你的DAG是串行执行(就像你代码里的would_succeed >> would_succeed_with_delay),第二个任务单独跑了2分钟,但整个DAG Run的总时长是从启动到第二个任务结束的时间——但Airflow不会在任务运行过程中主动中断它,只会在任务结束后判断总时长是否超标。这就导致即使单个任务超时,只要DAG Run完成时没触发总时长超时(或者说Airflow没在运行中终止任务),就会被标记为成功。
如果你想限制单个任务的运行时间,应该给每个任务单独设置execution_timeout参数,而不是在DAG层面设置dagrun_timeout。
修改后的代码示例如下:
import airflow from airflow import DAG from airflow.operators.python_operator import PythonOperator from datetime import timedelta import time args = { 'owner': 'me', 'start_date': airflow.utils.dates.days_ago(2), 'provide_context': True, } dag = DAG( 'test_timeout', schedule_interval=None, default_args=args, dagrun_timeout=timedelta(seconds=20), # 管控整个DAG Run的总时长 ) def this_passes(**kwargs): return def this_passes_with_delay(**kwargs): time.sleep(120) return would_succeed = PythonOperator( task_id='would_succeed', dag=dag, python_callable=this_passes, # email=to, # 假设to是你预先定义的变量 ) would_succeed_with_delay = PythonOperator( task_id='would_succeed_with_delay', dag=dag, python_callable=this_passes_with_delay, # email=to, execution_timeout=timedelta(seconds=20), # 给该任务单独设置超时限制 ) would_succeed >> would_succeed_with_delay
这样设置后,当would_succeed_with_delay任务运行超过20秒时,Airflow会主动终止该任务,并将其标记为失败,整个DAG Run也会随之失败。
补充说明:在Airflow 1.10.2版本中,dagrun_timeout的生效逻辑确实是“事后检查”,不会中断正在运行的任务。如果需要强制终止超时任务,任务级别的execution_timeout才是正确的选择。
内容的提问来源于stack exchange,提问作者Babak Tourani
相关产品推荐
相关产品推荐

