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