为何关闭DAG后SparkSubmitOperator启动5分钟后自行终止?
Airflow SparkSubmitOperator 任务被SIGTERM终止问题分析
问题描述
我创建了包含SparkSubmitOperator的DAG,操作流程为:开启DAG生成任务 → 关闭DAG避免生成新任务 → 手动运行任务。但任务启动5分钟后总会触发SIGTERM终止,该问题已重复出现5次以上;若保持DAG开启状态运行任务则一切正常。相关日志如下:
[2023-01-27 12:50:04,783] {local_task_job.py:188} WARNING - State of this instance has been externally set to None. Terminating instance. [2023-01-27 12:50:04,798] {process_utils.py:100} INFO - Sending Signals.SIGTERM to GPID 968 [2023-01-27 12:50:04,802] {taskinstance.py:1265} ERROR - Received SIGTERM. Terminating subprocesses. [2023-01-27 12:50:04,804] {spark_submit.py:657} INFO - Sending kill signal to spark-submit [2023-01-27 12:50:15,985] {spark_submit.py:674} INFO - YARN app killed with return code: 0 [2023-01-27 12:50:16,122] {taskinstance.py:1482} ERROR - Task failed with exception Traceback (most recent call last): File "/home/airflow/.local/lib/python3.6/site-packages/airflow/models/taskinstance.py", line 1138, in _run_raw_task self._prepare_and_execute_task_with_callbacks(context, task) File "/home/airflow/.local/lib/python3.6/site-packages/airflow/models/taskinstance.py", line 1311, in _prepare_and_execute_task_with_callbacks result = self._execute_task(context, task_copy) File "/home/airflow/.local/lib/python3.6/site-packages/airflow/models/taskinstance.py", line 1341, in _execute_task result = task_copy.execute(context=context) File "/home/airflow/.local/lib/python3.6/site-packages/airflow/providers/apache/spark/operators/spark_submit.py", line 183, in execute self._hook.submit(self._application) File "/home/airflow/.local/lib/python3.6/site-packages/airflow/providers/apache/spark/hooks/spark_submit.py", line 440, in submit self._process_spark_submit_log(iter(self._submit_sp.stdout)) # type: ignore File "/home/airflow/.local/lib/python3.6/site-packages/airflow/providers/apache/spark/hooks/spark_submit.py", line 494, in _process_spark_submit_log for line in itr: File "/home/airflow/.local/lib/python3.6/site-packages/airflow/models/taskinstance.py", line 1267, in signal_handler raise AirflowException("Task received SIGTERM signal") airflow.exceptions.AirflowException: Task received SIGTERM signal
原因分析
核心原因是Airflow Scheduler的任务清理机制:
- 当DAG被关闭后,Scheduler会将该DAG下的所有运行中任务标记为“不应继续执行”,并按照默认5分钟的检查周期(对应
scheduler_zombie_task_threshold配置值)发送SIGTERM信号终止任务。 - 手动启动的任务虽然初始执行,但Scheduler检测到其所属DAG处于关闭状态,就会触发清理逻辑,导致任务被终止。
- 而DAG开启时,Scheduler会判定任务为合法运行实例,不会执行清理操作,因此任务能正常完成。
解决方案
- 临时操作方案:手动运行任务前保持DAG开启,若需避免自动生成新任务,可将DAG的
max_active_runs设为1,或通过任务实例的“标记成功”/“清除”操作控制任务生命周期,而非直接关闭DAG。 - 长期配置方案:将DAG的
schedule_interval设置为None,这样DAG始终处于开启状态但不会自动调度生成任务,手动运行时也不会被Scheduler终止。 - 应急规避方案:若必须关闭DAG后运行任务,可修改Airflow配置中的
scheduler_zombie_task_threshold参数,延长清理周期,但此方法会影响Scheduler整体清理效率,不推荐常规使用。
内容的提问来源于stack exchange,提问作者Gumada Yaroslav
相关产品推荐
相关产品推荐

