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) {