This is an automated email from the ASF dual-hosted git repository.

benjobs pushed a commit to branch dev-2.1.3
in repository https://gitbox.apache.org/repos/asf/incubator-streampark.git


The following commit(s) were added to refs/heads/dev-2.1.3 by this push:
     new 6de2959bd [Improve] package name issue fixed.
6de2959bd is described below

commit 6de2959bd2e77d49c02910276d993f971ec5cea1
Author: benjobs <[email protected]>
AuthorDate: Sat Jan 13 18:23:09 2024 +0800

    [Improve] package name issue fixed.
---
 .../flink/client/impl/KubernetesNativeApplicationClient.scala       | 2 +-
 .../flink/client/impl/KubernetesNativeSessionClient.scala           | 6 +++---
 .../streampark/flink/client/trait/KubernetesNativeClientTrait.scala | 6 +++---
 3 files changed, 7 insertions(+), 7 deletions(-)

diff --git 
a/streampark-flink/streampark-flink-client/streampark-flink-client-core/src/main/scala/org/apache/streampark/flink/client/impl/KubernetesNativeApplicationClient.scala
 
b/streampark-flink/streampark-flink-client/streampark-flink-client-core/src/main/scala/org/apache/streampark/flink/client/impl/KubernetesNativeApplicationClient.scala
index f5623113e..7dcc250be 100644
--- 
a/streampark-flink/streampark-flink-client/streampark-flink-client-core/src/main/scala/org/apache/streampark/flink/client/impl/KubernetesNativeApplicationClient.scala
+++ 
b/streampark-flink/streampark-flink-client/streampark-flink-client-core/src/main/scala/org/apache/streampark/flink/client/impl/KubernetesNativeApplicationClient.scala
@@ -41,7 +41,7 @@ object KubernetesNativeApplicationClient extends 
KubernetesNativeClientTrait {
 
     // require parameters
     require(
-      StringUtils.isNotBlank(submitRequest.k8sSubmitParam.clusterId),
+      StringUtils.isNotBlank(submitRequest.clusterId),
       s"[flink-submit] submit flink job failed, clusterId is null, 
mode=${flinkConfig.get(DeploymentOptions.TARGET)}"
     )
 
diff --git 
a/streampark-flink/streampark-flink-client/streampark-flink-client-core/src/main/scala/org/apache/streampark/flink/client/impl/KubernetesNativeSessionClient.scala
 
b/streampark-flink/streampark-flink-client/streampark-flink-client-core/src/main/scala/org/apache/streampark/flink/client/impl/KubernetesNativeSessionClient.scala
index dce9a7865..7638be720 100644
--- 
a/streampark-flink/streampark-flink-client/streampark-flink-client-core/src/main/scala/org/apache/streampark/flink/client/impl/KubernetesNativeSessionClient.scala
+++ 
b/streampark-flink/streampark-flink-client/streampark-flink-client-core/src/main/scala/org/apache/streampark/flink/client/impl/KubernetesNativeSessionClient.scala
@@ -51,7 +51,7 @@ object KubernetesNativeSessionClient extends 
KubernetesNativeClientTrait with Lo
       flinkConfig: Configuration): SubmitResponse = {
     // require parameters
     require(
-      StringUtils.isNotBlank(submitRequest.k8sSubmitParam.clusterId),
+      StringUtils.isNotBlank(submitRequest.clusterId),
       s"[flink-submit] submit flink job failed, clusterId is null, 
mode=${flinkConfig.get(DeploymentOptions.TARGET)}"
     )
     super.trySubmit(submitRequest, flinkConfig, submitRequest.userJarFile)(
@@ -69,8 +69,8 @@ object KubernetesNativeSessionClient extends 
KubernetesNativeClientTrait with Lo
       // get jm rest url of flink session cluster
       val clusterKey = ClusterKey(
         FlinkK8sExecuteMode.SESSION,
-        submitRequest.k8sSubmitParam.kubernetesNamespace,
-        submitRequest.k8sSubmitParam.clusterId)
+        submitRequest.kubernetesNamespace,
+        submitRequest.clusterId)
       val jmRestUrl = KubernetesRetriever
         .retrieveFlinkRestUrl(clusterKey)
         .getOrElse(throw new Exception(
diff --git 
a/streampark-flink/streampark-flink-client/streampark-flink-client-core/src/main/scala/org/apache/streampark/flink/client/trait/KubernetesNativeClientTrait.scala
 
b/streampark-flink/streampark-flink-client/streampark-flink-client-core/src/main/scala/org/apache/streampark/flink/client/trait/KubernetesNativeClientTrait.scala
index ff65a998e..9e7468053 100644
--- 
a/streampark-flink/streampark-flink-client/streampark-flink-client-core/src/main/scala/org/apache/streampark/flink/client/trait/KubernetesNativeClientTrait.scala
+++ 
b/streampark-flink/streampark-flink-client/streampark-flink-client-core/src/main/scala/org/apache/streampark/flink/client/trait/KubernetesNativeClientTrait.scala
@@ -40,11 +40,11 @@ trait KubernetesNativeClientTrait extends FlinkClientTrait {
   override def setConfig(submitRequest: SubmitRequest, flinkConfig: 
Configuration): Unit = {
     // extract from submitRequest
     flinkConfig
-      .safeSet(KubernetesConfigOptions.CLUSTER_ID, 
submitRequest.k8sSubmitParam.clusterId)
-      .safeSet(KubernetesConfigOptions.NAMESPACE, 
submitRequest.k8sSubmitParam.kubernetesNamespace)
+      .safeSet(KubernetesConfigOptions.CLUSTER_ID, submitRequest.clusterId)
+      .safeSet(KubernetesConfigOptions.NAMESPACE, 
submitRequest.kubernetesNamespace)
       .safeSet(
         KubernetesConfigOptions.REST_SERVICE_EXPOSED_TYPE,
-        
covertToServiceExposedType(submitRequest.k8sSubmitParam.flinkRestExposedType))
+        covertToServiceExposedType(submitRequest.flinkRestExposedType))
 
     if (submitRequest.buildResult != null) {
       if (submitRequest.executionMode == 
ExecutionMode.KUBERNETES_NATIVE_APPLICATION) {

Reply via email to