Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
41 changes: 33 additions & 8 deletions spark/src/main/scala/org/apache/spark/Plugins.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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}
Expand Down Expand Up @@ -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
Expand Down
71 changes: 58 additions & 13 deletions spark/src/test/scala/org/apache/spark/CometPluginsSuite.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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 '$'
Expand All @@ -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)))
}
}
Expand Down
Loading