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

gyfora pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/flink-kubernetes-operator.git


The following commit(s) were added to refs/heads/main by this push:
     new 38a9f197 [FLINK-40474] Autoscaler parallelism overrides serialized in 
the operator's YAML dialect break job submission for Flink 1.x jobs (#1195)
38a9f197 is described below

commit 38a9f197465082a5f5987653b9497d7e5aef384a
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="&lt;mxfile 
host=&quot;Electron&quot; agent=&quot;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&quot; 
version=&quot;29.0.3&quot; pages=&quot;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="&lt;mxfile 
host=&quot;Electron&quot; agent=&quot;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&quot; 
version=&quot;29.0.3&quot; pages=&quot;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();

Reply via email to