diff --git a/externals/kyuubi-spark-sql-engine/src/main/scala/org/apache/kyuubi/engine/spark/SparkSQLEngine.scala b/externals/kyuubi-spark-sql-engine/src/main/scala/org/apache/kyuubi/engine/spark/SparkSQLEngine.scala index 8f157698c91..b684608d4c4 100644 --- a/externals/kyuubi-spark-sql-engine/src/main/scala/org/apache/kyuubi/engine/spark/SparkSQLEngine.scala +++ b/externals/kyuubi-spark-sql-engine/src/main/scala/org/apache/kyuubi/engine/spark/SparkSQLEngine.scala @@ -26,6 +26,7 @@ import scala.concurrent.duration.Duration import scala.util.control.NonFatal import com.google.common.annotations.VisibleForTesting +import org.apache.hadoop.fs.Path import org.apache.spark.{ui, SparkConf} import org.apache.spark.kyuubi.{SparkContextHelper, SparkSQLEngineEventListener, SparkSQLEngineListener} import org.apache.spark.kyuubi.SparkUtilsHelper.getLocalDir @@ -58,8 +59,7 @@ case class SparkSQLEngine(spark: SparkSession) extends Serverable("SparkSQLEngin @volatile private var lifetimeTerminatingChecker: Option[ScheduledExecutorService] = None @volatile private var stopEngineExec: Option[ThreadPoolExecutor] = None - private lazy val engineSavePath = - backendService.sessionManager.asInstanceOf[SparkSQLSessionManager].getEngineResultSavePath() + @volatile private var engineSavePath: Path = _ override def initialize(conf: KyuubiConf): Unit = { val listener = new SparkSQLEngineListener(this) @@ -90,11 +90,13 @@ case class SparkSQLEngine(spark: SparkSession) extends Serverable("SparkSQLEngin startFastFailChecker(maxInitTimeout) } - if (backendService.sessionManager.getConf.get(OPERATION_RESULT_SAVE_TO_FILE)) { - val fs = engineSavePath.getFileSystem(spark.sparkContext.hadoopConfiguration) - fs.mkdirs(engineSavePath) - fs.deleteOnExit(engineSavePath) - } + // Due to the fact that session-level enabling of `kyuubi.operation.result.saveToFile.enabled` + // is allowed, we always create and clean up this directory + engineSavePath = + backendService.sessionManager.asInstanceOf[SparkSQLSessionManager].getEngineResultSavePath() + val fs = engineSavePath.getFileSystem(spark.sparkContext.hadoopConfiguration) + fs.mkdirs(engineSavePath) + fs.deleteOnExit(engineSavePath) } override def stop(): Unit = if (shutdown.compareAndSet(false, true)) { @@ -110,13 +112,15 @@ case class SparkSQLEngine(spark: SparkSession) extends Serverable("SparkSQLEngin exec, Duration(60, TimeUnit.SECONDS)) }) - try { - val fs = engineSavePath.getFileSystem(spark.sparkContext.hadoopConfiguration) - if (fs.exists(engineSavePath)) { - fs.delete(engineSavePath, true) + if (engineSavePath != null) { + try { + val fs = engineSavePath.getFileSystem(spark.sparkContext.hadoopConfiguration) + if (fs.exists(engineSavePath)) { + fs.delete(engineSavePath, true) + } + } catch { + case e: Throwable => error(s"Error cleaning engine result save path: $engineSavePath", e) } - } catch { - case e: Throwable => error(s"Error cleaning engine result save path: $engineSavePath", e) } }