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)); }
