From a467fffd7b80627895200e50a59a07e018eeab42 Mon Sep 17 00:00:00 2001 From: JinyLeeChina <42576980+JinyLeeChina@users.noreply.github.com> Date: Mon, 29 Mar 2021 14:22:40 +0800 Subject: [PATCH] [Feature][JsonSplit] fix taskId in json (#5165) * modify checkDAGRing and ProcessService method * merge * modify dagRing * modify process instance for project home page * fix save process bug * codeStyle * Fix logical bug in saving process definition * codeSytle * Fix bug in interface of queryProcessDefinitionList * codeSytle * Fix api bug" * fix taskId in processDefinitionJson * fix json taskId * codeStyle Co-authored-by: JinyLeeChina <297062848@qq.com> --- .../ProcessDefinitionController.java | 9 ++-- .../common/utils/JSONUtils.java | 2 +- .../master/runner/MasterExecThread.java | 17 ++++--- .../service/process/ProcessService.java | 47 ++++++++++++++----- .../service/process/ProcessServiceTest.java | 15 ++++++ .../src/js/module/i18n/locale/zh_CN.js | 2 +- 6 files changed, 65 insertions(+), 27 deletions(-) diff --git a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/ProcessDefinitionController.java b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/ProcessDefinitionController.java index 47ab8ac480..d2cbe78029 100644 --- a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/ProcessDefinitionController.java +++ b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/ProcessDefinitionController.java @@ -292,7 +292,8 @@ public class ProcessDefinitionController extends BaseController { @RequestParam(value = "pageNo") int pageNo, @RequestParam(value = "pageSize") int pageSize, @RequestParam(value = "processDefinitionId") int processDefinitionId) { - + logger.info("login user {}, query process versions, project name: {}, pageNo: {}, pageSize: {}, processDefinitionId: {}", + loginUser.getUserName(), projectName, pageNo, pageSize, processDefinitionId); Map result = processDefinitionService.queryProcessDefinitionVersions(loginUser , projectName, pageNo, pageSize, processDefinitionId); @@ -320,7 +321,8 @@ public class ProcessDefinitionController extends BaseController { @ApiParam(name = "projectName", value = "PROJECT_NAME", required = true) @PathVariable String projectName, @RequestParam(value = "processDefinitionId") int processDefinitionId, @RequestParam(value = "version") long version) { - + logger.info("login user {}, switch process version, project name: {}, processDefinitionId: {}, version: {}", + loginUser.getUserName(), projectName, processDefinitionId, version); Map result = processDefinitionService.switchProcessDefinitionVersion(loginUser, projectName , processDefinitionId, version); return returnDataList(result); @@ -347,7 +349,8 @@ public class ProcessDefinitionController extends BaseController { @ApiParam(name = "projectName", value = "PROJECT_NAME", required = true) @PathVariable String projectName, @RequestParam(value = "processDefinitionId") int processDefinitionId, @RequestParam(value = "version") long version) { - + logger.info("login user {}, delete process definition, project name: {}, processDefinitionId: {}, version: {}", + loginUser.getUserName(), projectName, processDefinitionId, version); Map result = processDefinitionService.deleteByProcessDefinitionIdAndVersion(loginUser, projectName, processDefinitionId, version); return returnDataList(result); } diff --git a/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/utils/JSONUtils.java b/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/utils/JSONUtils.java index 73af57929b..b8c949b80c 100644 --- a/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/utils/JSONUtils.java +++ b/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/utils/JSONUtils.java @@ -206,7 +206,7 @@ public class JSONUtils { return null; } - return node.toString(); + return node.asText(); } /** diff --git a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/runner/MasterExecThread.java b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/runner/MasterExecThread.java index e727ae02db..d7e059da1d 100644 --- a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/runner/MasterExecThread.java +++ b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/runner/MasterExecThread.java @@ -78,7 +78,6 @@ import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutorService; import java.util.concurrent.Future; -import java.util.stream.Collectors; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -185,8 +184,8 @@ public class MasterExecThread implements Runnable { /** * constructor of MasterExecThread * - * @param processInstance processInstance - * @param processService processService + * @param processInstance processInstance + * @param processService processService * @param nettyRemotingClient nettyRemotingClient */ public MasterExecThread(ProcessInstance processInstance @@ -388,7 +387,7 @@ public class MasterExecThread implements Runnable { */ private void buildFlowDag() throws Exception { recoverNodeIdList = getStartTaskInstanceList(processInstance.getCommandParam()); - List taskNodeList = processService.genTaskNodeList(processInstance.getProcessDefinitionCode(), processInstance.getProcessDefinitionVersion()); + List taskNodeList = processService.genTaskNodeList(processInstance.getProcessDefinitionCode(), processInstance.getProcessDefinitionVersion(), new HashMap<>()); forbiddenTaskList.clear(); taskNodeList.stream().forEach(taskNode -> { if (taskNode.isForbidden()) { @@ -496,7 +495,7 @@ public class MasterExecThread implements Runnable { * encapsulation task * * @param processInstance process instance - * @param nodeName node name + * @param nodeName node name * @return TaskInstance */ private TaskInstance createTaskInstance(ProcessInstance processInstance, String nodeName, @@ -1274,10 +1273,10 @@ public class MasterExecThread implements Runnable { /** * generate flow dag * - * @param totalTaskNodeList total task node list - * @param startNodeNameList start node name list - * @param recoveryNodeNameList recovery node name list - * @param depNodeType depend node type + * @param totalTaskNodeList total task node list + * @param startNodeNameList start node name list + * @param recoveryNodeNameList recovery node name list + * @param depNodeType depend node type * @return ProcessDag process dag * @throws Exception exception */ diff --git a/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/process/ProcessService.java b/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/process/ProcessService.java index 9c9a4e7fbc..022002865d 100644 --- a/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/process/ProcessService.java +++ b/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/process/ProcessService.java @@ -120,6 +120,7 @@ import java.util.HashSet; import java.util.Iterator; import java.util.List; import java.util.Map; +import java.util.Map.Entry; import java.util.Objects; import java.util.Set; import java.util.stream.Collectors; @@ -1121,7 +1122,7 @@ public class ProcessService { /** * complement data needs transform parent parameter to child. */ - private String getSubWorkFlowParam(ProcessInstanceMap instanceMap, ProcessInstance parentProcessInstance,Map fatherParams) { + private String getSubWorkFlowParam(ProcessInstanceMap instanceMap, ProcessInstance parentProcessInstance, Map fatherParams) { // set sub work process command String processMapStr = JSONUtils.toJsonString(instanceMap); Map cmdParam = JSONUtils.toMap(processMapStr); @@ -1165,13 +1166,13 @@ public class ProcessService { Object localParams = subProcessParam.get(Constants.LOCAL_PARAMS); List allParam = JSONUtils.toList(JSONUtils.toJsonString(localParams), Property.class); Map globalMap = this.getGlobalParamMap(parentProcessInstance.getGlobalParams()); - Map fatherParams = new HashMap<>(); + Map fatherParams = new HashMap<>(); if (CollectionUtils.isNotEmpty(allParam)) { for (Property info : allParam) { fatherParams.put(info.getProp(), globalMap.get(info.getProp())); } } - String processParam = getSubWorkFlowParam(instanceMap, parentProcessInstance,fatherParams); + String processParam = getSubWorkFlowParam(instanceMap, parentProcessInstance, fatherParams); return new Command( commandType, @@ -2239,8 +2240,7 @@ public class ProcessService { } private void setTaskFromTaskNode(TaskNode taskNode, TaskDefinition taskDefinition) { - // TODO for the front-end UI, name with id - taskDefinition.setName(taskNode.getId() + "|" + taskNode.getName()); + taskDefinition.setName(taskNode.getName()); taskDefinition.setDescription(taskNode.getDesc()); taskDefinition.setTaskType(TaskType.of(taskNode.getType())); taskDefinition.setTaskParams(TaskType.of(taskNode.getType()) == TaskType.DEPENDENT ? taskNode.getDependence() : taskNode.getParams()); @@ -2340,7 +2340,7 @@ public class ProcessService { List taskNodeList = (processData.getTasks() == null) ? new ArrayList<>() : processData.getTasks(); Map taskNameAndCode = new HashMap<>(); for (TaskNode taskNode : taskNodeList) { - TaskDefinition taskDefinition = taskDefinitionMapper.queryByDefinitionName(projectCode, taskNode.getName()); + TaskDefinition taskDefinition = taskDefinitionMapper.queryByDefinitionCode(taskNode.getCode()); if (taskDefinition == null) { try { long code = SnowFlakeUtils.getInstance().nextId(); @@ -2448,7 +2448,8 @@ public class ProcessService { * @return dag graph */ public DAG genDagGraph(ProcessDefinition processDefinition) { - List taskNodeList = genTaskNodeList(processDefinition.getCode(), processDefinition.getVersion()); + Map locationMap = locationToMap(processDefinition.getLocations()); + List taskNodeList = genTaskNodeList(processDefinition.getCode(), processDefinition.getVersion(), locationMap); List processTaskRelations = processTaskRelationLogMapper.queryByProcessCodeAndVersion(processDefinition.getCode(), processDefinition.getVersion()); ProcessDag processDag = DagHelper.getProcessDag(taskNodeList, new ArrayList<>(processTaskRelations)); // Generate concrete Dag to be executed @@ -2459,7 +2460,8 @@ public class ProcessService { * generate ProcessData */ public ProcessData genProcessData(ProcessDefinition processDefinition) { - List taskNodes = genTaskNodeList(processDefinition.getCode(), processDefinition.getVersion()); + Map locationMap = locationToMap(processDefinition.getLocations()); + List taskNodes = genTaskNodeList(processDefinition.getCode(), processDefinition.getVersion(), locationMap); ProcessData processData = new ProcessData(); processData.setTasks(taskNodes); processData.setGlobalParams(JSONUtils.toList(processDefinition.getGlobalParams(), Property.class)); @@ -2468,7 +2470,7 @@ public class ProcessService { return processData; } - public List genTaskNodeList(Long processCode, int processVersion) { + public List genTaskNodeList(Long processCode, int processVersion, Map locationMap) { List processTaskRelations = processTaskRelationLogMapper.queryByProcessCodeAndVersion(processCode, processVersion); Set taskDefinitionSet = new HashSet<>(); Map taskNodeMap = new HashMap<>(); @@ -2501,10 +2503,9 @@ public class ProcessService { Map taskDefinitionLogMap = taskDefinitionLogs.stream().collect(Collectors.toMap(TaskDefinitionLog::getCode, log -> log)); taskNodeMap.forEach((k, v) -> { TaskDefinitionLog taskDefinitionLog = taskDefinitionLogMap.get(k); - // TODO split from name - v.setId(StringUtils.substringBefore(taskDefinitionLog.getName(), "|")); + v.setId(locationMap.get(taskDefinitionLog.getName())); v.setCode(taskDefinitionLog.getCode()); - v.setName(StringUtils.substringAfter(taskDefinitionLog.getName(), "|")); + v.setName(taskDefinitionLog.getName()); v.setDesc(taskDefinitionLog.getDescription()); v.setType(taskDefinitionLog.getTaskType().getDescp().toUpperCase()); v.setRunFlag(taskDefinitionLog.getFlag() == Flag.YES ? Constants.FLOWNODE_RUN_FLAG_NORMAL : Constants.FLOWNODE_RUN_FLAG_FORBIDDEN); @@ -2517,15 +2518,35 @@ public class ProcessService { v.setTimeout(JSONUtils.toJsonString(new TaskTimeoutParameter(taskDefinitionLog.getTimeoutFlag() == TimeoutFlag.OPEN, taskDefinitionLog.getTimeoutNotifyStrategy(), taskDefinitionLog.getTimeout()))); - // TODO name will be remove v.getPreTaskNodeList().forEach(task -> task.setName(taskDefinitionLogMap.get(task.getCode()).getName())); v.setPreTasks(JSONUtils.toJsonString(v.getPreTaskNodeList().stream().map(PreviousTaskNode::getName).collect(Collectors.toList()))); }); return new ArrayList<>(taskNodeMap.values()); } + /** + * parse locations + * + * @param locations processDefinition locations + * @return key:taskName,value:taskId + */ + public Map locationToMap(String locations) { + Map frontTaskIdAndNameMap = new HashMap<>(); + if (StringUtils.isBlank(locations)) { + return frontTaskIdAndNameMap; + } + ObjectNode jsonNodes = JSONUtils.parseObject(locations); + Iterator> fields = jsonNodes.fields(); + while (fields.hasNext()) { + Entry jsonNodeEntry = fields.next(); + frontTaskIdAndNameMap.put(JSONUtils.findValue(jsonNodeEntry.getValue(), "name"), jsonNodeEntry.getKey()); + } + return frontTaskIdAndNameMap; + } + /** * add authorized resources + * * @param ownResources own resources * @param userId userId */ diff --git a/dolphinscheduler-service/src/test/java/org/apache/dolphinscheduler/service/process/ProcessServiceTest.java b/dolphinscheduler-service/src/test/java/org/apache/dolphinscheduler/service/process/ProcessServiceTest.java index 1db3760366..15933747ae 100644 --- a/dolphinscheduler-service/src/test/java/org/apache/dolphinscheduler/service/process/ProcessServiceTest.java +++ b/dolphinscheduler-service/src/test/java/org/apache/dolphinscheduler/service/process/ProcessServiceTest.java @@ -503,4 +503,19 @@ public class ProcessServiceTest { Assert.assertEquals(processDefinitionJson, json); } + + @Test + public void locationToMap() { + String locations = "{\"tasks-64888\":{\"name\":\"test_a\",\"targetarr\":\"\",\"nodenumber\":\"1\",\"x\":134,\"y\":183}," + + "\"tasks-24501\":{\"name\":\"test_b\",\"targetarr\":\"tasks-64888\",\"nodenumber\":\"0\",\"x\":392,\"y\":184}," + + "\"tasks-81137\":{\"name\":\"test_c\",\"targetarr\":\"\",\"nodenumber\":\"1\",\"x\":122,\"y\":327}," + + "\"tasks-41367\":{\"name\":\"test_d\",\"targetarr\":\"tasks-81137\",\"nodenumber\":\"0\",\"x\":409,\"y\":324}}"; + Map frontTaskIdAndNameMap = new HashMap<>(); + frontTaskIdAndNameMap.put("test_a", "tasks-64888"); + frontTaskIdAndNameMap.put("test_b", "tasks-24501"); + frontTaskIdAndNameMap.put("test_c", "tasks-81137"); + frontTaskIdAndNameMap.put("test_d", "tasks-41367"); + Map locationToMap = processService.locationToMap(locations); + Assert.assertEquals(frontTaskIdAndNameMap, locationToMap); + } } diff --git a/dolphinscheduler-ui/src/js/module/i18n/locale/zh_CN.js b/dolphinscheduler-ui/src/js/module/i18n/locale/zh_CN.js index e9a3603e2a..4d8c9704dc 100755 --- a/dolphinscheduler-ui/src/js/module/i18n/locale/zh_CN.js +++ b/dolphinscheduler-ui/src/js/module/i18n/locale/zh_CN.js @@ -46,7 +46,7 @@ export default { Minute: '分', 'Delay execution time': '延时执行时间', 'Delay execution': '延时执行', - 'Forced success': '强制成功过', + 'Forced success': '强制成功', Cancel: '取消', 'Confirm add': '确认添加', 'The newly created sub-Process has not yet been executed and cannot enter the sub-Process': '新创建子工作流还未执行,不能进入子工作流',