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