diff --git a/streampark-common/src/main/java/org/apache/streampark/common/util/DeflaterUtils.java b/streampark-common/src/main/java/org/apache/streampark/common/util/DeflaterUtils.java index 3b2f33dd98..2da4d1e177 100644 --- a/streampark-common/src/main/java/org/apache/streampark/common/util/DeflaterUtils.java +++ b/streampark-common/src/main/java/org/apache/streampark/common/util/DeflaterUtils.java @@ -53,7 +53,13 @@ public static String zipString(String text) { } public static String unzipString(String zipString) { - byte[] decode = Base64.getDecoder().decode(zipString); + byte[] decode; + try { + decode = Base64.getDecoder().decode(zipString); + } catch (IllegalArgumentException e) { + LOG.warn("Failed to decode base64 string: {}", e.getMessage()); + return null; + } Inflater inflater = new Inflater(); inflater.setInput(decode); byte[] bytes = new byte[256]; diff --git a/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/component/FlinkCheckpointProcessor.java b/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/component/FlinkCheckpointProcessor.java index 2651f50056..53539ed60a 100644 --- a/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/component/FlinkCheckpointProcessor.java +++ b/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/component/FlinkCheckpointProcessor.java @@ -92,8 +92,11 @@ private void process(FlinkApplication application, @Nonnull CheckPoints.CheckPoi if (CheckPointStatusEnum.COMPLETED == status) { if (shouldStoreAsSavepoint(checkPointKey, checkPoint)) { savepointedCache.put(checkPointKey.getSavePointId(), DEFAULT_FLAG_BYTE); - saveSavepoint(checkPoint, application.getId()); - flinkAppHttpWatcher.cleanSavepoint(application); + try { + saveSavepoint(checkPoint, application.getId()); + } finally { + flinkAppHttpWatcher.cleanSavepoint(application); + } return; } 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 66942473b5..cfb539d1fb 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 @@ -118,6 +118,9 @@ public void doSetVersion() { public Map convertFlinkYamlAsMap() { String flinkYamlString = DeflaterUtils.unzipString(flinkConf); + if (flinkYamlString == null) { + return new java.util.HashMap<>(); + } if (isLegacyFlinkConf()) { return FlinkConfigurationUtils.loadLegacyFlinkConf(flinkYamlString); } @@ -168,6 +171,9 @@ public String getVersionOfLast() { public Properties getFlinkConfig() { String flinkYamlString = DeflaterUtils.unzipString(flinkConf); Properties flinkConfig = new Properties(); + if (flinkYamlString == null) { + return flinkConfig; + } Map config = FlinkConfigurationUtils.loadLegacyFlinkConf(flinkYamlString); flinkConfig.putAll(config); return flinkConfig;