Conversation
…executor containers The executor's memory usage log compared Comet's native footprint against spark.memory.offHeap.size plus an overhead derived from spark.executor.memoryOverheadFactor alone. Kubernetes falls back to spark.kubernetes.memoryOverheadFactor when that is unset, and spark-submit sets it to 0.4 for PySpark and SparkR applications in cluster mode, so the limit was too low for them. The limit also left out spark.executor.pyspark.memory, which YARN and Kubernetes add to a Python application's container, and the log warned on standalone clusters, whose workers never read the overhead settings and do not limit an executor's memory. Only YARN and Kubernetes now get a limit, the overhead follows their precedence, and the driver's startup warning about an unset overhead skips standalone clusters through the same check.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: The warning underestimated Kubernetes/PySpark container allowances and incorrectly warned on standalone clusters.
- Design approach: Match YARN and Kubernetes sizing rules and share
isContainerSizedFromOverheadbetween the driver and executor warnings. - Correctness / compatibility analysis: Checked configuration precedence, Python application detection, executor configuration propagation, and minimum overhead against Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0 sources. No introduced P1/P2 issues found within this review.
- Key design decisions: Explicit overhead takes precedence. Kubernetes alone uses its fallback factor. The configurable minimum applies from Spark 4.0. Unsupported cluster managers receive no inferred limit.
- Implementation sketch: Update the sizing helpers and warning gate, extend the existing tests, and update both memory guides. The shared predicate keeps the implementation small and consistent.
- Behavioral changes worth calling out: Changes affect diagnostics, preserving INFO usage logging and existing allocation, reservation, and FFI ownership behavior. Configuration resolution remains outside per-batch execution.
- Suggested improvements: None meeting the P1/P2 reporting threshold.
Reviewed the entire six-file diff from 61810e2839d54241336be44de8765370e183a5e2 to 9e26f96053511e2d3328f50044b7d45ca1253ac4. The PR remains non-draft. The discussion snapshot contained no reviews, comments, or threads.
Routed skills: review-comet-pr, review-comet-memory-pr, and review-comet-ffi-pr.
Exact-head CI: Comet CI passed. Logs confirm the changed sizing tests and all eight overhead-warning tests passed. Linux builds, lint, Rust tests, Comet suites, and TPC-H/TPC-DS checks succeeded. Spark SQL, Iceberg, and macOS jobs were skipped.
Validation: A disposable Scala harness compiled the exact-head sizing methods and passed 192 configuration comparisons against Spark 4.1.3's actual ResourceProfile, plus unsupported-master and malformed-configuration checks. No local full Comet/native rebuild or live YARN/Kubernetes deployment was performed. Other supported versions were checked through source comparison.
Which issue does this PR close?
Closes #6188.
Rationale for this change
The executor's native memory usage log warns when Comet's untracked native memory plus Spark's off-heap memory in use exceeds what the executor's container has outside the JVM heap.
CometExecIterator.executorMemoryOverheadderived the overhead part of that limit fromspark.executor.memoryOverhead, elsespark.executor.memoryOverheadFactor(default 0.1) ofspark.executor.memory. That is how YARN sizes the container, but it is not the whole story:BasicExecutorFeatureStepfalls back tospark.kubernetes.memoryOverheadFactorwhenspark.executor.memoryOverheadFactoris unset. In cluster mode,BasicDriverFeatureStep.getAdditionalPodSystemPropertiesalways writes that factor into the driver's configuration, 0.4 for a PySpark or SparkR application that did not set it, and executors receive the driver's configuration. For the issue's example (a PySpark application with an 8 GiB executor), the pod has 3276 MiB of overhead but the log compared against 819 MiB, so it could warn that the cluster manager may kill an executor that was still well inside its pod. A Kubernetes factor that the user set was ignored in the same way.spark.yarn.isPython) and Kubernetes (for resource typepython, which is only set in cluster mode) addspark.executor.pyspark.memoryto the container. The limit left it out.-Xmxfromspark.executor.memory, never reads the overhead settings and sets no memory limit. The log still warned there, and its advice to raisespark.executor.memoryOverheaddoes nothing on standalone. The driver's startup warning from fix: skip the executor memory overhead warning when a factor is set or in local mode #6198 has the same problem, and that PR left it for this issue.I checked the precedence in
ResourceProfile.getResourcesForClusterManager,YarnAllocator,BasicExecutorFeatureStep,BasicDriverFeatureStepand the standaloneMaster/ExecutorRunnerat v3.4.3, v3.5.9, v4.0.4, v4.1.3 and v4.2.0. It is the same in all of them except the minimum overhead:spark.executor.minMemoryOverheadonly exists from 4.0, and 3.4 and 3.5 use a fixed 384 MiB.On standalone clusters I chose to skip the warning rather than reword it. There is no container limit to compare against (the executor shares the host's memory with the worker's other executors), and the only remedy the warning could suggest has no effect there. The INFO usage line is still logged.
What changes are included in this PR?
CometExecIterator.isContainerSizedFromOverhead(master), which returns true only for YARN and Kubernetes. Local mode, standalone and any other cluster manager get no limit, so they get no warning.CometExecIterator.executorMemoryOverheadnow sizes the overhead the way the cluster manager does. It uses the explicit overhead if one is set. Otherwise it takes a factor of the executor memory, with a minimum. The factor isspark.executor.memoryOverheadFactor, falling back tospark.kubernetes.memoryOverheadFactor(Kubernetes only) and then to 0.1. The minimum isspark.executor.minMemoryOverheadon Spark 4.0 and later, and 384 MiB before that. It doesn't need spark-submit's defaults, because executors inherit the factor the driver was given.CometDriverPlugin.isKubernetesMemoryOverheadFactorSetdoes need those defaults, to tell a user-set factor from spark-submit's, so that logic stays in the plugin.CometExecIterator.nativeMemoryLimitaddsspark.executor.pyspark.memorywhen the cluster manager adds it to the container, and the warning's description of the limit now mentions it. The warning text is only changed after the lines that fix: count pool overcommit as untracked memory in the native memory usage log #6271 edits.CometDriverPlugin.warnIfExecutorMemoryOverheadUnsetuses the same check, so it no longer warns on standalone clusters.tuning/memory.md): explain Kubernetes' fallback tospark.kubernetes.memoryOverheadFactor, note that a standalone cluster ignores the overhead, and replace "as Spark sizes the default container" with how the warning derives its limit and where it is not logged. Contributor guide (memory_management.md): the Kubernetes pod formula no longer hard-codes the 0.1 factor.How are these changes tested?
The overhead resolution is still a pure function of the
SparkConf, master included. It is unit-tested inCometExecIteratorLifecycleSuite:the executor memory overhead is sized as YARN sizes the containercovers an explicit overhead (with and without a unit), the executor factor, the 384 MiB minimum,spark.executor.minMemoryOverhead(honored on 4.0 and later, ignored before), YARN ignoring the Kubernetes factor, and a value that does not parse.the executor memory overhead is sized as Kubernetes sizes the pod(new) covers the 0.4 factor passed on for PySpark and SparkR (the issue's 3276 MiB), a Kubernetes factor the user set, the 0.1 default in client mode, the executor factor and an explicit overhead taking precedence, the minimum, and a factor that does not parse.the executor memory overhead is unknown without a container sized from it(new) covers local, local-cluster, standalone and Mesos masters.the native memory limit is the container's memory outside the JVM heapcovers PySpark memory being added on YARN withspark.yarn.isPythonand on Kubernetes with resource typepython, and not for Java, R or Kubernetes client mode.CometPluginsMemoryOverheadWarningSuite,does not warn in local mode or on a standalone clusternow also checksspark://host:7077.I wrote the tests first. On the unchanged code, three of the four
CometExecIteratorLifecycleSuitetests failed:Some(858783744) did not equal Some(3435134976), that is 819 MiB instead of 3276 MiB.The plugin's standalone case also failed without the
Plugins.scalachange, because the warning was logged. With the fix, those four tests and all eightCometPluginsMemoryOverheadWarningSuitetests pass on the default Spark 4.1 profile. The syntactic scalafix check and a scalastyle-inclusivetest-compilealso pass.