From 8d547d256665b0669c0cf6d0d970c78edd5e325f Mon Sep 17 00:00:00 2001 From: wuchao Date: Wed, 8 May 2024 16:22:49 +0800 Subject: [PATCH] [Improvement][parameter] Improvement startup parameters/global parameters data type --- .../api/controller/ExecutorController.java | 14 ++- .../api/service/ExecutorService.java | 6 +- .../api/service/impl/ExecutorServiceImpl.java | 13 ++- .../ExecuteFunctionControllerTest.java | 8 +- .../expand/CuringParamsServiceImpl.java | 26 +++-- .../components/dag/dag-save-modal.tsx | 30 ++++- .../projects/workflow/components/dag/types.ts | 1 + .../definition/components/start-modal.tsx | 110 +++++++++++------- .../definition/components/use-modal.ts | 10 +- .../workflow/definition/create/index.tsx | 2 +- .../workflow/definition/detail/index.tsx | 2 +- .../workflow/instance/detail/index.tsx | 2 +- 12 files changed, 141 insertions(+), 83 deletions(-) diff --git a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/ExecutorController.java b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/ExecutorController.java index a2a300b946..0d9ae055fe 100644 --- a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/ExecutorController.java +++ b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/ExecutorController.java @@ -44,6 +44,7 @@ import org.apache.dolphinscheduler.common.enums.WarningType; import org.apache.dolphinscheduler.common.utils.JSONUtils; import org.apache.dolphinscheduler.dao.entity.User; import org.apache.dolphinscheduler.extract.master.dto.WorkflowExecuteDto; +import org.apache.dolphinscheduler.plugin.task.api.model.Property; import org.apache.commons.lang3.StringUtils; @@ -163,9 +164,9 @@ public class ExecutorController extends BaseController { if (timeout == null) { timeout = Constants.MAX_TASK_TIMEOUT; } - Map startParamMap = null; + List startParamList = null; if (startParams != null) { - startParamMap = JSONUtils.toMap(startParams); + startParamList = JSONUtils.toList(startParams, Property.class); } if (complementDependentMode == null) { @@ -175,7 +176,7 @@ public class ExecutorController extends BaseController { Map result = execService.execProcessInstance(loginUser, projectCode, processDefinitionCode, scheduleTime, execType, failureStrategy, startNodeList, taskDependType, warningType, warningGroupId, runMode, processInstancePriority, - workerGroup, tenantCode, environmentCode, timeout, startParamMap, expectedParallelismNumber, dryRun, + workerGroup, tenantCode, environmentCode, timeout, startParamList, expectedParallelismNumber, dryRun, testFlag, complementDependentMode, version, allLevelDependent, executionOrder); return returnDataList(result); @@ -262,9 +263,9 @@ public class ExecutorController extends BaseController { timeout = Constants.MAX_TASK_TIMEOUT; } - Map startParamMap = null; + List startParamList = null; if (startParams != null) { - startParamMap = JSONUtils.toMap(startParams); + startParamList = JSONUtils.toList(startParams, Property.class); } if (complementDependentMode == null) { @@ -283,7 +284,8 @@ public class ExecutorController extends BaseController { result = execService.execProcessInstance(loginUser, projectCode, processDefinitionCode, scheduleTime, execType, failureStrategy, startNodeList, taskDependType, warningType, warningGroupId, runMode, processInstancePriority, - workerGroup, tenantCode, environmentCode, timeout, startParamMap, expectedParallelismNumber, dryRun, + workerGroup, tenantCode, environmentCode, timeout, startParamList, expectedParallelismNumber, + dryRun, testFlag, complementDependentMode, null, allLevelDependent, executionOrder); diff --git a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/ExecutorService.java b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/ExecutorService.java index 3166d3e718..55ca697412 100644 --- a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/ExecutorService.java +++ b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/ExecutorService.java @@ -30,7 +30,9 @@ import org.apache.dolphinscheduler.common.enums.WarningType; import org.apache.dolphinscheduler.dao.entity.ProcessDefinition; import org.apache.dolphinscheduler.dao.entity.User; import org.apache.dolphinscheduler.extract.master.dto.WorkflowExecuteDto; +import org.apache.dolphinscheduler.plugin.task.api.model.Property; +import java.util.List; import java.util.Map; /** @@ -57,7 +59,7 @@ public interface ExecutorService { * @param environmentCode environment code * @param runMode run mode * @param timeout timeout - * @param startParams the global param values which pass to new process instance + * @param startParamList the global param values which pass to new process instance * @param expectedParallelismNumber the expected parallelism number when execute complement in parallel mode * @param executionOrder the execution order when complementing data * @return execute process instance code @@ -71,7 +73,7 @@ public interface ExecutorService { Priority processInstancePriority, String workerGroup, String tenantCode, Long environmentCode, Integer timeout, - Map startParams, Integer expectedParallelismNumber, + List startParamList, Integer expectedParallelismNumber, int dryRun, int testFlag, ComplementDependentMode complementDependentMode, Integer version, boolean allLevelDependent, ExecutionOrder executionOrder); diff --git a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/ExecutorServiceImpl.java b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/ExecutorServiceImpl.java index 5b576ce95a..7508afdc85 100644 --- a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/ExecutorServiceImpl.java +++ b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/ExecutorServiceImpl.java @@ -90,6 +90,7 @@ import org.apache.dolphinscheduler.extract.master.transportor.StreamingTaskTrigg import org.apache.dolphinscheduler.extract.master.transportor.StreamingTaskTriggerResponse; import org.apache.dolphinscheduler.extract.master.transportor.WorkflowInstanceStateChangeEvent; import org.apache.dolphinscheduler.plugin.task.api.TaskConstants; +import org.apache.dolphinscheduler.plugin.task.api.model.Property; import org.apache.dolphinscheduler.service.command.CommandService; import org.apache.dolphinscheduler.service.cron.CronUtils; import org.apache.dolphinscheduler.service.exceptions.CronParseException; @@ -203,7 +204,7 @@ public class ExecutorServiceImpl extends BaseServiceImpl implements ExecutorServ * @param environmentCode environment code * @param runMode run mode * @param timeout timeout - * @param startParams the global param values which pass to new process instance + * @param startParamList the global param values which pass to new process instance * @param expectedParallelismNumber the expected parallelism number when execute complement in parallel mode * @param testFlag testFlag * @param executionOrder the execution order when complementing data @@ -219,7 +220,7 @@ public class ExecutorServiceImpl extends BaseServiceImpl implements ExecutorServ Priority processInstancePriority, String workerGroup, String tenantCode, Long environmentCode, Integer timeout, - Map startParams, Integer expectedParallelismNumber, + List startParamList, Integer expectedParallelismNumber, int dryRun, int testFlag, ComplementDependentMode complementDependentMode, Integer version, boolean allLevelDependent, ExecutionOrder executionOrder) { @@ -269,7 +270,7 @@ public class ExecutorServiceImpl extends BaseServiceImpl implements ExecutorServ startNodeList, cronTime, warningType, loginUser.getId(), warningGroupId, runMode, processInstancePriority, workerGroup, tenantCode, - environmentCode, startParams, expectedParallelismNumber, dryRun, testFlag, + environmentCode, startParamList, expectedParallelismNumber, dryRun, testFlag, complementDependentMode, allLevelDependent, executionOrder); if (create > 0) { @@ -731,7 +732,7 @@ public class ExecutorServiceImpl extends BaseServiceImpl implements ExecutorServ WarningType warningType, int executorId, Integer warningGroupId, RunMode runMode, Priority processInstancePriority, String workerGroup, String tenantCode, Long environmentCode, - Map startParams, Integer expectedParallelismNumber, int dryRun, + List startParamList, Integer expectedParallelismNumber, int dryRun, int testFlag, ComplementDependentMode complementDependentMode, boolean allLevelDependent, ExecutionOrder executionOrder) { @@ -760,8 +761,8 @@ public class ExecutorServiceImpl extends BaseServiceImpl implements ExecutorServ if (warningType != null) { command.setWarningType(warningType); } - if (startParams != null && startParams.size() > 0) { - cmdParam.put(CMD_PARAM_START_PARAMS, JSONUtils.toJsonString(startParams)); + if (startParamList != null && startParamList.size() > 0) { + cmdParam.put(CMD_PARAM_START_PARAMS, JSONUtils.toJsonString(startParamList)); } command.setCommandParam(JSONUtils.toJsonString(cmdParam)); command.setExecutorId(executorId); diff --git a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/controller/ExecuteFunctionControllerTest.java b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/controller/ExecuteFunctionControllerTest.java index 2c1a62ccc6..5b77acc099 100644 --- a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/controller/ExecuteFunctionControllerTest.java +++ b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/controller/ExecuteFunctionControllerTest.java @@ -38,8 +38,13 @@ import org.apache.dolphinscheduler.common.enums.RunMode; import org.apache.dolphinscheduler.common.enums.TaskDependType; import org.apache.dolphinscheduler.common.enums.WarningType; import org.apache.dolphinscheduler.dao.entity.User; +import org.apache.dolphinscheduler.plugin.task.api.enums.DataType; +import org.apache.dolphinscheduler.plugin.task.api.enums.Direct; +import org.apache.dolphinscheduler.plugin.task.api.model.Property; +import java.util.Collections; import java.util.HashMap; +import java.util.List; import java.util.Map; import org.junit.jupiter.api.Test; @@ -75,7 +80,8 @@ public class ExecuteFunctionControllerTest extends AbstractControllerTest { final String tenantCode = "root"; final Long environmentCode = 4L; final Integer timeout = 5; - final ImmutableMap startParams = ImmutableMap.of("start", "params"); + final List startParams = + Collections.singletonList(new Property("start", Direct.IN, DataType.VARCHAR, "params")); final Integer expectedParallelismNumber = 6; final int dryRun = 7; final int testFlag = 0; diff --git a/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/expand/CuringParamsServiceImpl.java b/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/expand/CuringParamsServiceImpl.java index 9b869fdf18..07084524df 100644 --- a/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/expand/CuringParamsServiceImpl.java +++ b/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/expand/CuringParamsServiceImpl.java @@ -54,6 +54,7 @@ import java.util.Iterator; import java.util.List; import java.util.Map; import java.util.Set; +import java.util.function.Function; import java.util.stream.Collectors; import javax.annotation.Nullable; @@ -63,6 +64,9 @@ import lombok.NonNull; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Component; +import com.google.gson.JsonElement; +import com.google.gson.JsonParser; + @Component public class CuringParamsServiceImpl implements CuringParamsService { @@ -151,8 +155,16 @@ public class CuringParamsServiceImpl implements CuringParamsService { } String startParamJson = cmdParam.get(CommandKeyConstants.CMD_PARAM_START_PARAMS); Map startParamMap = JSONUtils.toMap(startParamJson); - return startParamMap.entrySet().stream().collect(Collectors.toMap(Map.Entry::getKey, - entry -> new Property(entry.getKey(), Direct.IN, DataType.VARCHAR, entry.getValue()))); + JsonElement jsonElement = JsonParser.parseString(startParamJson); + // check whether it is json + boolean isJson = jsonElement.isJsonObject(); + if (isJson) { + return startParamMap.entrySet().stream().collect(Collectors.toMap(Map.Entry::getKey, + entry -> new Property(entry.getKey(), Direct.IN, DataType.VARCHAR, entry.getValue()))); + } else { + List propList = JSONUtils.toList(startParamJson, Property.class); + return propList.stream().collect(Collectors.toMap(Property::getProp, Function.identity())); + } } @Override @@ -181,8 +193,7 @@ public class CuringParamsServiceImpl implements CuringParamsService { Map prepareParamsMap = new HashMap<>(); // assign value to definedParams here - Map globalParamsMap = setGlobalParamsMap(processInstance); - Map globalParams = ParameterUtils.getUserDefParamsMap(globalParamsMap); + Map globalParams = setGlobalParamsMap(processInstance); // combining local and global parameters Map localParams = parameters.getInputLocalParametersMap(); @@ -287,15 +298,16 @@ public class CuringParamsServiceImpl implements CuringParamsService { Long.toString(taskInstance.getProcessInstance().getProcessDefinition().getProjectCode())); return params; } - private Map setGlobalParamsMap(ProcessInstance processInstance) { - Map globalParamsMap = new HashMap<>(16); + private Map setGlobalParamsMap(ProcessInstance processInstance) { + Map globalParamsMap = new HashMap<>(16); // global params string String globalParamsStr = processInstance.getGlobalParams(); if (globalParamsStr != null) { List globalParamsList = JSONUtils.toList(globalParamsStr, Property.class); globalParamsMap - .putAll(globalParamsList.stream().collect(Collectors.toMap(Property::getProp, Property::getValue))); + .putAll(globalParamsList.stream() + .collect(Collectors.toMap(Property::getProp, Function.identity()))); } return globalParamsMap; } diff --git a/dolphinscheduler-ui/src/views/projects/workflow/components/dag/dag-save-modal.tsx b/dolphinscheduler-ui/src/views/projects/workflow/components/dag/dag-save-modal.tsx index 9a55f8d3b2..bfd76e79b1 100644 --- a/dolphinscheduler-ui/src/views/projects/workflow/components/dag/dag-save-modal.tsx +++ b/dolphinscheduler-ui/src/views/projects/workflow/components/dag/dag-save-modal.tsx @@ -154,7 +154,8 @@ export default defineComponent({ formValue.value.globalParams = process.globalParamList.map((param) => ({ key: param.prop, value: param.value, - direct: param.direct + direct: param.direct, + type: param.type })) } } @@ -239,6 +240,7 @@ export default defineComponent({ return { key: '', direct: 'IN', + type: 'VARCHAR', value: '' } }} @@ -246,16 +248,16 @@ export default defineComponent({ > {{ default: (param: { - value: { key: string; direct: string; value: string } + value: { key: string; direct: string; type: string; value: string } }) => ( - + - + - + + + + - {this.startParamsList.length === 0 ? ( - - - - - - ) : ( - - {this.startParamsList.map((item, index) => ( - - - this.updateParamsList(index, param) - } - /> - this.removeStartParams(index)} - class='btn-delete-custom-parameter' - > - - - - - - - - - - - ))} - - )} + { + return { + key: '', + direct: 'IN', + type: 'VARCHAR', + value: '' + } + }} + class='input-startup-params' + > + {{ + default: (param: { + value: { prop: string; direct: string; type: string; value: string } + }) => ( + + + + + + + + + + + + + + + ) + }} +