From 229927061abfb534ff2febc7f1ab3e602337106a Mon Sep 17 00:00:00 2001 From: Angerszhuuuu Date: Mon, 3 Aug 2026 18:52:52 +0800 Subject: [PATCH 1/4] [SPARK-58513][CORE] Init UnifiedMemeoryManager should check driver config in driver , check executor config in executor --- .../scala/org/apache/spark/SparkEnv.scala | 5 ++- .../spark/memory/UnifiedMemoryManager.scala | 13 ++++--- .../memory/UnifiedMemoryManagerSuite.scala | 34 +++++++++++++++++-- 3 files changed, 44 insertions(+), 8 deletions(-) diff --git a/core/src/main/scala/org/apache/spark/SparkEnv.scala b/core/src/main/scala/org/apache/spark/SparkEnv.scala index ca48ee473eb0f..877a269943762 100644 --- a/core/src/main/scala/org/apache/spark/SparkEnv.scala +++ b/core/src/main/scala/org/apache/spark/SparkEnv.scala @@ -498,7 +498,10 @@ class SparkEnv ( } else { conf.clone.set(MEMORY_OFFHEAP_ENABLED, false).set(MEMORY_OFFHEAP_SIZE, 0L) } - _memoryManager = UnifiedMemoryManager(memoryManagerConf, numUsableCores) + _memoryManager = UnifiedMemoryManager( + memoryManagerConf, + numUsableCores, + isDriver = SparkContext.isDriver(executorId)) } } diff --git a/core/src/main/scala/org/apache/spark/memory/UnifiedMemoryManager.scala b/core/src/main/scala/org/apache/spark/memory/UnifiedMemoryManager.scala index 6b278c47f32f1..22e38a9051c2f 100644 --- a/core/src/main/scala/org/apache/spark/memory/UnifiedMemoryManager.scala +++ b/core/src/main/scala/org/apache/spark/memory/UnifiedMemoryManager.scala @@ -446,8 +446,11 @@ object UnifiedMemoryManager extends Logging { } } - def apply(conf: SparkConf, numCores: Int): UnifiedMemoryManager = { - val maxMemory = getMaxMemory(conf) + def apply( + conf: SparkConf, + numCores: Int, + isDriver: Boolean = true): UnifiedMemoryManager = { + val maxMemory = getMaxMemory(conf, isDriver) new UnifiedMemoryManager( conf, maxHeapMemory = maxMemory, @@ -459,12 +462,12 @@ object UnifiedMemoryManager extends Logging { /** * Return the total amount of memory shared between execution and storage, in bytes. */ - private def getMaxMemory(conf: SparkConf): Long = { + private def getMaxMemory(conf: SparkConf, isDriver: Boolean): Long = { val systemMemory = conf.get(TEST_MEMORY) val reservedMemory = conf.getLong(TEST_RESERVED_MEMORY.key, if (conf.contains(IS_TESTING)) 0 else RESERVED_SYSTEM_MEMORY_BYTES) val minSystemMemory = (reservedMemory * 1.5).ceil.toLong - if (systemMemory < minSystemMemory) { + if (isDriver && systemMemory < minSystemMemory) { throw new SparkIllegalArgumentException( errorClass = "INVALID_DRIVER_MEMORY", messageParameters = Map( @@ -473,7 +476,7 @@ object UnifiedMemoryManager extends Logging { "config" -> config.DRIVER_MEMORY.key)) } // SPARK-12759 Check executor memory to fail fast if memory is insufficient - if (conf.contains(config.EXECUTOR_MEMORY)) { + if (!isDriver && conf.contains(config.EXECUTOR_MEMORY)) { val executorMemory = conf.getSizeAsBytes(config.EXECUTOR_MEMORY.key) if (executorMemory < minSystemMemory) { throw new SparkIllegalArgumentException( diff --git a/core/src/test/scala/org/apache/spark/memory/UnifiedMemoryManagerSuite.scala b/core/src/test/scala/org/apache/spark/memory/UnifiedMemoryManagerSuite.scala index 9f0e622b1d515..dd8adc1f956fd 100644 --- a/core/src/test/scala/org/apache/spark/memory/UnifiedMemoryManagerSuite.scala +++ b/core/src/test/scala/org/apache/spark/memory/UnifiedMemoryManagerSuite.scala @@ -249,16 +249,46 @@ class UnifiedMemoryManagerSuite extends MemoryManagerSuite with PrivateMethodTes .set(TEST_MEMORY, systemMemory) .set(TEST_RESERVED_MEMORY, reservedMemory) - val mm = UnifiedMemoryManager(conf, numCores = 1) + val mm = UnifiedMemoryManager(conf, numCores = 1, isDriver = false) // Try using an executor memory that's too small val conf2 = conf.clone().set(EXECUTOR_MEMORY.key, (reservedMemory / 2).toString) val exception = intercept[IllegalArgumentException] { - UnifiedMemoryManager(conf2, numCores = 1) + UnifiedMemoryManager(conf2, numCores = 1, isDriver = false) } assert(exception.getMessage.contains("increase executor memory")) } + test("executor does not validate driver heap") { + val systemMemory = 400L * 1024 + val reservedMemory = 300L * 1024 + val memoryFraction = 0.8 + val conf = new SparkConf() + .set(MEMORY_FRACTION, memoryFraction) + .set(TEST_MEMORY, systemMemory) + .set(TEST_RESERVED_MEMORY, reservedMemory) + .set(EXECUTOR_MEMORY.key, (500L * 1024).toString) + + val mm = UnifiedMemoryManager(conf, numCores = 1, isDriver = false) + val expectedMaxMemory = ((systemMemory - reservedMemory) * memoryFraction).toLong + assert(mm.maxHeapMemory === expectedMaxMemory) + } + + test("driver does not validate executor memory") { + val systemMemory = 1024L * 1024 + val reservedMemory = 300L * 1024 + val memoryFraction = 0.8 + val conf = new SparkConf() + .set(MEMORY_FRACTION, memoryFraction) + .set(TEST_MEMORY, systemMemory) + .set(TEST_RESERVED_MEMORY, reservedMemory) + .set(EXECUTOR_MEMORY.key, (reservedMemory / 2).toString) + + val mm = UnifiedMemoryManager(conf, numCores = 1) + val expectedMaxMemory = ((systemMemory - reservedMemory) * memoryFraction).toLong + assert(mm.maxHeapMemory === expectedMaxMemory) + } + test("execution can evict cached blocks when there are multiple active tasks (SPARK-12155)") { val conf = new SparkConf() .set(MEMORY_FRACTION, 1.0) From 66d311dae122e2122dd08682a375713ac7908e16 Mon Sep 17 00:00:00 2001 From: Angerszhuuuu Date: Mon, 3 Aug 2026 18:58:17 +0800 Subject: [PATCH 2/4] Update UnifiedMemoryManagerSuite.scala --- .../org/apache/spark/memory/UnifiedMemoryManagerSuite.scala | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/core/src/test/scala/org/apache/spark/memory/UnifiedMemoryManagerSuite.scala b/core/src/test/scala/org/apache/spark/memory/UnifiedMemoryManagerSuite.scala index dd8adc1f956fd..d8a9e9e1e5dc5 100644 --- a/core/src/test/scala/org/apache/spark/memory/UnifiedMemoryManagerSuite.scala +++ b/core/src/test/scala/org/apache/spark/memory/UnifiedMemoryManagerSuite.scala @@ -259,7 +259,7 @@ class UnifiedMemoryManagerSuite extends MemoryManagerSuite with PrivateMethodTes assert(exception.getMessage.contains("increase executor memory")) } - test("executor does not validate driver heap") { + test("SPARK-58513: executor does not validate driver heap") { val systemMemory = 400L * 1024 val reservedMemory = 300L * 1024 val memoryFraction = 0.8 @@ -274,7 +274,7 @@ class UnifiedMemoryManagerSuite extends MemoryManagerSuite with PrivateMethodTes assert(mm.maxHeapMemory === expectedMaxMemory) } - test("driver does not validate executor memory") { + test("SPARK-58513: driver does not validate executor memory") { val systemMemory = 1024L * 1024 val reservedMemory = 300L * 1024 val memoryFraction = 0.8 From d72521b17bd915d0c90ca0136ad3c9c417458434 Mon Sep 17 00:00:00 2001 From: Angerszhuuuu Date: Tue, 4 Aug 2026 10:23:19 +0800 Subject: [PATCH 3/4] Update UnifiedMemoryManager.scala --- .../org/apache/spark/memory/UnifiedMemoryManager.scala | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/core/src/main/scala/org/apache/spark/memory/UnifiedMemoryManager.scala b/core/src/main/scala/org/apache/spark/memory/UnifiedMemoryManager.scala index 22e38a9051c2f..5e9cf86a6e2d0 100644 --- a/core/src/main/scala/org/apache/spark/memory/UnifiedMemoryManager.scala +++ b/core/src/main/scala/org/apache/spark/memory/UnifiedMemoryManager.scala @@ -446,10 +446,14 @@ object UnifiedMemoryManager extends Logging { } } + def apply(conf: SparkConf, numCores: Int): UnifiedMemoryManager = { + apply(conf, numCores, isDriver = true) + } + def apply( conf: SparkConf, numCores: Int, - isDriver: Boolean = true): UnifiedMemoryManager = { + isDriver: Boolean): UnifiedMemoryManager = { val maxMemory = getMaxMemory(conf, isDriver) new UnifiedMemoryManager( conf, From 3ac46439997ea04530b653c8f6a616b5e7afb07a Mon Sep 17 00:00:00 2001 From: Angerszhuuuu Date: Tue, 4 Aug 2026 11:50:48 +0800 Subject: [PATCH 4/4] update --- core/src/main/scala/org/apache/spark/SparkEnv.scala | 2 +- .../apache/spark/memory/UnifiedMemoryManager.scala | 12 +++++++----- .../spark/memory/UnifiedMemoryManagerSuite.scala | 8 ++++---- 3 files changed, 12 insertions(+), 10 deletions(-) diff --git a/core/src/main/scala/org/apache/spark/SparkEnv.scala b/core/src/main/scala/org/apache/spark/SparkEnv.scala index 877a269943762..e71b3c89cecc8 100644 --- a/core/src/main/scala/org/apache/spark/SparkEnv.scala +++ b/core/src/main/scala/org/apache/spark/SparkEnv.scala @@ -501,7 +501,7 @@ class SparkEnv ( _memoryManager = UnifiedMemoryManager( memoryManagerConf, numUsableCores, - isDriver = SparkContext.isDriver(executorId)) + isDriver = Some(SparkContext.isDriver(executorId))) } } diff --git a/core/src/main/scala/org/apache/spark/memory/UnifiedMemoryManager.scala b/core/src/main/scala/org/apache/spark/memory/UnifiedMemoryManager.scala index 5e9cf86a6e2d0..4e7ed6638290e 100644 --- a/core/src/main/scala/org/apache/spark/memory/UnifiedMemoryManager.scala +++ b/core/src/main/scala/org/apache/spark/memory/UnifiedMemoryManager.scala @@ -447,13 +447,13 @@ object UnifiedMemoryManager extends Logging { } def apply(conf: SparkConf, numCores: Int): UnifiedMemoryManager = { - apply(conf, numCores, isDriver = true) + apply(conf, numCores, isDriver = None) } def apply( conf: SparkConf, numCores: Int, - isDriver: Boolean): UnifiedMemoryManager = { + isDriver: Option[Boolean]): UnifiedMemoryManager = { val maxMemory = getMaxMemory(conf, isDriver) new UnifiedMemoryManager( conf, @@ -466,12 +466,14 @@ object UnifiedMemoryManager extends Logging { /** * Return the total amount of memory shared between execution and storage, in bytes. */ - private def getMaxMemory(conf: SparkConf, isDriver: Boolean): Long = { + private def getMaxMemory(conf: SparkConf, isDriver: Option[Boolean]): Long = { val systemMemory = conf.get(TEST_MEMORY) val reservedMemory = conf.getLong(TEST_RESERVED_MEMORY.key, if (conf.contains(IS_TESTING)) 0 else RESERVED_SYSTEM_MEMORY_BYTES) val minSystemMemory = (reservedMemory * 1.5).ceil.toLong - if (isDriver && systemMemory < minSystemMemory) { + val checkDriverMemory = isDriver.isEmpty || isDriver.contains(true) + val checkExecutorMemory = isDriver.isEmpty || isDriver.contains(false) + if (checkDriverMemory && systemMemory < minSystemMemory) { throw new SparkIllegalArgumentException( errorClass = "INVALID_DRIVER_MEMORY", messageParameters = Map( @@ -480,7 +482,7 @@ object UnifiedMemoryManager extends Logging { "config" -> config.DRIVER_MEMORY.key)) } // SPARK-12759 Check executor memory to fail fast if memory is insufficient - if (!isDriver && conf.contains(config.EXECUTOR_MEMORY)) { + if (checkExecutorMemory && conf.contains(config.EXECUTOR_MEMORY)) { val executorMemory = conf.getSizeAsBytes(config.EXECUTOR_MEMORY.key) if (executorMemory < minSystemMemory) { throw new SparkIllegalArgumentException( diff --git a/core/src/test/scala/org/apache/spark/memory/UnifiedMemoryManagerSuite.scala b/core/src/test/scala/org/apache/spark/memory/UnifiedMemoryManagerSuite.scala index d8a9e9e1e5dc5..5501eb124fffa 100644 --- a/core/src/test/scala/org/apache/spark/memory/UnifiedMemoryManagerSuite.scala +++ b/core/src/test/scala/org/apache/spark/memory/UnifiedMemoryManagerSuite.scala @@ -249,12 +249,12 @@ class UnifiedMemoryManagerSuite extends MemoryManagerSuite with PrivateMethodTes .set(TEST_MEMORY, systemMemory) .set(TEST_RESERVED_MEMORY, reservedMemory) - val mm = UnifiedMemoryManager(conf, numCores = 1, isDriver = false) + val mm = UnifiedMemoryManager(conf, numCores = 1) // Try using an executor memory that's too small val conf2 = conf.clone().set(EXECUTOR_MEMORY.key, (reservedMemory / 2).toString) val exception = intercept[IllegalArgumentException] { - UnifiedMemoryManager(conf2, numCores = 1, isDriver = false) + UnifiedMemoryManager(conf2, numCores = 1) } assert(exception.getMessage.contains("increase executor memory")) } @@ -269,7 +269,7 @@ class UnifiedMemoryManagerSuite extends MemoryManagerSuite with PrivateMethodTes .set(TEST_RESERVED_MEMORY, reservedMemory) .set(EXECUTOR_MEMORY.key, (500L * 1024).toString) - val mm = UnifiedMemoryManager(conf, numCores = 1, isDriver = false) + val mm = UnifiedMemoryManager(conf, numCores = 1, isDriver = Some(false)) val expectedMaxMemory = ((systemMemory - reservedMemory) * memoryFraction).toLong assert(mm.maxHeapMemory === expectedMaxMemory) } @@ -284,7 +284,7 @@ class UnifiedMemoryManagerSuite extends MemoryManagerSuite with PrivateMethodTes .set(TEST_RESERVED_MEMORY, reservedMemory) .set(EXECUTOR_MEMORY.key, (reservedMemory / 2).toString) - val mm = UnifiedMemoryManager(conf, numCores = 1) + val mm = UnifiedMemoryManager(conf, numCores = 1, isDriver = Some(true)) val expectedMaxMemory = ((systemMemory - reservedMemory) * memoryFraction).toLong assert(mm.maxHeapMemory === expectedMaxMemory) }