andygrove opened a new pull request, #6375:
URL: https://github.com/apache/datafusion-comet/pull/6375

   ## 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 #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 #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.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to