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

zhongjiajie pushed a commit to branch 3.2.1-prepare
in repository https://gitbox.apache.org/repos/asf/dolphinscheduler.git

commit 57a4325821a439bb40071fccc9a05a49f7cb1bf4
Author: Jay Chung <[email protected]>
AuthorDate: Tue Feb 6 18:32:26 2024 +0800

    Revert "Fix k8sTaskExecutionContext setting configYaml (#15116)"
    
    This reverts commit ce11674668cf26747f34c3474763c32e1ba5b6d1.
---
 .../plugin/task/api/am/KubernetesApplicationManager.java    |  3 +++
 .../plugin/task/api/k8s/impl/K8sTaskExecutor.java           |  4 +++-
 .../plugin/task/api/parameters/AbstractParameters.java      | 13 -------------
 .../plugin/task/api/parameters/K8sTaskParameters.java       | 13 +++++++++++--
 .../apache/dolphinscheduler/plugin/task/k8s/K8sTask.java    |  6 ++----
 5 files changed, 19 insertions(+), 20 deletions(-)

diff --git 
a/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/am/KubernetesApplicationManager.java
 
b/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/am/KubernetesApplicationManager.java
index ac1ce69f76..a18637a4ff 100644
--- 
a/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/am/KubernetesApplicationManager.java
+++ 
b/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/am/KubernetesApplicationManager.java
@@ -22,6 +22,7 @@ import static 
org.apache.dolphinscheduler.plugin.task.api.TaskConstants.UNIQUE_L
 
 import org.apache.dolphinscheduler.common.enums.ResourceManagerType;
 import org.apache.dolphinscheduler.common.thread.ThreadUtils;
+import org.apache.dolphinscheduler.common.utils.JSONUtils;
 import org.apache.dolphinscheduler.plugin.task.api.K8sTaskExecutionContext;
 import org.apache.dolphinscheduler.plugin.task.api.TaskException;
 import org.apache.dolphinscheduler.plugin.task.api.enums.TaskExecutionStatus;
@@ -131,6 +132,8 @@ public class KubernetesApplicationManager implements 
ApplicationManager {
     private KubernetesClient getClient(KubernetesApplicationManagerContext 
kubernetesApplicationManagerContext) {
         K8sTaskExecutionContext k8sTaskExecutionContext =
                 
kubernetesApplicationManagerContext.getK8sTaskExecutionContext();
+        k8sTaskExecutionContext
+                
.setConfigYaml(JSONUtils.getNodeString(k8sTaskExecutionContext.getConnectionParams(),
 "kubeConfig"));
         return 
cacheClientMap.computeIfAbsent(kubernetesApplicationManagerContext.getLabelValue(),
                 key -> new KubernetesClientBuilder()
                         
.withConfig(Config.fromKubeconfig(k8sTaskExecutionContext.getConfigYaml())).build());
diff --git 
a/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/k8s/impl/K8sTaskExecutor.java
 
b/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/k8s/impl/K8sTaskExecutor.java
index 986c9dc8a7..4d41a85bbb 100644
--- 
a/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/k8s/impl/K8sTaskExecutor.java
+++ 
b/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/k8s/impl/K8sTaskExecutor.java
@@ -300,7 +300,9 @@ public class K8sTaskExecutor extends 
AbstractK8sTaskExecutor {
                 return result;
             }
             K8sTaskExecutionContext k8sTaskExecutionContext = 
taskRequest.getK8sTaskExecutionContext();
-            String configYaml = k8sTaskExecutionContext.getConfigYaml();
+            String connectionParams = 
k8sTaskExecutionContext.getConnectionParams();
+            String kubeConfig = JSONUtils.getNodeString(connectionParams, 
"kubeConfig");
+            String configYaml = kubeConfig;
             k8sUtils.buildClient(configYaml);
             submitJob2k8s(k8sParameterStr);
             parsePodLogOutput();
diff --git 
a/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/parameters/AbstractParameters.java
 
b/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/parameters/AbstractParameters.java
index 6ca1be7d7a..a57eececf5 100644
--- 
a/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/parameters/AbstractParameters.java
+++ 
b/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/parameters/AbstractParameters.java
@@ -18,12 +18,9 @@
 package org.apache.dolphinscheduler.plugin.task.api.parameters;
 
 import org.apache.dolphinscheduler.common.utils.JSONUtils;
-import org.apache.dolphinscheduler.plugin.task.api.K8sTaskExecutionContext;
 import org.apache.dolphinscheduler.plugin.task.api.enums.Direct;
-import org.apache.dolphinscheduler.plugin.task.api.enums.ResourceType;
 import org.apache.dolphinscheduler.plugin.task.api.model.Property;
 import org.apache.dolphinscheduler.plugin.task.api.model.ResourceInfo;
-import 
org.apache.dolphinscheduler.plugin.task.api.parameters.resource.DataSourceParameters;
 import 
org.apache.dolphinscheduler.plugin.task.api.parameters.resource.ResourceParametersHelper;
 
 import org.apache.commons.collections4.CollectionUtils;
@@ -89,16 +86,6 @@ public abstract class AbstractParameters implements 
IParameters {
         return localParametersMaps;
     }
 
-    public K8sTaskExecutionContext 
generateK8sTaskExecutionContext(ResourceParametersHelper parametersHelper,
-                                                                   int 
datasource) {
-        DataSourceParameters dataSourceParameters =
-                (DataSourceParameters) 
parametersHelper.getResourceParameters(ResourceType.DATASOURCE, datasource);
-        K8sTaskExecutionContext k8sTaskExecutionContext = new 
K8sTaskExecutionContext();
-        k8sTaskExecutionContext.setConnectionParams(
-                Objects.nonNull(dataSourceParameters) ? 
dataSourceParameters.getConnectionParams() : null);
-        return k8sTaskExecutionContext;
-    }
-
     /**
      * get input local parameters map if the param direct is IN
      * @return parameters map
diff --git 
a/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/parameters/K8sTaskParameters.java
 
b/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/parameters/K8sTaskParameters.java
index 4f045abe19..d3d5e4963b 100644
--- 
a/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/parameters/K8sTaskParameters.java
+++ 
b/dolphinscheduler-task-plugin/dolphinscheduler-task-api/src/main/java/org/apache/dolphinscheduler/plugin/task/api/parameters/K8sTaskParameters.java
@@ -17,16 +17,19 @@
 
 package org.apache.dolphinscheduler.plugin.task.api.parameters;
 
+import org.apache.dolphinscheduler.plugin.task.api.K8sTaskExecutionContext;
 import org.apache.dolphinscheduler.plugin.task.api.enums.ResourceType;
 import org.apache.dolphinscheduler.plugin.task.api.model.Label;
 import 
org.apache.dolphinscheduler.plugin.task.api.model.NodeSelectorExpression;
 import org.apache.dolphinscheduler.plugin.task.api.model.ResourceInfo;
+import 
org.apache.dolphinscheduler.plugin.task.api.parameters.resource.DataSourceParameters;
 import 
org.apache.dolphinscheduler.plugin.task.api.parameters.resource.ResourceParametersHelper;
 
 import org.apache.commons.lang3.StringUtils;
 
 import java.util.ArrayList;
 import java.util.List;
+import java.util.Objects;
 
 import lombok.Data;
 import lombok.extern.slf4j.Slf4j;
@@ -55,12 +58,18 @@ public class K8sTaskParameters extends AbstractParameters {
     public boolean checkParameters() {
         return StringUtils.isNotEmpty(image);
     }
-
+    public K8sTaskExecutionContext 
generateExtendedContext(ResourceParametersHelper parametersHelper) {
+        DataSourceParameters dataSourceParameters =
+                (DataSourceParameters) 
parametersHelper.getResourceParameters(ResourceType.DATASOURCE, datasource);
+        K8sTaskExecutionContext k8sTaskExecutionContext = new 
K8sTaskExecutionContext();
+        k8sTaskExecutionContext.setConnectionParams(
+                Objects.nonNull(dataSourceParameters) ? 
dataSourceParameters.getConnectionParams() : null);
+        return k8sTaskExecutionContext;
+    }
     @Override
     public List<ResourceInfo> getResourceFilesList() {
         return new ArrayList<>();
     }
-
     @Override
     public ResourceParametersHelper getResources() {
         ResourceParametersHelper resources = super.getResources();
diff --git 
a/dolphinscheduler-task-plugin/dolphinscheduler-task-k8s/src/main/java/org/apache/dolphinscheduler/plugin/task/k8s/K8sTask.java
 
b/dolphinscheduler-task-plugin/dolphinscheduler-task-k8s/src/main/java/org/apache/dolphinscheduler/plugin/task/k8s/K8sTask.java
index fdb39d7c28..fceb29e163 100644
--- 
a/dolphinscheduler-task-plugin/dolphinscheduler-task-k8s/src/main/java/org/apache/dolphinscheduler/plugin/task/k8s/K8sTask.java
+++ 
b/dolphinscheduler-task-plugin/dolphinscheduler-task-k8s/src/main/java/org/apache/dolphinscheduler/plugin/task/k8s/K8sTask.java
@@ -70,16 +70,14 @@ public class K8sTask extends AbstractK8sTask {
         }
 
         k8sTaskExecutionContext =
-                
k8sTaskParameters.generateK8sTaskExecutionContext(taskExecutionContext.getResourceParametersHelper(),
-                        k8sTaskParameters.getDatasource());
+                
k8sTaskParameters.generateExtendedContext(taskExecutionContext.getResourceParametersHelper());
+        taskRequest.setK8sTaskExecutionContext(k8sTaskExecutionContext);
         k8sConnectionParam =
                 (K8sConnectionParam) 
DataSourceUtils.buildConnectionParams(DbType.valueOf(k8sTaskParameters.getType()),
                         k8sTaskExecutionContext.getConnectionParams());
         String kubeConfig = k8sConnectionParam.getKubeConfig();
         k8sTaskParameters.setNamespace(k8sConnectionParam.getNamespace());
         k8sTaskParameters.setKubeConfig(kubeConfig);
-        k8sTaskExecutionContext.setConfigYaml(kubeConfig);
-        taskRequest.setK8sTaskExecutionContext(k8sTaskExecutionContext);
         log.info("Initialize k8s task params:{}", 
JSONUtils.toPrettyJsonString(k8sTaskParameters));
     }
 

Reply via email to