From 355b5c67e12e5050c81230ba547d9e5b2dfeb68c Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Thu, 24 Sep 2026 09:44:58 -0600 Subject: [PATCH] fix: skip the executor memory overhead warning when a factor is set or in local mode The driver plugin warned that spark.executor.memoryOverhead was unset even when the application sized the overhead with spark.executor.memoryOverheadFactor, which the tuning guide recommends for large executors, and in local mode, where there is no executor container to size. Treat the overhead as set when spark.executor.memoryOverhead or spark.executor.memoryOverheadFactor is set, or when spark.kubernetes.memoryOverheadFactor differs from the value spark-submit passes to a cluster-mode driver on its own (0.4 for PySpark and SparkR applications, 0.1 otherwise), and skip local masters. The warning now names both settings. Closes #6189 --- .../main/scala/org/apache/spark/Plugins.scala | 41 ++++++++--- .../org/apache/spark/CometPluginsSuite.scala | 71 +++++++++++++++---- 2 files changed, 91 insertions(+), 21 deletions(-) diff --git a/spark/src/main/scala/org/apache/spark/Plugins.scala b/spark/src/main/scala/org/apache/spark/Plugins.scala index d0402085965..7891735b869 100644 --- a/spark/src/main/scala/org/apache/spark/Plugins.scala +++ b/spark/src/main/scala/org/apache/spark/Plugins.scala @@ -22,9 +22,11 @@ package org.apache.spark import java.{util => ju} import java.util.Collections +import scala.util.Try + import org.apache.spark.api.plugin.{DriverPlugin, ExecutorPlugin, PluginContext, SparkPlugin} import org.apache.spark.internal.Logging -import org.apache.spark.internal.config.EXECUTOR_MEMORY_OVERHEAD +import org.apache.spark.internal.config.{EXECUTOR_MEMORY_OVERHEAD, EXECUTOR_MEMORY_OVERHEAD_FACTOR} import org.apache.spark.sql.internal.StaticSQLConf import org.apache.comet.{COMET_VERSION, CometSparkSessionExtensions, NativeBase} @@ -167,16 +169,39 @@ object CometDriverPlugin extends Logging { val cometExecEnabled = getBooleanConf(conf, CometConf.COMET_EXEC_ENABLED) val cometShuffleEnabled = getBooleanConf(conf, CometConf.COMET_SHUFFLE_ENABLED) val cometActive = cometEnabled && (cometExecEnabled || cometShuffleEnabled) + // Local mode, local-cluster included, has no executor container to size + val localMode = conf.get("spark.master", "").startsWith("local") - if (cometActive && !conf.contains(EXECUTOR_MEMORY_OVERHEAD.key)) { + if (cometActive && !localMode && !isExecutorMemoryOverheadSet(conf)) { logWarning( - s"${EXECUTOR_MEMORY_OVERHEAD.key} is not set. Comet allocates outside the JVM heap, and " + - "the part of that which no memory pool tracks is not covered by " + - "spark.executor.memory or spark.memory.offHeap.size, so Spark's default overhead can " + - "leave the executor short and the cluster manager may kill it. Set " + - s"${EXECUTOR_MEMORY_OVERHEAD.key} before creating the SparkContext; it cannot be set " + - s"later. ${CometConf.TUNING_GUIDE}.") + s"Neither ${EXECUTOR_MEMORY_OVERHEAD.key} nor ${EXECUTOR_MEMORY_OVERHEAD_FACTOR.key} is " + + "set. Comet allocates outside the JVM heap, and the part of that which no memory pool " + + "tracks is not covered by spark.executor.memory or spark.memory.offHeap.size, so " + + "Spark's default overhead can leave the executor short and the cluster manager may " + + "kill it. Set one of them before creating the SparkContext; neither can be set later. " + + s"${CometConf.TUNING_GUIDE}.") + } + } + + // Whether the application sized the executor memory overhead itself, as an amount or as a + // factor of spark.executor.memory, rather than leaving it at Spark's default. + private def isExecutorMemoryOverheadSet(conf: SparkConf): Boolean = + conf.contains(EXECUTOR_MEMORY_OVERHEAD.key) || + conf.contains(EXECUTOR_MEMORY_OVERHEAD_FACTOR.key) || + isKubernetesMemoryOverheadFactorSet(conf) + + // Kubernetes falls back to spark.kubernetes.memoryOverheadFactor when + // spark.executor.memoryOverheadFactor is unset. In cluster mode spark-submit sets it for the + // driver even when the application did not, to 0.4 for PySpark and SparkR applications and 0.1 + // for the rest, so only a different value shows that the application set it. + private def isKubernetesMemoryOverheadFactorSet(conf: SparkConf): Boolean = { + val submitDefault = conf.get("spark.kubernetes.resource.type", "java") match { + case "python" | "r" => 0.4 + case _ => 0.1 } + conf + .getOption("spark.kubernetes.memoryOverheadFactor") + .exists(factor => !Try(factor.toDouble).toOption.contains(submitDefault)) } // spark.comet.exec.memoryPool.fraction was documented as holding back part of the off-heap pool diff --git a/spark/src/test/scala/org/apache/spark/CometPluginsSuite.scala b/spark/src/test/scala/org/apache/spark/CometPluginsSuite.scala index 4673740dafc..e65ae80ccd6 100644 --- a/spark/src/test/scala/org/apache/spark/CometPluginsSuite.scala +++ b/spark/src/test/scala/org/apache/spark/CometPluginsSuite.scala @@ -166,7 +166,8 @@ class CometPluginsDefaultSuite extends CometTestBase { class CometPluginsMemoryOverheadWarningSuite extends CometTestBase { - private val warning = "spark.executor.memoryOverhead is not set" + private val warning = + "Neither spark.executor.memoryOverhead nor spark.executor.memoryOverheadFactor is set" private def warningsFor(conf: SparkConf): Seq[String] = { // Logging derives the logger name by stripping the object's trailing '$' @@ -178,26 +179,70 @@ class CometPluginsMemoryOverheadWarningSuite extends CometTestBase { appender.loggingEvents.map(_.getMessage.getFormattedMessage).toSeq } + private val kubernetesMaster = "k8s://https://kubernetes.default.svc:443" + + private def cometConf(master: String = "yarn"): SparkConf = + new SparkConf() + .set("spark.master", master) + .set("spark.comet.enabled", "true") + .set("spark.comet.exec.enabled", "true") + test("warns when executor memory overhead is unset and Comet is active") { - val conf = new SparkConf() - conf.set("spark.comet.enabled", "true") - conf.set("spark.comet.exec.enabled", "true") - assert(warningsFor(conf).exists(_.contains(warning))) + assert(warningsFor(cometConf()).exists(_.contains(warning))) } test("does not warn when executor memory overhead is set") { - val conf = new SparkConf() - conf.set("spark.comet.enabled", "true") - conf.set("spark.comet.exec.enabled", "true") - conf.set("spark.executor.memoryOverhead", "2g") + val conf = cometConf().set("spark.executor.memoryOverhead", "2g") + assert(!warningsFor(conf).exists(_.contains(warning))) + } + + test("does not warn when the executor memory overhead factor is set") { + val conf = cometConf().set("spark.executor.memoryOverheadFactor", "0.2") + assert(!warningsFor(conf).exists(_.contains(warning))) + } + + test("does not warn when the Kubernetes memory overhead factor is set") { + // A factor other than the one spark-submit would have passed on, whatever the application type + Seq( + Map("spark.kubernetes.memoryOverheadFactor" -> "0.3"), + Map( + "spark.kubernetes.resource.type" -> "java", + "spark.kubernetes.memoryOverheadFactor" -> "0.4"), + Map( + "spark.kubernetes.resource.type" -> "python", + "spark.kubernetes.memoryOverheadFactor" -> "0.5")).foreach { settings => + val conf = cometConf(kubernetesMaster).setAll(settings) + assert(!warningsFor(conf).exists(_.contains(warning)), settings) + } + } + + test("warns on Kubernetes when the memory overhead factor is the one spark-submit passed on") { + // In cluster mode spark-submit sets spark.kubernetes.memoryOverheadFactor for the driver even + // when the application did not: 0.4 for PySpark and SparkR applications, 0.1 for the rest + Seq("java" -> "0.1", "python" -> "0.4", "r" -> "0.4").foreach { case (resourceType, factor) => + val conf = cometConf(kubernetesMaster) + .set("spark.kubernetes.resource.type", resourceType) + .set("spark.kubernetes.memoryOverheadFactor", factor) + assert(warningsFor(conf).exists(_.contains(warning)), resourceType) + } + } + + test("does not fail on a Kubernetes memory overhead factor that does not parse") { + // Spark rejects the value itself when it sizes the executor pods + val conf = cometConf(kubernetesMaster).set("spark.kubernetes.memoryOverheadFactor", "lots") assert(!warningsFor(conf).exists(_.contains(warning))) } + test("does not warn in local mode") { + Seq("local", "local[4]", "local-cluster[2,1,1024]").foreach { master => + assert(!warningsFor(cometConf(master)).exists(_.contains(warning)), master) + } + } + test("does not warn when Comet is not executing anything") { - val conf = new SparkConf() - conf.set("spark.comet.enabled", "true") - conf.set("spark.comet.exec.enabled", "false") - conf.set("spark.comet.shuffle.enabled", "false") + val conf = cometConf() + .set("spark.comet.exec.enabled", "false") + .set("spark.comet.shuffle.enabled", "false") assert(!warningsFor(conf).exists(_.contains(warning))) } }