diff --git a/spark/src/main/scala/org/apache/spark/Plugins.scala b/spark/src/main/scala/org/apache/spark/Plugins.scala index d040208596..7891735b86 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 4673740daf..e65ae80ccd 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))) } }