This is an automated email from the ASF dual-hosted git repository. gyfora pushed a commit to branch release-1.16 in repository https://gitbox.apache.org/repos/asf/flink-kubernetes-operator.git
commit 17a29a75ab593c682e0c54883d9a5ee96dcbf233 Author: Dennis-Mircea Ciupitu <[email protected]> AuthorDate: Thu Aug 27 15:23:38 2026 +0300 [FLINK-40474] Autoscaler parallelism overrides serialized in the operator's YAML dialect break job submission for Flink 1.x jobs (#1195) --- .../autoscaler-parallelism-pipeline.svg | 2 +- .../jdbc/state/JdbcAutoScalerStateStore.java | 2 +- .../flink/autoscaler/tuning/ConfigChanges.java | 3 +- .../flink/autoscaler/tuning/ConfigChangesTest.java | 61 ++++++++++++++++++++++ .../autoscaler/KubernetesScalingRealizer.java | 43 ++++++++++++--- .../state/KubernetesAutoScalerStateStore.java | 2 +- .../autoscaler/KubernetesScalingRealizerTest.java | 56 ++++++++++++++++++++ .../state/KubernetesAutoScalerStateStoreTest.java | 19 +++++++ 8 files changed, 177 insertions(+), 11 deletions(-) diff --git a/docs/static/img/custom-resource/autoscaler-parallelism-pipeline.svg b/docs/static/img/custom-resource/autoscaler-parallelism-pipeline.svg index 3667740f..cc4dfa05 100644 --- a/docs/static/img/custom-resource/autoscaler-parallelism-pipeline.svg +++ b/docs/static/img/custom-resource/autoscaler-parallelism-pipeline.svg @@ -1,4 +1,4 @@ <?xml version="1.0" encoding="UTF-8"?> <!-- Do not edit this file with editors other than draw.io --> <!DOCTYPE svg PUBLIC "-//W3C//DTD SVG 1.1//EN" "http://www.w3.org/Graphics/SVG/1.1/DTD/svg11.dtd"> -<svg xmlns="http://www.w3.org/2000/svg" style="background: transparent; background-color: transparent; color-scheme: light;" xmlns:xlink="http://www.w3.org/1999/xlink" version="1.1" width="854px" height="425px" viewBox="0 0 854 425" content="<mxfile host="Electron" agent="Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) draw.io/29.0.3 Chrome/140.0.7339.249 Electron/38.7.0 Safari/537.36" version="29.0.3" pages="10 [...] \ No newline at end of file +<svg xmlns="http://www.w3.org/2000/svg" style="background: transparent; background-color: transparent; color-scheme: light;" xmlns:xlink="http://www.w3.org/1999/xlink" version="1.1" width="854px" height="425px" viewBox="0 0 854 425" content="<mxfile host="Electron" agent="Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) draw.io/29.0.3 Chrome/140.0.7339.249 Electron/38.7.0 Safari/537.36" version="29.0.3" pages="10 [...] \ No newline at end of file diff --git a/flink-autoscaler-plugin-jdbc/src/main/java/org/apache/flink/autoscaler/jdbc/state/JdbcAutoScalerStateStore.java b/flink-autoscaler-plugin-jdbc/src/main/java/org/apache/flink/autoscaler/jdbc/state/JdbcAutoScalerStateStore.java index 6bc48878..0e382c94 100644 --- a/flink-autoscaler-plugin-jdbc/src/main/java/org/apache/flink/autoscaler/jdbc/state/JdbcAutoScalerStateStore.java +++ b/flink-autoscaler-plugin-jdbc/src/main/java/org/apache/flink/autoscaler/jdbc/state/JdbcAutoScalerStateStore.java @@ -297,7 +297,7 @@ public class JdbcAutoScalerStateStore<KEY, Context extends JobAutoScalerContext< } private static String serializeParallelismOverrides(Map<String, String> overrides) { - return ConfigurationUtils.convertValue(overrides, String.class); + return ConfigurationUtils.convertValue(overrides, String.class, false); } private static Map<String, String> deserializeParallelismOverrides(String overrides) { diff --git a/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/tuning/ConfigChanges.java b/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/tuning/ConfigChanges.java index f3e92e32..55fcba1d 100644 --- a/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/tuning/ConfigChanges.java +++ b/flink-autoscaler/src/main/java/org/apache/flink/autoscaler/tuning/ConfigChanges.java @@ -39,7 +39,8 @@ public class ConfigChanges { @Getter private final Set<String> removals = new HashSet<>(); public <T> ConfigChanges addOverride(ConfigOption<T> configOption, T value) { - overrides.put(configOption.key(), ConfigurationUtils.convertValue(value, String.class)); + overrides.put( + configOption.key(), ConfigurationUtils.convertValue(value, String.class, false)); return this; } diff --git a/flink-autoscaler/src/test/java/org/apache/flink/autoscaler/tuning/ConfigChangesTest.java b/flink-autoscaler/src/test/java/org/apache/flink/autoscaler/tuning/ConfigChangesTest.java new file mode 100644 index 00000000..91c7bd62 --- /dev/null +++ b/flink-autoscaler/src/test/java/org/apache/flink/autoscaler/tuning/ConfigChangesTest.java @@ -0,0 +1,61 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.autoscaler.tuning; + +import org.apache.flink.configuration.GlobalConfiguration; +import org.apache.flink.configuration.MemorySize; +import org.apache.flink.configuration.PipelineOptions; +import org.apache.flink.configuration.TaskManagerOptions; + +import org.junit.jupiter.api.Test; + +import java.util.LinkedHashMap; +import java.util.List; + +import static org.assertj.core.api.Assertions.assertThat; + +/** Tests for {@link ConfigChanges}. */ +class ConfigChangesTest { + + @Test + void testOverridesSerializedInLegacyDialect() { + boolean previousDialect = GlobalConfiguration.isStandardYaml(); + GlobalConfiguration.setStandardYaml(true); + try { + var globalJobParameters = new LinkedHashMap<String, String>(); + globalJobParameters.put("k1", "v1"); + globalJobParameters.put("k2", "v2"); + var changes = + new ConfigChanges() + .addOverride(PipelineOptions.JARS, List.of("a.jar", "b.jar")) + .addOverride(PipelineOptions.GLOBAL_JOB_PARAMETERS, globalJobParameters) + .addOverride( + TaskManagerOptions.TOTAL_PROCESS_MEMORY, + MemorySize.ofMebiBytes(1024)); + + assertThat(changes.getOverrides().get(PipelineOptions.JARS.key())) + .isEqualTo("a.jar;b.jar"); + assertThat(changes.getOverrides().get(PipelineOptions.GLOBAL_JOB_PARAMETERS.key())) + .isEqualTo("k1:v1,k2:v2"); + assertThat(changes.getOverrides().get(TaskManagerOptions.TOTAL_PROCESS_MEMORY.key())) + .isEqualTo("1 gb"); + } finally { + GlobalConfiguration.setStandardYaml(previousDialect); + } + } +} diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/autoscaler/KubernetesScalingRealizer.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/autoscaler/KubernetesScalingRealizer.java index 3d93f62f..e0a2336d 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/autoscaler/KubernetesScalingRealizer.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/autoscaler/KubernetesScalingRealizer.java @@ -168,21 +168,50 @@ public class KubernetesScalingRealizer private static String getOverrideString( KubernetesJobAutoScalerContext context, Map<String, String> newOverrides) { if (context.getResource().getStatus().getReconciliationStatus().isBeforeFirstDeployment()) { - return ConfigurationUtils.convertValue(newOverrides, String.class); + // Use the legacy parser as it is understood by every supported Flink versions + return toLegacyOverrideString(newOverrides); } var conf = context.getResourceContext().getObserveConfig(); - var currentOverrides = - conf.getOptional(PipelineOptions.PARALLELISM_OVERRIDES).orElse(Map.of()); + String currentString = conf.getValue(PipelineOptions.PARALLELISM_OVERRIDES); + var currentOverrides = parseDialectTolerant(currentString); // Check that the overrides actually changed and not just the String representation. // This way we prevent reconciling a NOOP config change which would unnecessarily redeploy // the pipeline. if (currentOverrides.equals(newOverrides)) { - // If overrides are identical, use the previous string as-is. - return conf.getValue(PipelineOptions.PARALLELISM_OVERRIDES); - } else { - return ConfigurationUtils.convertValue(newOverrides, String.class); + if (currentString == null || parsesAsLegacyDialect(currentString, newOverrides)) { + return currentString; + } + } + return toLegacyOverrideString(newOverrides); + } + + /** + * Parse a stored override string accepting both dialects, regardless of the dialect flag the + * observe {@code Configuration} instance happens to carry: the standard parser falls back to + * the legacy pattern, so it understands every format an operator version may have written. + */ + private static Map<String, String> parseDialectTolerant(@Nullable String value) { + if (value == null) { + return Map.of(); + } + try { + return ConfigurationUtils.convertValue(value, Map.class, true); + } catch (Exception e) { + return Map.of(); + } + } + + private static String toLegacyOverrideString(Map<String, String> overrides) { + return ConfigurationUtils.convertValue(overrides, String.class, false); + } + + private static boolean parsesAsLegacyDialect(String value, Map<String, String> expected) { + try { + return expected.equals(ConfigurationUtils.convertValue(value, Map.class, false)); + } catch (Exception e) { + return false; } } } diff --git a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/autoscaler/state/KubernetesAutoScalerStateStore.java b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/autoscaler/state/KubernetesAutoScalerStateStore.java index 16c01e0e..8a59ca0b 100644 --- a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/autoscaler/state/KubernetesAutoScalerStateStore.java +++ b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/autoscaler/state/KubernetesAutoScalerStateStore.java @@ -307,7 +307,7 @@ public class KubernetesAutoScalerStateStore } private static String serializeParallelismOverrides(Map<String, String> overrides) { - return ConfigurationUtils.convertValue(overrides, String.class); + return ConfigurationUtils.convertValue(overrides, String.class, false); } private static Map<String, String> deserializeParallelismOverrides(String overrides) { diff --git a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/autoscaler/KubernetesScalingRealizerTest.java b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/autoscaler/KubernetesScalingRealizerTest.java index 25e7ba19..6da46339 100644 --- a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/autoscaler/KubernetesScalingRealizerTest.java +++ b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/autoscaler/KubernetesScalingRealizerTest.java @@ -18,6 +18,7 @@ package org.apache.flink.kubernetes.operator.autoscaler; import org.apache.flink.autoscaler.tuning.ConfigChanges; +import org.apache.flink.configuration.ConfigurationUtils; import org.apache.flink.configuration.GlobalConfiguration; import org.apache.flink.configuration.MemorySize; import org.apache.flink.configuration.PipelineOptions; @@ -62,6 +63,61 @@ public class KubernetesScalingRealizerTest { overrides -> assertThat(overrides).isEqualTo("b:2,a:1")); } + @Test + public void testApplyOverridesUsesLegacyDialectWhenOperatorRunsStandardYaml() { + boolean previousDialect = GlobalConfiguration.isStandardYaml(); + GlobalConfiguration.setStandardYaml(true); + try { + KubernetesJobAutoScalerContext ctx = + TestingKubernetesAutoscalerUtils.createContext("test", null); + + new KubernetesScalingRealizer() + .realizeParallelismOverrides(ctx, Map.of("a", "1", "b", "2")); + + String overrides = + ctx.getResource() + .getSpec() + .getFlinkConfiguration() + .asFlatMap() + .get(PipelineOptions.PARALLELISM_OVERRIDES.key()); + Map<String, String> legacyParsed = + ConfigurationUtils.convertValue(overrides, Map.class, false); + assertThat(legacyParsed).isEqualTo(Map.of("a", "1", "b", "2")); + } finally { + GlobalConfiguration.setStandardYaml(previousDialect); + } + } + + @Test + public void testStandardDialectOverrideStringIsHealedOnReconcile() { + KubernetesJobAutoScalerContext ctx = + TestingKubernetesAutoscalerUtils.createContext("test", null); + FlinkDeployment resource = (FlinkDeployment) ctx.getResource(); + + resource.getSpec() + .getFlinkConfiguration() + .put(PipelineOptions.PARALLELISM_OVERRIDES.key(), "{a: '1', b: '2'}"); + resource.getStatus() + .getReconciliationStatus() + .serializeAndSetLastReconciledSpec(resource.getSpec(), resource); + resource.getSpec() + .getFlinkConfiguration() + .remove(PipelineOptions.PARALLELISM_OVERRIDES.key()); + + LinkedHashMap<String, String> newOverrides = new LinkedHashMap<>(); + newOverrides.put("a", "1"); + newOverrides.put("b", "2"); + new KubernetesScalingRealizer().realizeParallelismOverrides(ctx, newOverrides); + + assertThat( + ctx.getResource() + .getSpec() + .getFlinkConfiguration() + .asFlatMap() + .get(PipelineOptions.PARALLELISM_OVERRIDES.key())) + .isEqualTo("a:1,b:2"); + } + @Test public void testAutoscalerOverridesStringDoesNotChangeUnlessOverridesChange() { // Create an overrides map which returns the keys in a deterministic order diff --git a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/autoscaler/state/KubernetesAutoScalerStateStoreTest.java b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/autoscaler/state/KubernetesAutoScalerStateStoreTest.java index 3096ad07..0563fd7c 100644 --- a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/autoscaler/state/KubernetesAutoScalerStateStoreTest.java +++ b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/autoscaler/state/KubernetesAutoScalerStateStoreTest.java @@ -23,6 +23,7 @@ import org.apache.flink.autoscaler.metrics.EvaluatedScalingMetric; import org.apache.flink.autoscaler.metrics.ScalingMetric; import org.apache.flink.autoscaler.state.AbstractAutoScalerStateStoreTest; import org.apache.flink.autoscaler.state.AutoScalerStateStore; +import org.apache.flink.configuration.GlobalConfiguration; import org.apache.flink.kubernetes.operator.autoscaler.KubernetesJobAutoScalerContext; import org.apache.flink.runtime.jobgraph.JobVertexID; @@ -291,6 +292,24 @@ public class KubernetesAutoScalerStateStoreTest .containsExactlyInAnyOrderEntriesOf(Map.of(v1, "4")); } + @Test + void testParallelismOverridesStoredInLegacyDialect() throws Exception { + boolean previousDialect = GlobalConfiguration.isStandardYaml(); + GlobalConfiguration.setStandardYaml(true); + try { + stateStore.storeParallelismOverrides(ctx, Map.of("a", "1")); + stateStore.flush(ctx); + assertThat( + configMapStore.getSerializedState( + ctx, KubernetesAutoScalerStateStore.PARALLELISM_OVERRIDES_KEY)) + .hasValue("a:1"); + assertThat(stateStore.getParallelismOverrides(ctx)) + .containsExactlyInAnyOrderEntriesOf(Map.of("a", "1")); + } finally { + GlobalConfiguration.setStandardYaml(previousDialect); + } + } + @Test protected void testDiscardAllState() throws Exception { super.testDiscardAllState();
