This is an automated email from the ASF dual-hosted git repository.

xtsong pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/flink.git


The following commit(s) were added to refs/heads/master by this push:
     new 9cc5ab9caf3 [FLINK-30081][runtime] Reset unused options like 
jvm-overhead.min/max for local executor
9cc5ab9caf3 is described below

commit 9cc5ab9caf368ef336599e7d48f679c8c9750f49
Author: Mingliang Liu <[email protected]>
AuthorDate: Sat Apr 13 01:24:59 2024 -0700

    [FLINK-30081][runtime] Reset unused options like jvm-overhead.min/max for 
local executor
    
    This closes #22547
---
 .../taskexecutor/TaskExecutorResourceUtils.java    | 39 +++++++++-------------
 .../TaskExecutorResourceUtilsTest.java             | 19 +++++++++++
 2 files changed, 35 insertions(+), 23 deletions(-)

diff --git 
a/flink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/TaskExecutorResourceUtils.java
 
b/flink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/TaskExecutorResourceUtils.java
index c1d80510c7f..1461afb6012 100644
--- 
a/flink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/TaskExecutorResourceUtils.java
+++ 
b/flink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/TaskExecutorResourceUtils.java
@@ -49,7 +49,8 @@ public class TaskExecutorResourceUtils {
                     TaskManagerOptions.NETWORK_MEMORY_MAX,
                     TaskManagerOptions.MANAGED_MEMORY_SIZE);
 
-    private static final List<ConfigOption<?>> UNUSED_CONFIG_OPTIONS =
+    @VisibleForTesting
+    static final List<ConfigOption<?>> UNUSED_CONFIG_OPTIONS =
             Arrays.asList(
                     TaskManagerOptions.TOTAL_PROCESS_MEMORY,
                     TaskManagerOptions.TOTAL_FLINK_MEMORY,
@@ -189,7 +190,8 @@ public class TaskExecutorResourceUtils {
     }
 
     public static Configuration adjustForLocalExecution(Configuration config) {
-        UNUSED_CONFIG_OPTIONS.forEach(option -> 
warnOptionHasNoEffectIfSet(config, option));
+        UNUSED_CONFIG_OPTIONS.forEach(
+                option -> warnAndRemoveOptionHasNoEffectIfSet(config, option));
 
         setConfigOptionToPassedMaxIfNotSet(
                 config, TaskManagerOptions.CPU_CORES, 
LOCAL_EXECUTION_CPU_CORES);
@@ -201,24 +203,20 @@ public class TaskExecutorResourceUtils {
         adjustNetworkMemoryForLocalExecution(config);
         setConfigOptionToDefaultIfNotSet(
                 config, TaskManagerOptions.MANAGED_MEMORY_SIZE, 
DEFAULT_MANAGED_MEMORY_SIZE);
-        silentlySetConfigOptionIfNotSet(
-                config,
+
+        // Set valid default values for unused config options which should 
have been removed.
+        config.set(
                 TaskManagerOptions.FRAMEWORK_HEAP_MEMORY,
                 TaskManagerOptions.FRAMEWORK_HEAP_MEMORY.defaultValue());
-        silentlySetConfigOptionIfNotSet(
-                config,
+        config.set(
                 TaskManagerOptions.FRAMEWORK_OFF_HEAP_MEMORY,
                 TaskManagerOptions.FRAMEWORK_OFF_HEAP_MEMORY.defaultValue());
-        silentlySetConfigOptionIfNotSet(
-                config,
-                TaskManagerOptions.JVM_METASPACE,
-                TaskManagerOptions.JVM_METASPACE.defaultValue());
-        silentlySetConfigOptionIfNotSet(
-                config,
+        config.set(
+                TaskManagerOptions.JVM_METASPACE, 
TaskManagerOptions.JVM_METASPACE.defaultValue());
+        config.set(
                 TaskManagerOptions.JVM_OVERHEAD_MAX,
                 TaskManagerOptions.JVM_OVERHEAD_MAX.defaultValue());
-        silentlySetConfigOptionIfNotSet(
-                config,
+        config.set(
                 TaskManagerOptions.JVM_OVERHEAD_MIN,
                 TaskManagerOptions.JVM_OVERHEAD_MAX.defaultValue());
 
@@ -244,20 +242,15 @@ public class TaskExecutorResourceUtils {
                 config, TaskManagerOptions.NETWORK_MEMORY_MAX, 
DEFAULT_SHUFFLE_MEMORY_SIZE);
     }
 
-    private static void warnOptionHasNoEffectIfSet(Configuration config, 
ConfigOption<?> option) {
+    private static void warnAndRemoveOptionHasNoEffectIfSet(
+            Configuration config, ConfigOption<?> option) {
         if (config.contains(option)) {
             LOG.warn(
                     "The resource configuration option {} is set but it will 
have no effect for local execution, "
                             + "only the following options matter for the 
resource configuration: {}",
                     option,
-                    UNUSED_CONFIG_OPTIONS);
-        }
-    }
-
-    private static <T> void silentlySetConfigOptionIfNotSet(
-            Configuration config, ConfigOption<T> option, T value) {
-        if (!config.contains(option)) {
-            config.set(option, value);
+                    CONFIG_OPTIONS);
+            config.removeConfig(option);
         }
     }
 
diff --git 
a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskExecutorResourceUtilsTest.java
 
b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskExecutorResourceUtilsTest.java
index 30eb8cec90d..ce74981ea25 100644
--- 
a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskExecutorResourceUtilsTest.java
+++ 
b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskExecutorResourceUtilsTest.java
@@ -123,6 +123,25 @@ class TaskExecutorResourceUtilsTest {
                 .isEqualTo(networkMemorySize);
     }
 
+    @Test
+    void testUnusedOptionsAreIgnoredForLocalExecution() {
+        Configuration configuration = new Configuration();
+        configuration.set(TaskManagerOptions.TOTAL_PROCESS_MEMORY, 
MemorySize.ofMebiBytes(2024));
+        configuration.set(TaskManagerOptions.TOTAL_FLINK_MEMORY, 
MemorySize.ofMebiBytes(2024));
+        configuration.set(TaskManagerOptions.FRAMEWORK_HEAP_MEMORY, 
MemorySize.ofMebiBytes(2024));
+        configuration.set(
+                TaskManagerOptions.FRAMEWORK_OFF_HEAP_MEMORY, 
MemorySize.ofMebiBytes(2024));
+        configuration.set(TaskManagerOptions.JVM_METASPACE, 
MemorySize.ofMebiBytes(2024));
+        configuration.set(TaskManagerOptions.JVM_OVERHEAD_MIN, 
MemorySize.ofMebiBytes(2024));
+        configuration.set(TaskManagerOptions.JVM_OVERHEAD_MAX, 
MemorySize.ofMebiBytes(2024));
+        configuration.set(TaskManagerOptions.JVM_OVERHEAD_FRACTION, 2024.0f);
+
+        TaskExecutorResourceUtils.adjustForLocalExecution(configuration);
+
+        assertThat(configuration)
+                
.isEqualTo(TaskExecutorResourceUtils.adjustForLocalExecution(new 
Configuration()));
+    }
+
     @Test
     void testCalculateTotalFlinkMemoryWithAllFactorsBeingSet() {
         Configuration config = new Configuration();

Reply via email to