Skip to content
Open
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
8 changes: 6 additions & 2 deletions docs/source/contributor-guide/memory_management.md
Original file line number Diff line number Diff line change
Expand Up @@ -455,11 +455,15 @@ On Kubernetes, Spark sizes the executor pod from `ResourceProfile`:
```text
pod memory request = pod memory limit
= spark.executor.memory
+ spark.executor.memoryOverhead (default max(0.1 * executor.memory, 384 MiB))
+ spark.executor.memoryOverhead (default max(factor * executor.memory, 384 MiB))
+ spark.memory.offHeap.size
+ pyspark memory (Python applications only)
+ spark.executor.pyspark.memory (Python applications in cluster mode only)
```

The factor is `spark.executor.memoryOverheadFactor` or, when that is unset,
`spark.kubernetes.memoryOverheadFactor`. Both default to 0.1, but in cluster mode spark-submit sets
the Kubernetes factor to 0.4 for a PySpark or SparkR application that did not set it.

Both the request and the limit are set to this same value, so the pod's cgroup `memory.max` is a
hard ceiling on the sum of everything in the container. That cgroup counts, among other things:

Expand Down
24 changes: 14 additions & 10 deletions docs/source/user-guide/latest/tuning/memory.md
Original file line number Diff line number Diff line change
Expand Up @@ -99,11 +99,12 @@ already drawing on.
Work out what the executor already gets before choosing a value. When
`spark.executor.memoryOverhead` is unset, Spark derives the overhead as
`max(spark.executor.memoryOverheadFactor * spark.executor.memory, 384 MiB)`. The factor defaults to
`0.1`, except for PySpark and SparkR applications submitted to Kubernetes in cluster mode, where it
defaults to `0.4`. On Spark 4.0 and later the floor is configurable through
`spark.executor.minMemoryOverhead`. Setting `spark.executor.memoryOverhead` **replaces** the derived
value rather than adding to it, so a value below what is derived today shrinks the container instead
of growing it.
`0.1`. On Kubernetes, when `spark.executor.memoryOverheadFactor` is unset, Spark uses
`spark.kubernetes.memoryOverheadFactor` instead, which also defaults to `0.1`, except for PySpark
and SparkR applications submitted in cluster mode, where it defaults to `0.4`. On Spark 4.0 and
later the floor is configurable through `spark.executor.minMemoryOverhead`. Setting
`spark.executor.memoryOverhead` **replaces** the derived value rather than adding to it, so a value
below what is derived today shrinks the container instead of growing it.

For a small executor, `2g` is a reasonable starting point. A 4 GiB executor derives only 409 MiB, so
this is a real increase:
Expand All @@ -126,7 +127,9 @@ Raise the value further if executors are killed by the cluster manager (on Kuber
To measure how much Comet needs rather than guessing, see [Sizing the Overhead from the Memory Usage Log].

Note that on Kubernetes and YARN the overhead is added to the container size, so raising it reduces
how many executors fit on a node.
how many executors fit on a node. A standalone cluster ignores the overhead: its workers start
executors without a memory limit and count only `spark.executor.memory` against the memory they
offer.

[Sizing the Overhead from the Memory Usage Log]: #sizing-the-overhead-from-the-memory-usage-log

Expand Down Expand Up @@ -179,10 +182,11 @@ occupy until Spark hands it out, so a quiet log is not a sign that the overhead
size it from the most untracked memory as described above. The overhead also has to hold the JVM's
own non-heap memory, so by the time the warning appears the executor has likely outgrown its
container. It warns the first time this happens, and again each time it happens after dropping back
below. The overhead it uses is `spark.executor.memoryOverhead` if set, otherwise
`spark.executor.memoryOverheadFactor` of `spark.executor.memory` with a minimum of
`spark.executor.minMemoryOverhead`, as Spark sizes the default container. There is no warning in
local mode.
below. It derives the overhead the way YARN and Kubernetes size the executor's container, as
described in [Configuring Executor Memory Overhead]. For a PySpark application it also counts
`spark.executor.pyspark.memory` as part of the container when the cluster manager adds it, which
YARN does, and Kubernetes does in cluster mode. There is no warning in local mode, on a standalone
cluster, whose workers do not limit an executor's memory, or with any other cluster manager.

Look more closely before raising the overhead if the untracked memory keeps growing through a run
rather than levelling off: memory that is not being released will exhaust any overhead eventually.
Expand Down
73 changes: 56 additions & 17 deletions spark/src/main/scala/org/apache/comet/CometExecIterator.scala
Original file line number Diff line number Diff line change
Expand Up @@ -456,24 +456,51 @@ object CometExecIterator extends Logging {
}

/**
* The executor's memory overhead in bytes, sized the way Spark sizes the default resource
* profile's container: `spark.executor.memoryOverhead` if set, otherwise
* `spark.executor.memoryOverheadFactor` of `spark.executor.memory`, but at least
* `spark.executor.minMemoryOverhead`. None in local mode, where there is no container, or if
* the settings do not parse. An executor running a non-default resource profile may have a
* Whether the cluster manager for `master` runs each executor in a container whose memory it
* sizes from the memory overhead settings: YARN and Kubernetes. Local mode, local-cluster
* included, has no executor container, and a standalone worker starts its executors with no
* memory limit and never reads those settings. Comet does not know how any other cluster
* manager sizes its executors.
*/
def isContainerSizedFromOverhead(master: String): Boolean =
master == "yarn" || master.startsWith("k8s://")

/**
* The executor's memory overhead in bytes, sized the way YARN and Kubernetes size the default
* resource profile's container: `spark.executor.memoryOverhead` if set, otherwise a factor of
* `spark.executor.memory`, but at least `spark.executor.minMemoryOverhead` (a fixed 384 MiB
* before Spark 4.0 added that setting). The factor is `spark.executor.memoryOverheadFactor`,
* 0.1 by default, except that Kubernetes falls back to `spark.kubernetes.memoryOverheadFactor`
* when it is unset. In cluster mode spark-submit always sets that for the driver, to 0.4 for a
* PySpark or SparkR application that did not set it, and the executors get the driver's value.
* None where no container is sized from the overhead (see [[isContainerSizedFromOverhead]]), or
* if the settings do not parse. An executor running a non-default resource profile may have a
* different overhead.
*/
def executorMemoryOverhead(conf: SparkConf): Option[Long] = {
if (conf.get("spark.master", "").startsWith("local")) {
val master = conf.get("spark.master", "")
if (!isContainerSizedFromOverhead(master)) {
None
} else {
try {
val overheadMiB = conf.getOption("spark.executor.memoryOverhead") match {
case Some(_) => conf.getSizeAsMb("spark.executor.memoryOverhead")
case None =>
val executorMiB = conf.getSizeAsMb("spark.executor.memory", "1g")
val factor = conf.getDouble("spark.executor.memoryOverheadFactor", 0.1)
val minimumMiB = conf.getSizeAsMb("spark.executor.minMemoryOverhead", "384m")
val kubernetesFactor = if (master.startsWith("k8s://")) {
conf.getOption("spark.kubernetes.memoryOverheadFactor")
} else {
None
}
val factor = conf
.getOption("spark.executor.memoryOverheadFactor")
.orElse(kubernetesFactor)
.fold(0.1)(_.toDouble)
val minimumMiB = if (CometSparkSessionExtensions.isSpark40Plus) {
conf.getSizeAsMb("spark.executor.minMemoryOverhead", "384m")
} else {
384L
}
math.max((executorMiB * factor).toLong, minimumMiB)
}
Some(ByteUnit.MiB.toBytes(overheadMiB))
Expand All @@ -484,19 +511,30 @@ object CometExecIterator extends Logging {
}

/**
* The memory the executor's container has for native memory: `spark.memory.offHeap.size` plus
* the memory overhead; see [[executorMemoryOverhead]]. None, so that nothing is compared
* against it, in local mode, when off-heap memory is disabled (a testing-only mode in which
* Comet's reservations do not come from Spark's off-heap pool), or if the settings do not
* parse.
* The memory the executor's container has outside the JVM heap: `spark.memory.offHeap.size`,
* plus the memory overhead (see [[executorMemoryOverhead]]), plus
* `spark.executor.pyspark.memory` for an application that spark-submit marked as Python, with
* `spark.yarn.isPython` for YARN and with `spark.kubernetes.resource.type` in Kubernetes
* cluster mode. None, so that nothing is compared against it, where the overhead is None, when
* off-heap memory is disabled (a testing-only mode in which Comet's reservations do not come
* from Spark's off-heap pool), or if the settings do not parse.
*/
def nativeMemoryLimit(conf: SparkConf): Option[Long] = {
if (!CometSparkSessionExtensions.isOffHeapEnabled(conf)) {
None
} else {
executorMemoryOverhead(conf).flatMap { overhead =>
try {
Some(overhead + conf.getSizeAsBytes("spark.memory.offHeap.size", "0"))
val pythonApp = if (conf.get("spark.master", "") == "yarn") {
conf.getBoolean("spark.yarn.isPython", false)
} else {
conf.get("spark.kubernetes.resource.type", "") == "python"
}
val pysparkMiB =
if (pythonApp) conf.getSizeAsMb("spark.executor.pyspark.memory", "0") else 0L
Some(
overhead + conf.getSizeAsBytes("spark.memory.offHeap.size", "0") +
ByteUnit.MiB.toBytes(pysparkMiB))
} catch {
case NonFatal(_) => None
}
Expand Down Expand Up @@ -608,9 +646,10 @@ object CometExecIterator extends Logging {
s"Arrow) plus Spark's off-heap memory in use (${toMiB(sparkOffHeapUsed)}, including " +
s"Comet's reservations) is ${toMiB(footprint)}, more than the ${toMiB(limitBytes)} the " +
"executor's container has outside the JVM heap (spark.memory.offHeap.size plus the " +
"memory overhead), which also has to hold the JVM's own non-heap memory. The cluster " +
"manager may kill this executor for exceeding its container limit. Raise " +
s"spark.executor.memoryOverhead. ${CometConf.TUNING_GUIDE}.")
"memory overhead, and spark.executor.pyspark.memory for a PySpark application), which " +
"also has to hold the JVM's own non-heap memory. The cluster manager may kill this " +
"executor for exceeding its container limit. Raise spark.executor.memoryOverhead. " +
s"${CometConf.TUNING_GUIDE}.")
} else {
None
}
Expand Down
9 changes: 5 additions & 4 deletions spark/src/main/scala/org/apache/spark/Plugins.scala
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,7 @@ import org.apache.spark.internal.Logging
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}
import org.apache.comet.{COMET_VERSION, CometExecIterator, CometSparkSessionExtensions, NativeBase}
import org.apache.comet.{CometConf, ConfigEntry}
import org.apache.comet.CometConf.{COMET_ICEBERG_WRITE_REPORT_DIR, COMET_METRICS_ENABLED, COMET_ONHEAP_ENABLED}
import org.apache.comet.CometKryoRegistrator
Expand Down Expand Up @@ -172,10 +172,11 @@ 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")
// Only YARN and Kubernetes size executors from the overhead, not local mode or standalone
val sizedFromOverhead =
CometExecIterator.isContainerSizedFromOverhead(conf.get("spark.master", ""))

if (cometActive && !localMode && !isExecutorMemoryOverheadSet(conf)) {
if (cometActive && sizedFromOverhead && !isExecutorMemoryOverheadSet(conf)) {
logWarning(
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 " +
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -453,7 +453,9 @@ class CometExecIteratorLifecycleSuite extends CometTestBase {
== fourGiB)
}

test("the native memory limit is the off-heap size plus the memory overhead") {
private val kubernetesMaster = "k8s://https://kubernetes.default.svc:443"

test("the native memory limit is the container's memory outside the JVM heap") {
import CometExecIterator.nativeMemoryLimit
val mib = 1024L * 1024
val offHeap = new SparkConf(false)
Expand All @@ -464,9 +466,28 @@ class CometExecIteratorLifecycleSuite extends CometTestBase {
assert(nativeMemoryLimit(offHeap) == Some((4096 + 1638) * mib))
assert(nativeMemoryLimit(offHeap.clone.set("spark.memory.offHeap.enabled", "false")).isEmpty)
assert(nativeMemoryLimit(offHeap.clone.set("spark.master", "local[*]")).isEmpty)

// The container also holds spark.executor.pyspark.memory for an application that
// spark-submit marked as Python: with spark.yarn.isPython for YARN, and with
// spark.kubernetes.resource.type in Kubernetes cluster mode.
val pyspark = offHeap.clone.set("spark.executor.pyspark.memory", "2g")
assert(nativeMemoryLimit(pyspark) == Some((4096 + 1638) * mib))
assert(
nativeMemoryLimit(pyspark.clone.set("spark.yarn.isPython", "true"))
== Some((4096 + 1638 + 2048) * mib))
val kubernetes = pyspark.clone
.set("spark.master", kubernetesMaster)
.set("spark.kubernetes.memoryOverheadFactor", "0.1")
// Nothing sets the resource type in client mode, and Kubernetes then leaves it out.
assert(nativeMemoryLimit(kubernetes) == Some((4096 + 1638) * mib))
Seq("java" -> 0, "r" -> 0, "python" -> 2048).foreach { case (resourceType, pysparkMiB) =>
val conf = kubernetes.clone.set("spark.kubernetes.resource.type", resourceType)
assert(nativeMemoryLimit(conf) == Some((4096 + 1638 + pysparkMiB) * mib), resourceType)
}
}

test("the executor memory overhead is sized as Spark sizes the container") {
test("the executor memory overhead is sized as YARN sizes the container") {
import org.apache.comet.CometSparkSessionExtensions.isSpark40Plus
import CometExecIterator.executorMemoryOverhead
val mib = 1024L * 1024
def conf(settings: (String, String)*): SparkConf =
Expand All @@ -485,8 +506,80 @@ class CometExecIteratorLifecycleSuite extends CometTestBase {
conf(
"spark.executor.memory" -> "10g",
"spark.executor.memoryOverheadFactor" -> "0.25")) == Some(2560 * mib))
// Spark 4.0 made the 384 MiB minimum configurable; earlier versions ignore the setting.
assert(
executorMemoryOverhead(
conf("spark.executor.memory" -> "4g", "spark.executor.minMemoryOverhead" -> "1g")) ==
Some((if (isSpark40Plus) 1024 else 409) * mib))
// Only Kubernetes reads the Kubernetes factor.
assert(
executorMemoryOverhead(
conf(
"spark.executor.memory" -> "8g",
"spark.kubernetes.memoryOverheadFactor" -> "0.4")) == Some(819 * mib))
assert(executorMemoryOverhead(conf("spark.executor.memoryOverhead" -> "lots")).isEmpty)
assert(executorMemoryOverhead(new SparkConf(false).set("spark.master", "local[4]")).isEmpty)
}

test("the executor memory overhead is sized as Kubernetes sizes the pod") {
import CometExecIterator.executorMemoryOverhead
val mib = 1024L * 1024
def conf(settings: (String, String)*): SparkConf =
new SparkConf(false)
.set("spark.master", kubernetesMaster)
.set("spark.executor.memory", "8g")
.setAll(settings)

// In cluster mode spark-submit passes spark.kubernetes.memoryOverheadFactor on to the
// executors, 0.4 for a PySpark or SparkR application that did not set it, and Kubernetes
// uses it when spark.executor.memoryOverheadFactor is unset.
Seq("python", "r").foreach { resourceType =>
val pod = conf(
"spark.kubernetes.resource.type" -> resourceType,
"spark.kubernetes.memoryOverheadFactor" -> "0.4")
assert(executorMemoryOverhead(pod) == Some(3276 * mib), resourceType)
}
// A factor the application set is passed on in the same way.
assert(
executorMemoryOverhead(conf("spark.kubernetes.memoryOverheadFactor" -> "0.3")) ==
Some(2457 * mib))
// Nothing sets it in client mode, where it defaults to 0.1.
assert(executorMemoryOverhead(conf()) == Some(819 * mib))
assert(
executorMemoryOverhead(
conf(
"spark.kubernetes.memoryOverheadFactor" -> "0.4",
"spark.executor.memoryOverheadFactor" -> "0.2")) == Some(1638 * mib))
assert(
executorMemoryOverhead(
conf(
"spark.kubernetes.memoryOverheadFactor" -> "0.4",
"spark.executor.memoryOverhead" -> "1g")) == Some(1024 * mib))
assert(
executorMemoryOverhead(
conf(
"spark.executor.memory" -> "512m",
"spark.kubernetes.memoryOverheadFactor" -> "0.4")) == Some(384 * mib))
assert(
executorMemoryOverhead(conf("spark.kubernetes.memoryOverheadFactor" -> "lots")).isEmpty)
}

test("the executor memory overhead is unknown without a container sized from it") {
import CometExecIterator.executorMemoryOverhead
// Local mode has no executor container, and a standalone worker starts executors without
// a memory limit, never reading the overhead settings.
val masters =
Seq(
"local",
"local[4]",
"local-cluster[2,1,1024]",
"spark://host:7077",
"mesos://host:5050")
masters.foreach { master =>
val conf = new SparkConf(false)
.set("spark.master", master)
.set("spark.executor.memoryOverhead", "2g")
assert(executorMemoryOverhead(conf).isEmpty, master)
}
}

test("the memory usage log interval disables the log on a value it cannot use") {
Expand Down
4 changes: 2 additions & 2 deletions spark/src/test/scala/org/apache/spark/CometPluginsSuite.scala
Original file line number Diff line number Diff line change
Expand Up @@ -251,8 +251,8 @@ class CometPluginsMemoryOverheadWarningSuite extends CometTestBase {
assert(!warningsFor(conf).exists(_.contains(warning)))
}

test("does not warn in local mode") {
Seq("local", "local[4]", "local-cluster[2,1,1024]").foreach { master =>
test("does not warn in local mode or on a standalone cluster") {
Seq("local", "local[4]", "local-cluster[2,1,1024]", "spark://host:7077").foreach { master =>
assert(!warningsFor(cometConf(master)).exists(_.contains(warning)), master)
}
}
Expand Down
Loading