From b0e81686af84f3231a4bc7d871df82d21c7de339 Mon Sep 17 00:00:00 2001 From: wforget <643348094@qq.com> Date: Wed, 12 Aug 2026 10:12:05 +0800 Subject: [PATCH 1/4] [MINOR] Refactor engineSavePath handling in SparkSQLEngine to use Option for better safety --- .../kyuubi/engine/spark/SparkSQLEngine.scala | 26 ++++++++++++------- 1 file changed, 17 insertions(+), 9 deletions(-) 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..b384d6f1ac5 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: Option[Path] = None override def initialize(conf: KyuubiConf): Unit = { val listener = new SparkSQLEngineListener(this) @@ -91,9 +91,13 @@ case class SparkSQLEngine(spark: SparkSession) extends Serverable("SparkSQLEngin } if (backendService.sessionManager.getConf.get(OPERATION_RESULT_SAVE_TO_FILE)) { - val fs = engineSavePath.getFileSystem(spark.sparkContext.hadoopConfiguration) - fs.mkdirs(engineSavePath) - fs.deleteOnExit(engineSavePath) + engineSavePath = Some(backendService.sessionManager + .asInstanceOf[SparkSQLSessionManager].getEngineResultSavePath()) + engineSavePath.foreach { path => + val fs = path.getFileSystem(spark.sparkContext.hadoopConfiguration) + fs.mkdirs(path) + fs.deleteOnExit(path) + } } } @@ -111,12 +115,16 @@ case class SparkSQLEngine(spark: SparkSession) extends Serverable("SparkSQLEngin Duration(60, TimeUnit.SECONDS)) }) try { - val fs = engineSavePath.getFileSystem(spark.sparkContext.hadoopConfiguration) - if (fs.exists(engineSavePath)) { - fs.delete(engineSavePath, true) + engineSavePath.foreach { path => + val fs = path.getFileSystem(spark.sparkContext.hadoopConfiguration) + if (fs.exists(path)) { + fs.delete(path, true) + } } } catch { - case e: Throwable => error(s"Error cleaning engine result save path: $engineSavePath", e) + case e: Throwable => + error(s"Error cleaning engine result save path:" + + s" ${engineSavePath.map(_.toString).getOrElse("")}", e) } } From 9682c77814b2ff898c126ffece5b7238851fb65f Mon Sep 17 00:00:00 2001 From: Cheng Pan Date: Wed, 12 Aug 2026 10:56:20 +0800 Subject: [PATCH 2/4] Apply suggestion --- .../scala/org/apache/kyuubi/engine/spark/SparkSQLEngine.scala | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) 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 b384d6f1ac5..02e41b6f1c3 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 @@ -123,8 +123,8 @@ case class SparkSQLEngine(spark: SparkSession) extends Serverable("SparkSQLEngin } } catch { case e: Throwable => - error(s"Error cleaning engine result save path:" + - s" ${engineSavePath.map(_.toString).getOrElse("")}", e) + error("Error cleaning engine result save path: " + + engineSavePath.map(_.toString).getOrElse(""), e) } } From 6315834842db84f5a84f2d995f86f29da65277b1 Mon Sep 17 00:00:00 2001 From: wforget <643348094@qq.com> Date: Wed, 12 Aug 2026 15:30:44 +0800 Subject: [PATCH 3/4] address comment --- .../kyuubi/engine/spark/SparkSQLEngine.scala | 30 +++++++------------ 1 file changed, 11 insertions(+), 19 deletions(-) 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 02e41b6f1c3..3a388372052 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,7 +26,6 @@ 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 @@ -59,7 +58,8 @@ case class SparkSQLEngine(spark: SparkSession) extends Serverable("SparkSQLEngin @volatile private var lifetimeTerminatingChecker: Option[ScheduledExecutorService] = None @volatile private var stopEngineExec: Option[ThreadPoolExecutor] = None - @volatile private var engineSavePath: Option[Path] = None + private val engineSavePath = + backendService.sessionManager.asInstanceOf[SparkSQLSessionManager].getEngineResultSavePath() override def initialize(conf: KyuubiConf): Unit = { val listener = new SparkSQLEngineListener(this) @@ -90,15 +90,11 @@ case class SparkSQLEngine(spark: SparkSession) extends Serverable("SparkSQLEngin startFastFailChecker(maxInitTimeout) } - if (backendService.sessionManager.getConf.get(OPERATION_RESULT_SAVE_TO_FILE)) { - engineSavePath = Some(backendService.sessionManager - .asInstanceOf[SparkSQLSessionManager].getEngineResultSavePath()) - engineSavePath.foreach { path => - val fs = path.getFileSystem(spark.sparkContext.hadoopConfiguration) - fs.mkdirs(path) - fs.deleteOnExit(path) - } - } + // Due to the fact that session-level enabling of `kyuubi.operation.result.saveToFile.enabled` + // is allowed, we always create and clean up this directory + val fs = engineSavePath.getFileSystem(spark.sparkContext.hadoopConfiguration) + fs.mkdirs(engineSavePath) + fs.deleteOnExit(engineSavePath) } override def stop(): Unit = if (shutdown.compareAndSet(false, true)) { @@ -115,16 +111,12 @@ case class SparkSQLEngine(spark: SparkSession) extends Serverable("SparkSQLEngin Duration(60, TimeUnit.SECONDS)) }) try { - engineSavePath.foreach { path => - val fs = path.getFileSystem(spark.sparkContext.hadoopConfiguration) - if (fs.exists(path)) { - fs.delete(path, true) - } + val fs = engineSavePath.getFileSystem(spark.sparkContext.hadoopConfiguration) + if (fs.exists(engineSavePath)) { + fs.delete(engineSavePath, true) } } catch { - case e: Throwable => - error("Error cleaning engine result save path: " + - engineSavePath.map(_.toString).getOrElse(""), e) + case e: Throwable => error(s"Error cleaning engine result save path: $engineSavePath", e) } } From 639ad06c0bd69f37d414d5e95b2b5a8e611f5dc3 Mon Sep 17 00:00:00 2001 From: wforget <643348094@qq.com> Date: Wed, 12 Aug 2026 16:12:14 +0800 Subject: [PATCH 4/4] fix tests --- .../kyuubi/engine/spark/SparkSQLEngine.scala | 20 +++++++++++-------- 1 file changed, 12 insertions(+), 8 deletions(-) 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 3a388372052..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 val engineSavePath = - backendService.sessionManager.asInstanceOf[SparkSQLSessionManager].getEngineResultSavePath() + @volatile private var engineSavePath: Path = _ override def initialize(conf: KyuubiConf): Unit = { val listener = new SparkSQLEngineListener(this) @@ -92,6 +92,8 @@ case class SparkSQLEngine(spark: SparkSession) extends Serverable("SparkSQLEngin // 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) @@ -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) } }