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 0b7fc32f [FLINK-40055] Serialize managed cluster config in the target 
Flink version's YAML dialect (#1152)
0b7fc32f is described below

commit 0b7fc32fc5cd02f589fec7974f1d2752e0af3908
Author: Gerk Elznik <[email protected]>
AuthorDate: Wed Aug 5 09:03:40 2026 -0600

    [FLINK-40055] Serialize managed cluster config in the target Flink 
version's YAML dialect (#1152)
---
 .../decorators/FlinkConfMountDecorator.java        | 18 +++++-----
 .../operator/service/AbstractFlinkService.java     | 40 ++++++++++++++++------
 .../operator/service/AbstractFlinkServiceTest.java | 30 ++++++++++++++++
 3 files changed, 69 insertions(+), 19 deletions(-)

diff --git 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/kubeclient/decorators/FlinkConfMountDecorator.java
 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/kubeclient/decorators/FlinkConfMountDecorator.java
index 2d5cbe6d..c2a83c2c 100644
--- 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/kubeclient/decorators/FlinkConfMountDecorator.java
+++ 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/kubeclient/decorators/FlinkConfMountDecorator.java
@@ -157,7 +157,8 @@ public class FlinkConfMountDecorator extends 
AbstractKubernetesStepDecorator {
         // For Flink versions that use the standard config we have to set the 
standardYaml flag in
         // the Configuration object manually instead of simply cloning, 
otherwise it would simply
         // inherit it from the base config (which would always be false 
currently).
-        Configuration clusterSideConfig = new 
Configuration(useStandardYamlConfig());
+        Configuration clusterSideConfig =
+                new 
Configuration(useStandardYamlConfig(flinkConfig.get(FLINK_VERSION)));
         clusterSideConfig.addAll(flinkConfig);
         // Remove some configuration options that should not be taken to 
cluster side.
         
clusterSideConfig.removeConfig(KubernetesConfigOptions.KUBE_CONFIG_FILE);
@@ -175,7 +176,7 @@ public class FlinkConfMountDecorator extends 
AbstractKubernetesStepDecorator {
     private void validateConfigKeysForV2(Configuration clusterSideConfig) {
 
         // Only validate Flink 2.0 yaml configs
-        if (!useStandardYamlConfig()) {
+        if (!useStandardYamlConfig(clusterSideConfig.get(FLINK_VERSION))) {
             return;
         }
 
@@ -231,7 +232,10 @@ public class FlinkConfMountDecorator extends 
AbstractKubernetesStepDecorator {
      * @return conf file name
      */
     public String getFlinkConfFilename() {
-        return useStandardYamlConfig() ? "config.yaml" : "flink-conf.yaml";
+        return useStandardYamlConfig(
+                        
kubernetesComponentConf.getFlinkConfiguration().get(FLINK_VERSION))
+                ? "config.yaml"
+                : "flink-conf.yaml";
     }
 
     /**
@@ -239,12 +243,10 @@ public class FlinkConfMountDecorator extends 
AbstractKubernetesStepDecorator {
      * flink-conf.yaml. While technically 1.19+ could use this we don't want 
to change the behaviour
      * for already released Flink versions, so only switch to new yaml from 
Flink 2.0 onwards.
      *
+     * @param flinkVersion Flink version of the target deployment
      * @return True for Flink version 2.0 and above
      */
-    boolean useStandardYamlConfig() {
-        return kubernetesComponentConf
-                .getFlinkConfiguration()
-                .get(FLINK_VERSION)
-                .isEqualOrNewer(FlinkVersion.v2_0);
+    public static boolean useStandardYamlConfig(FlinkVersion flinkVersion) {
+        return flinkVersion != null && 
flinkVersion.isEqualOrNewer(FlinkVersion.v2_0);
     }
 }
diff --git 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkService.java
 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkService.java
index 1d48d10d..08467957 100644
--- 
a/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkService.java
+++ 
b/flink-kubernetes-operator/src/main/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkService.java
@@ -31,6 +31,7 @@ import org.apache.flink.configuration.SecurityOptions;
 import org.apache.flink.core.execution.RestoreMode;
 import org.apache.flink.kubernetes.configuration.KubernetesConfigOptions;
 import 
org.apache.flink.kubernetes.kubeclient.decorators.ExternalServiceDecorator;
+import 
org.apache.flink.kubernetes.kubeclient.decorators.FlinkConfMountDecorator;
 import org.apache.flink.kubernetes.operator.api.AbstractFlinkResource;
 import org.apache.flink.kubernetes.operator.api.FlinkDeployment;
 import org.apache.flink.kubernetes.operator.api.FlinkSessionJob;
@@ -938,7 +939,7 @@ public abstract class AbstractFlinkService implements 
FlinkService {
                                     ? RestoreMode.DEFAULT
                                     : null,
                             
conf.get(FLINK_VERSION).isEqualOrNewer(FlinkVersion.v1_17)
-                                    ? conf.toMap()
+                                    ? configToMapWithVersionDialect(conf, 
flinkVersion)
                                     : null);
             LOG.info("Submitting job: {} to session cluster.", jobID);
             clusterClient
@@ -1042,21 +1043,38 @@ public abstract class AbstractFlinkService implements 
FlinkService {
                 timeout);
     }
 
+    /**
+     * Remove operator-only keys, preserving raw values so the write 
boundaries can serialize them
+     * in the target Flink version's YAML dialect.
+     */
     @VisibleForTesting
     protected static Configuration removeOperatorConfigs(Configuration config) 
{
-        Configuration newConfig = new Configuration();
-        config.toMap()
-                .forEach(
-                        (k, v) -> {
-                            if (!k.startsWith(K8S_OP_CONF_PREFIX)
-                                    && 
!k.startsWith(AutoScalerOptions.AUTOSCALER_CONF_PREFIX)) {
-                                newConfig.setString(k, v);
-                            }
-                        });
-
+        Configuration newConfig = new Configuration(config);
+        for (String key : config.keySet()) {
+            if (key.startsWith(K8S_OP_CONF_PREFIX)
+                    || 
key.startsWith(AutoScalerOptions.AUTOSCALER_CONF_PREFIX)) {
+                newConfig.removeKey(key);
+            }
+        }
         return newConfig;
     }
 
+    /**
+     * Serialize the config in the YAML dialect of the given Flink version so 
the receiving cluster
+     * can parse the values regardless of the operator's own config format.
+     *
+     * @param conf Config to serialize
+     * @param flinkVersion Flink version of the receiving cluster
+     * @return Map of config entries in the target version's string format
+     */
+    @VisibleForTesting
+    protected static Map<String, String> configToMapWithVersionDialect(
+            Configuration conf, FlinkVersion flinkVersion) {
+        var copy = new 
Configuration(FlinkConfMountDecorator.useStandardYamlConfig(flinkVersion));
+        copy.addAll(conf);
+        return copy.toMap();
+    }
+
     private void validateHaMetadataExists(Configuration conf) {
         if (!isHaMetadataAvailable(conf)) {
             throw new UpgradeFailureException(
diff --git 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkServiceTest.java
 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkServiceTest.java
index ebfe0616..064c0926 100644
--- 
a/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkServiceTest.java
+++ 
b/flink-kubernetes-operator/src/test/java/org/apache/flink/kubernetes/operator/service/AbstractFlinkServiceTest.java
@@ -25,8 +25,10 @@ import org.apache.flink.api.java.tuple.Tuple4;
 import org.apache.flink.client.program.rest.RestClusterClient;
 import org.apache.flink.configuration.CheckpointingOptions;
 import org.apache.flink.configuration.Configuration;
+import org.apache.flink.configuration.GlobalConfiguration;
 import org.apache.flink.configuration.JobManagerOptions;
 import org.apache.flink.configuration.MemorySize;
+import org.apache.flink.configuration.PipelineOptions;
 import org.apache.flink.configuration.TaskManagerOptions;
 import org.apache.flink.core.execution.SavepointFormatType;
 import org.apache.flink.kubernetes.configuration.KubernetesConfigOptions;
@@ -1089,6 +1091,34 @@ public class AbstractFlinkServiceTest {
         assertTrue(newConf.containsKey(regularKey));
     }
 
+    @Test
+    public void removeOperatorConfigsKeepsTypedValuesTest() {
+        // Simulate an operator process running on standard config.yaml.
+        boolean previousDialect = GlobalConfiguration.isStandardYaml();
+        GlobalConfiguration.setStandardYaml(true);
+        try {
+            var jars = List.of("local:///opt/flink/job.jar");
+            var deployConfig = new Configuration(true);
+            deployConfig.set(PipelineOptions.JARS, jars);
+            deployConfig.setString("kubernetes.operator.myKey", "v");
+
+            var newConf = 
AbstractFlinkService.removeOperatorConfigs(deployConfig);
+
+            assertFalse(newConf.containsKey("kubernetes.operator.myKey"));
+            // Typed values must survive so write boundaries can render them 
per Flink version.
+            assertEquals(jars, newConf.get(PipelineOptions.JARS));
+            assertEquals(
+                    Map.of("pipeline.jars", "local:///opt/flink/job.jar"),
+                    AbstractFlinkService.configToMapWithVersionDialect(
+                            newConf, FlinkVersion.v1_20));
+            assertEquals(
+                    Map.of("pipeline.jars", "['local:///opt/flink/job.jar']"),
+                    
AbstractFlinkService.configToMapWithVersionDialect(newConf, FlinkVersion.v2_0));
+        } finally {
+            GlobalConfiguration.setStandardYaml(previousDialect);
+        }
+    }
+
     @Test
     public void getMetricsTest() throws Exception {
         var jobId = new JobID();

Reply via email to