From 38ec3143ae53d88bcd72b8cc4c416fc8637d88ca Mon Sep 17 00:00:00 2001 From: hutiefang Date: Sun, 21 Jun 2026 02:48:32 +0800 Subject: [PATCH] [ISSUE-4323][Bug] Fix flink conf path recovery --- .../console/core/entity/FlinkEnv.java | 34 +++++++- .../console/core/entity/FlinkEnvTest.java | 86 +++++++++++++++++++ 2 files changed, 117 insertions(+), 3 deletions(-) create mode 100644 streampark-console/streampark-console-service/src/test/java/org/apache/streampark/console/core/entity/FlinkEnvTest.java diff --git a/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/entity/FlinkEnv.java b/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/entity/FlinkEnv.java index 7b13727a93..797c52789c 100644 --- a/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/entity/FlinkEnv.java +++ b/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/entity/FlinkEnv.java @@ -36,6 +36,8 @@ import java.io.File; import java.io.Serializable; import java.nio.charset.StandardCharsets; +import java.nio.file.InvalidPathException; +import java.nio.file.Paths; import java.util.Date; import java.util.Map; import java.util.Properties; @@ -117,7 +119,7 @@ public void doSetVersion() { } public Map convertFlinkYamlAsMap() { - String flinkYamlString = DeflaterUtils.unzipString(flinkConf); + String flinkYamlString = getFlinkConfYaml(); if (isLegacyFlinkConf()) { return FlinkConfigurationUtils.loadLegacyFlinkConf(flinkYamlString); } @@ -133,7 +135,7 @@ public FlinkVersion getFlinkVersion() { } public void unzipFlinkConf() { - this.flinkConf = DeflaterUtils.unzipString(this.flinkConf); + this.flinkConf = getFlinkConfYaml(); } public String getLargeVersion() { @@ -166,13 +168,39 @@ public String getVersionOfLast() { @JsonIgnore public Properties getFlinkConfig() { - String flinkYamlString = DeflaterUtils.unzipString(flinkConf); + String flinkYamlString = getFlinkConfYaml(); Properties flinkConfig = new Properties(); Map config = FlinkConfigurationUtils.loadLegacyFlinkConf(flinkYamlString); flinkConfig.putAll(config); return flinkConfig; } + private String getFlinkConfYaml() { + try { + return DeflaterUtils.unzipString(flinkConf); + } catch (IllegalArgumentException e) { + if (isStoredFlinkHome()) { + doSetFlinkConf(); + return DeflaterUtils.unzipString(flinkConf); + } + throw e; + } + } + + private boolean isStoredFlinkHome() { + if (StringUtils.isBlank(flinkConf) || StringUtils.isBlank(flinkHome)) { + return false; + } + if (StringUtils.equals(flinkConf, flinkHome)) { + return true; + } + try { + return Paths.get(flinkConf).normalize().equals(Paths.get(flinkHome).normalize()); + } catch (InvalidPathException e) { + return false; + } + } + private Float getVersionNumber() { if (StringUtils.isNotBlank(this.version)) { return Float.parseFloat(getVersionOfFirst() + "." + getVersionOfMiddle()); diff --git a/streampark-console/streampark-console-service/src/test/java/org/apache/streampark/console/core/entity/FlinkEnvTest.java b/streampark-console/streampark-console-service/src/test/java/org/apache/streampark/console/core/entity/FlinkEnvTest.java new file mode 100644 index 0000000000..9d769a2b02 --- /dev/null +++ b/streampark-console/streampark-console-service/src/test/java/org/apache/streampark/console/core/entity/FlinkEnvTest.java @@ -0,0 +1,86 @@ +/* + * 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.streampark.console.core.entity; + +import org.apache.streampark.common.util.DeflaterUtils; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.Properties; + +import static org.assertj.core.api.Assertions.assertThat; + +class FlinkEnvTest { + + @TempDir + private Path tempDir; + + @Test + void unzipFlinkConfShouldRecoverWhenStoredValueIsFlinkHome() throws IOException { + FlinkEnv flinkEnv = flinkEnvWithStoredFlinkHome(); + + flinkEnv.unzipFlinkConf(); + + assertThat(flinkEnv.getFlinkConf()) + .contains("execution.checkpointing.savepoint-dir: hdfs:///savepoints"); + } + + @Test + void getFlinkConfigShouldRecoverWhenStoredValueIsFlinkHome() throws IOException { + FlinkEnv flinkEnv = flinkEnvWithStoredFlinkHome(); + + Properties flinkConfig = flinkEnv.getFlinkConfig(); + + assertThat(flinkConfig) + .containsEntry("execution.checkpointing.savepoint-dir", "hdfs:///savepoints"); + assertThat(DeflaterUtils.unzipString(flinkEnv.getFlinkConf())) + .contains("execution.checkpointing.savepoint-dir: hdfs:///savepoints"); + } + + @Test + void getFlinkConfigShouldRecoverWhenStoredValueHasTrailingSlash() throws IOException { + FlinkEnv flinkEnv = flinkEnvWithStoredFlinkHome(); + flinkEnv.setFlinkConf(flinkEnv.getFlinkHome() + "/"); + + Properties flinkConfig = flinkEnv.getFlinkConfig(); + + assertThat(flinkConfig) + .containsEntry("execution.checkpointing.savepoint-dir", "hdfs:///savepoints"); + } + + private FlinkEnv flinkEnvWithStoredFlinkHome() throws IOException { + Path flinkHome = tempDir.resolve("flink-1.16.3"); + Path confDir = flinkHome.resolve("conf"); + Files.createDirectories(confDir); + Files.write( + confDir.resolve("flink-conf.yaml"), + "execution.checkpointing.savepoint-dir: hdfs:///savepoints\n" + .getBytes(StandardCharsets.UTF_8)); + + FlinkEnv flinkEnv = new FlinkEnv(); + flinkEnv.setFlinkHome(flinkHome.toString()); + flinkEnv.setFlinkConf(flinkHome.toString()); + flinkEnv.setVersion("1.16.3"); + return flinkEnv; + } +}