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();