andygrove commented on code in PR #6066:
URL: https://github.com/apache/datafusion-comet/pull/6066#discussion_r4065213540
##########
spark/src/main/scala/org/apache/comet/CometExecIterator.scala:
##########
@@ -392,37 +391,26 @@ object CometExecIterator extends Logging {
}
def getMemoryConfig(conf: SparkConf): MemoryConfig = {
- val numCores = numDriverOrExecutorCores(conf)
- val coresPerTask = conf.get("spark.task.cpus", "1").toInt
// there are different paths for on-heap vs off-heap mode
val offHeapMode = CometSparkSessionExtensions.isOffHeapEnabled(conf)
if (offHeapMode) {
// in off-heap mode, Comet uses unified memory management to share
off-heap memory with Spark
val offHeapSize =
ByteUnit.MiB.toBytes(conf.getSizeAsMb("spark.memory.offHeap.size"))
val memoryFraction = CometConf.COMET_OFFHEAP_MEMORY_POOL_FRACTION.get()
val memoryLimit = (offHeapSize * memoryFraction).toLong
- val memoryLimitPerTask = (memoryLimit.toDouble * coresPerTask /
numCores).toLong
val memoryPoolType = COMET_OFFHEAP_MEMORY_POOL_TYPE.get()
logDebug(
s"memoryPoolType=$memoryPoolType, " +
s"offHeapSize=${toMB(offHeapSize)}, " +
s"memoryFraction=$memoryFraction, " +
- s"memoryLimit=${toMB(memoryLimit)}, " +
- s"memoryLimitPerTask=${toMB(memoryLimitPerTask)}")
- MemoryConfig(offHeapMode, memoryPoolType = memoryPoolType, memoryLimit,
memoryLimitPerTask)
+ s"memoryLimit=${toMB(memoryLimit)}")
+ MemoryConfig(offHeapMode, memoryPoolType, memoryLimit)
} else {
- // we'll use the built-in memory pool from DF, and initializes with
`memory_limit`
- // and `memory_fraction` below.
- val memoryLimit =
CometSparkSessionExtensions.getCometMemoryOverhead(conf)
- // example 16GB maxMemory * 16 cores with 4 cores per task results
- // in memory_limit_per_task = 16 GB * 4 / 16 = 16 GB / 4 = 4GB
- val memoryLimitPerTask = (memoryLimit.toDouble * coresPerTask /
numCores).toLong
- val memoryPoolType = COMET_ONHEAP_MEMORY_POOL_TYPE.get()
- logDebug(
- s"memoryPoolType=$memoryPoolType, " +
- s"memoryLimit=${toMB(memoryLimit)}, " +
- s"memoryLimitPerTask=${toMB(memoryLimitPerTask)}")
- MemoryConfig(offHeapMode, memoryPoolType = memoryPoolType, memoryLimit,
memoryLimitPerTask)
+ // On-heap mode exists only so that the Spark SQL tests can run against
Comet without
+ // changing Spark's memory configuration, and native memory cannot be
charged to Spark's
+ // on-heap pool, so nothing is accounted. See the memory management
contributor guide.
+ logDebug("on-heap mode: native memory is unbounded and unaccounted")
+ MemoryConfig(offHeapMode, memoryPoolType = "unbounded", memoryLimit = 0)
Review Comment:
@sunchao Hmm .. technically, this is correct, but this is a test-only config
intended for running Spark SQL test suite and it is disabled by default. Our
docs already say that no user should use on-heap mode. I'd like to push back
against this review issue. WDYT?
--
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]