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

为何关闭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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 18:45:59