diff --git a/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/entity/SparkApplication.java b/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/entity/SparkApplication.java index 2f59d1cd17..1719279386 100644 --- a/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/entity/SparkApplication.java +++ b/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/entity/SparkApplication.java @@ -127,8 +127,8 @@ public class SparkApplication extends BaseEntity implements ApplicationEntitySup /** spark docker base image */ private String k8sContainerImage; - /** k8s image pull policy */ - private int k8sImagePullPolicy; + /** k8s image pull policy — nullable, matching t_spark_app.k8s_image_pull_policy */ + private Integer k8sImagePullPolicy; /** k8s spark service account */ private String k8sServiceAccount; diff --git a/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/entity/SparkEnv.java b/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/entity/SparkEnv.java index dff43e0e80..c153f245a1 100644 --- a/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/entity/SparkEnv.java +++ b/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/entity/SparkEnv.java @@ -73,10 +73,11 @@ public class SparkEnv implements Serializable { public void doSetSparkConf() throws ApiDetailException { try { File yaml = new File(this.sparkHome.concat("/conf/spark-defaults.conf")); - if (yaml.exists()) { - String sparkConf = FileUtils.readFileToString(yaml, StandardCharsets.UTF_8); - this.sparkConf = DeflaterUtils.zipString(sparkConf); - } + // A stock Spark distribution ships only spark-defaults.conf.template, so an absent file + // is the normal case, not an error — but t_spark_env.spark_conf is NOT NULL, so it still + // has to be written as an empty conf rather than left null. + String sparkConf = yaml.exists() ? FileUtils.readFileToString(yaml, StandardCharsets.UTF_8) : ""; + this.sparkConf = DeflaterUtils.zipString(sparkConf); } catch (Exception e) { throw new ApiDetailException(e); } diff --git a/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/enums/SparkAppStateEnum.java b/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/enums/SparkAppStateEnum.java index 9e98476701..93bb89846f 100644 --- a/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/enums/SparkAppStateEnum.java +++ b/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/enums/SparkAppStateEnum.java @@ -78,6 +78,9 @@ public enum SparkAppStateEnum { } public static SparkAppStateEnum of(Integer state) { + if (state == null) { + return SparkAppStateEnum.OTHER; + } for (SparkAppStateEnum appState : values()) { if (appState.value == state) { return appState; diff --git a/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/request/flink/FlinkAppCreateRequest.java b/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/request/flink/FlinkAppCreateRequest.java index ce5839097a..de60c0203b 100644 --- a/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/request/flink/FlinkAppCreateRequest.java +++ b/streampark-console/streampark-console-service/src/main/java/org/apache/streampark/console/core/request/flink/FlinkAppCreateRequest.java @@ -54,6 +54,8 @@ public class FlinkAppCreateRequest implements Serializable { private String flinkSql; + private Long sqlId; + @NotNull @ApiParam(description = "Application type", required = true) private Integer appType;