This is an automated email from the ASF dual-hosted git repository.
wenhemin pushed a commit to branch json_split
in repository https://gitbox.apache.org/repos/asf/incubator-dolphinscheduler.git
The following commit(s) were added to refs/heads/json_split by this push:
new a467fff [Feature][JsonSplit] fix taskId in json (#5165)
a467fff is described below
commit a467fffd7b80627895200e50a59a07e018eeab42
Author: JinyLeeChina <[email protected]>
AuthorDate: Mon Mar 29 14:22:40 2021 +0800
[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 <[email protected]>
---
.../controller/ProcessDefinitionController.java | 9 +++--
.../dolphinscheduler/common/utils/JSONUtils.java | 2 +-
.../server/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 47ab8ac..d2cbe78 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<String, Object> 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<String, Object> 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<String, Object> 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 73af579..b8c949b 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 e727ae0..d7e059d 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<TaskNode> taskNodeList =
processService.genTaskNodeList(processInstance.getProcessDefinitionCode(),
processInstance.getProcessDefinitionVersion());
+ List<TaskNode> 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 9c9a4e7..0220028 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<String,String> fatherParams) {
+ private String getSubWorkFlowParam(ProcessInstanceMap instanceMap,
ProcessInstance parentProcessInstance, Map<String, String> fatherParams) {
// set sub work process command
String processMapStr = JSONUtils.toJsonString(instanceMap);
Map<String, String> cmdParam = JSONUtils.toMap(processMapStr);
@@ -1165,13 +1166,13 @@ public class ProcessService {
Object localParams = subProcessParam.get(Constants.LOCAL_PARAMS);
List<Property> allParam =
JSONUtils.toList(JSONUtils.toJsonString(localParams), Property.class);
Map<String, String> globalMap =
this.getGlobalParamMap(parentProcessInstance.getGlobalParams());
- Map<String,String> fatherParams = new HashMap<>();
+ Map<String, String> 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<TaskNode> taskNodeList = (processData.getTasks() == null) ? new
ArrayList<>() : processData.getTasks();
Map<String, TaskDefinition> 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<String, TaskNode, TaskNodeRelation>
genDagGraph(ProcessDefinition processDefinition) {
- List<TaskNode> taskNodeList =
genTaskNodeList(processDefinition.getCode(), processDefinition.getVersion());
+ Map<String, String> locationMap =
locationToMap(processDefinition.getLocations());
+ List<TaskNode> taskNodeList =
genTaskNodeList(processDefinition.getCode(), processDefinition.getVersion(),
locationMap);
List<ProcessTaskRelationLog> 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<TaskNode> taskNodes =
genTaskNodeList(processDefinition.getCode(), processDefinition.getVersion());
+ Map<String, String> locationMap =
locationToMap(processDefinition.getLocations());
+ List<TaskNode> 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<TaskNode> genTaskNodeList(Long processCode, int
processVersion) {
+ public List<TaskNode> genTaskNodeList(Long processCode, int
processVersion, Map<String, String> locationMap) {
List<ProcessTaskRelationLog> processTaskRelations =
processTaskRelationLogMapper.queryByProcessCodeAndVersion(processCode,
processVersion);
Set<TaskDefinition> taskDefinitionSet = new HashSet<>();
Map<Long, TaskNode> taskNodeMap = new HashMap<>();
@@ -2501,10 +2503,9 @@ public class ProcessService {
Map<Long, TaskDefinitionLog> 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,7 +2518,6 @@ 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())));
});
@@ -2525,7 +2525,28 @@ public class ProcessService {
}
/**
+ * parse locations
+ *
+ * @param locations processDefinition locations
+ * @return key:taskName,value:taskId
+ */
+ public Map<String, String> locationToMap(String locations) {
+ Map<String, String> frontTaskIdAndNameMap = new HashMap<>();
+ if (StringUtils.isBlank(locations)) {
+ return frontTaskIdAndNameMap;
+ }
+ ObjectNode jsonNodes = JSONUtils.parseObject(locations);
+ Iterator<Entry<String, JsonNode>> fields = jsonNodes.fields();
+ while (fields.hasNext()) {
+ Entry<String, JsonNode> 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 1db3760..1593374 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<String, String> 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<String, String> 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 e9a3603..4d8c970 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': '新创建子工作流还未执行,不能进入子工作流',