Skip to content

fix: derive the native memory limit the way YARN and Kubernetes size executor containers - #6375

Open
andygrove wants to merge 1 commit into
apache:mainfrom
andygrove:memory-log-k8s-overhead-factor
Open

andygrove wants to merge 1 commit into
apache:mainfrom
andygrove:memory-log-k8s-overhead-factor

Conversation

@andygrove

Copy link
Copy Markdown
Member

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.executorMemoryOverhead derived the overhead part of that limit from spark.executor.memoryOverhead, else spark.executor.memoryOverheadFactor (default 0.1) of spark.executor.memory. That is how YARN sizes the container, but it is not the whole story:

  • On Kubernetes, BasicExecutorFeatureStep falls back to spark.kubernetes.memoryOverheadFactor when spark.executor.memoryOverheadFactor is unset. In cluster mode, BasicDriverFeatureStep.getAdditionalPodSystemProperties always 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.
  • YARN (for an application with spark.yarn.isPython) and Kubernetes (for resource type python, which is only set in cluster mode) add spark.executor.pyspark.memory to the container. The limit left it out.
  • A standalone worker starts an executor with -Xmx from spark.executor.memory, never reads the overhead settings and sets no memory limit. The log still warned there, and its advice to raise spark.executor.memoryOverhead does 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, BasicDriverFeatureStep and the standalone Master/ExecutorRunner at 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.minMemoryOverhead only 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?

  • New 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.executorMemoryOverhead now 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 is spark.executor.memoryOverheadFactor, falling back to spark.kubernetes.memoryOverheadFactor (Kubernetes only) and then to 0.1. The minimum is spark.executor.minMemoryOverhead on 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.isKubernetesMemoryOverheadFactorSet does need those defaults, to tell a user-set factor from spark-submit's, so that logic stays in the plugin.
  • CometExecIterator.nativeMemoryLimit adds spark.executor.pyspark.memory when 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.warnIfExecutorMemoryOverheadUnset uses the same check, so it no longer warns on standalone clusters.
  • User guide (tuning/memory.md): explain Kubernetes' fallback to spark.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 in CometExecIteratorLifecycleSuite:

  • the executor memory overhead is sized as YARN sizes the container covers 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 heap covers PySpark memory being added on YARN with spark.yarn.isPython and on Kubernetes with resource type python, and not for Java, R or Kubernetes client mode.
  • In CometPluginsMemoryOverheadWarningSuite, does not warn in local mode or on a standalone cluster now also checks spark://host:7077.

I wrote the tests first. On the unchanged code, three of the four CometExecIteratorLifecycleSuite tests failed:

  • Kubernetes: Some(858783744) did not equal Some(3435134976), that is 819 MiB instead of 3276 MiB.
  • Standalone: a 2 GiB overhead instead of none.
  • Native memory limit: 5734 MiB instead of 7782 MiB, because the PySpark memory was missing.

The plugin's standalone case also failed without the Plugins.scala change, because the warning was logged. With the fix, those four tests and all eight CometPluginsMemoryOverheadWarningSuite tests pass on the default Spark 4.1 profile. The syntactic scalafix check and a scalastyle-inclusive test-compile also pass.

…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.
@github-actions github-actions Bot added the bug Something isn't working label Sep 29, 2026

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 isContainerSizedFromOverhead between 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.

@rich7420 rich7420 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

bug Something isn't working

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Native memory usage log underestimates the executor overhead for PySpark and SparkR on Kubernetes, and warns on standalone clusters

3 participants