This is an automated email from the ASF dual-hosted git repository.
wenhemin pushed a commit to branch json_split_two
in repository https://gitbox.apache.org/repos/asf/dolphinscheduler.git
The following commit(s) were added to refs/heads/json_split_two by this push:
new 3a4c608 [Feature][JsonSplit-api] modify viewTree of ProcessDefiniton
(#5881)
3a4c608 is described below
commit 3a4c6086bf28c5734392a6921475a9b454692b7e
Author: JinyLeeChina <[email protected]>
AuthorDate: Mon Jul 26 11:08:54 2021 +0800
[Feature][JsonSplit-api] modify viewTree of ProcessDefiniton (#5881)
* modify viewTree of ProcessDefiniton
* modify ut
* modify ut
* modify viewTree of ProcessDefiniton
Co-authored-by: JinyLeeChina <[email protected]>
---
.../controller/ProcessDefinitionController.java | 12 ++---
.../api/service/ProcessDefinitionService.java | 4 +-
.../service/impl/ProcessDefinitionServiceImpl.java | 53 ++++++++--------------
.../api/service/ProcessDefinitionServiceTest.java | 7 ++-
.../service/process/ProcessService.java | 21 +++++----
5 files changed, 42 insertions(+), 55 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 e72cf6b..a5dd7ae 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
@@ -473,18 +473,18 @@ public class ProcessDefinitionController extends
BaseController {
}
/**
- * encapsulation treeview structure
+ * encapsulation tree view structure
*
* @param loginUser login user
* @param projectCode project code
- * @param id process definition id
+ * @param code process definition code
* @param limit limit
* @return tree view json data
*/
@ApiOperation(value = "viewTree", notes = "VIEW_TREE_NOTES")
@ApiImplicitParams({
- @ApiImplicitParam(name = "processId", value =
"PROCESS_DEFINITION_ID", required = true, dataType = "Int", example = "100"),
- @ApiImplicitParam(name = "limit", value = "LIMIT", required =
true, dataType = "Int", example = "100")
+ @ApiImplicitParam(name = "code", value = "PROCESS_DEFINITION_CODE",
required = true, dataType = "Long", example = "100"),
+ @ApiImplicitParam(name = "limit", value = "LIMIT", required = true,
dataType = "Int", example = "100")
})
@GetMapping(value = "/view-tree")
@ResponseStatus(HttpStatus.OK)
@@ -492,9 +492,9 @@ public class ProcessDefinitionController extends
BaseController {
@AccessLogAnnotation(ignoreRequestArgs = "loginUser")
public Result viewTree(@ApiIgnore @RequestAttribute(value =
Constants.SESSION_USER) User loginUser,
@ApiParam(name = "projectCode", value =
"PROJECT_CODE", required = true) @PathVariable long projectCode,
- @RequestParam("processId") Integer id,
+ @RequestParam("code") long code,
@RequestParam("limit") Integer limit) throws
Exception {
- Map<String, Object> result = processDefinitionService.viewTree(id,
limit);
+ Map<String, Object> result = processDefinitionService.viewTree(code,
limit);
return returnDataList(result);
}
diff --git
a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/ProcessDefinitionService.java
b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/ProcessDefinitionService.java
index e8b170f..a4a4b98 100644
---
a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/ProcessDefinitionService.java
+++
b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/ProcessDefinitionService.java
@@ -277,12 +277,12 @@ public interface ProcessDefinitionService {
/**
* Encapsulates the TreeView structure
*
- * @param processId process definition id
+ * @param code process definition code
* @param limit limit
* @return tree view json data
* @throws Exception exception
*/
- Map<String, Object> viewTree(Integer processId,
+ Map<String, Object> viewTree(long code,
Integer limit) throws Exception;
/**
diff --git
a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/ProcessDefinitionServiceImpl.java
b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/ProcessDefinitionServiceImpl.java
index dc32d3c..239396a 100644
---
a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/ProcessDefinitionServiceImpl.java
+++
b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/ProcessDefinitionServiceImpl.java
@@ -1272,43 +1272,32 @@ public class ProcessDefinitionServiceImpl extends
BaseServiceImpl implements Pro
/**
* Encapsulates the TreeView structure
*
- * @param processId process definition id
+ * @param code process definition code
* @param limit limit
* @return tree view json data
*/
@Override
- public Map<String, Object> viewTree(Integer processId, Integer limit) {
+ public Map<String, Object> viewTree(long code, Integer limit) {
Map<String, Object> result = new HashMap<>();
-
- ProcessDefinition processDefinition =
processDefinitionMapper.selectById(processId);
+ ProcessDefinition processDefinition =
processDefinitionMapper.queryByCode(code);
if (null == processDefinition) {
logger.info("process define not exists");
- putMsg(result, Status.PROCESS_DEFINE_NOT_EXIST, processDefinition);
+ putMsg(result, Status.PROCESS_DEFINE_NOT_EXIST, code);
return result;
}
DAG<String, TaskNode, TaskNodeRelation> dag =
processService.genDagGraph(processDefinition);
- /**
- * nodes that is running
- */
+ // nodes that is running
Map<String, List<TreeViewDto>> runningNodeMap = new
ConcurrentHashMap<>();
- /**
- * nodes that is waiting torun
- */
+ //nodes that is waiting to run
Map<String, List<TreeViewDto>> waitingRunningNodeMap = new
ConcurrentHashMap<>();
- /**
- * List of process instances
- */
- List<ProcessInstance> processInstanceList =
processInstanceService.queryByProcessDefineCode(processDefinition.getCode(),
limit);
- List<TaskDefinitionLog> taskDefinitionList =
processService.queryTaskDefinitionList(processDefinition.getCode(),
- processDefinition.getVersion());
- Map<Long, TaskDefinition> taskDefinitionMap = new HashedMap();
- taskDefinitionList.forEach(taskDefinitionLog ->
taskDefinitionMap.put(taskDefinitionLog.getCode(), taskDefinitionLog));
-
- for (ProcessInstance processInstance : processInstanceList) {
-
processInstance.setDuration(DateUtils.format2Duration(processInstance.getStartTime(),
processInstance.getEndTime()));
- }
+ // List of process instances
+ List<ProcessInstance> processInstanceList =
processInstanceService.queryByProcessDefineCode(code, limit);
+ processInstanceList.forEach(processInstance ->
processInstance.setDuration(DateUtils.format2Duration(processInstance.getStartTime(),
processInstance.getEndTime())));
+ List<TaskDefinitionLog> taskDefinitionList =
processService.queryTaskDefinitionListByProcess(code,
processDefinition.getVersion());
+ Map<Long, TaskDefinitionLog> taskDefinitionMap =
taskDefinitionList.stream()
+ .collect(Collectors.toMap(TaskDefinitionLog::getCode,
taskDefinitionLog -> taskDefinitionLog));
if (limit > processInstanceList.size()) {
limit = processInstanceList.size();
@@ -1318,13 +1307,12 @@ public class ProcessDefinitionServiceImpl extends
BaseServiceImpl implements Pro
parentTreeViewDto.setName("DAG");
parentTreeViewDto.setType("");
// Specify the process definition, because it is a TreeView for a
process definition
-
for (int i = limit - 1; i >= 0; i--) {
ProcessInstance processInstance = processInstanceList.get(i);
-
Date endTime = processInstance.getEndTime() == null ? new Date() :
processInstance.getEndTime();
- parentTreeViewDto.getInstances().add(new
Instance(processInstance.getId(), processInstance.getName(), "",
processInstance.getState().toString()
- , processInstance.getStartTime(), endTime,
processInstance.getHost(), DateUtils.format2Readable(endTime.getTime() -
processInstance.getStartTime().getTime())));
+ parentTreeViewDto.getInstances().add(new
Instance(processInstance.getId(), processInstance.getName(), "",
+ processInstance.getState().toString(),
processInstance.getStartTime(), endTime, processInstance.getHost(),
+ DateUtils.format2Readable(endTime.getTime() -
processInstance.getStartTime().getTime())));
}
List<TreeViewDto> parentTreeViewDtoList = new ArrayList<>();
@@ -1335,7 +1323,7 @@ public class ProcessDefinitionServiceImpl extends
BaseServiceImpl implements Pro
}
while (Stopper.isRunning()) {
- Set<String> postNodeList = null;
+ Set<String> postNodeList;
Iterator<Map.Entry<String, List<TreeViewDto>>> iter =
runningNodeMap.entrySet().iterator();
while (iter.hasNext()) {
Map.Entry<String, List<TreeViewDto>> en = iter.next();
@@ -1358,16 +1346,15 @@ public class ProcessDefinitionServiceImpl extends
BaseServiceImpl implements Pro
Date endTime = taskInstance.getEndTime() == null ? new
Date() : taskInstance.getEndTime();
int subProcessId = 0;
- /**
- * if process is sub process, the return sub id, or
sub id=0
- */
+ // if process is sub process, the return sub id, or
sub id=0
if (taskInstance.isSubProcess()) {
TaskDefinition taskDefinition =
taskDefinitionMap.get(taskInstance.getTaskCode());
subProcessId =
Integer.parseInt(JSONUtils.parseObject(
taskDefinition.getTaskParams()).path(CMD_PARAM_SUB_PROCESS_DEFINE_ID).asText());
}
- treeViewDto.getInstances().add(new
Instance(taskInstance.getId(), taskInstance.getName(),
taskInstance.getTaskType(), taskInstance.getState().toString()
- , taskInstance.getStartTime(),
taskInstance.getEndTime(), taskInstance.getHost(),
DateUtils.format2Readable(endTime.getTime() - startTime.getTime()),
subProcessId));
+ treeViewDto.getInstances().add(new
Instance(taskInstance.getId(), taskInstance.getName(),
taskInstance.getTaskType(),
+ taskInstance.getState().toString(),
taskInstance.getStartTime(), taskInstance.getEndTime(), taskInstance.getHost(),
+ DateUtils.format2Readable(endTime.getTime() -
startTime.getTime()), subProcessId));
}
}
for (TreeViewDto pTreeViewDto : parentTreeViewDtoList) {
diff --git
a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/ProcessDefinitionServiceTest.java
b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/ProcessDefinitionServiceTest.java
index d2ed647..bf89ed1 100644
---
a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/ProcessDefinitionServiceTest.java
+++
b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/ProcessDefinitionServiceTest.java
@@ -747,7 +747,6 @@ public class ProcessDefinitionServiceTest {
//process definition not exist
ProcessDefinition processDefinition = getProcessDefinition();
processDefinition.setProcessDefinitionJson(SHELL_JSON);
- Mockito.when(processDefineMapper.selectById(46)).thenReturn(null);
Map<String, Object> processDefinitionNullRes =
processDefinitionService.viewTree(46, 10);
Assert.assertEquals(Status.PROCESS_DEFINE_NOT_EXIST,
processDefinitionNullRes.get(Constants.STATUS));
@@ -771,7 +770,7 @@ public class ProcessDefinitionServiceTest {
taskInstance.setHost("192.168.xx.xx");
//task instance not exist
-
Mockito.when(processDefineMapper.selectById(46)).thenReturn(processDefinition);
+
Mockito.when(processDefineMapper.queryByCode(46L)).thenReturn(processDefinition);
Mockito.when(processService.genDagGraph(processDefinition)).thenReturn(new
DAG<>());
Map<String, Object> taskNullRes =
processDefinitionService.viewTree(46, 10);
Assert.assertEquals(Status.SUCCESS, taskNullRes.get(Constants.STATUS));
@@ -783,7 +782,7 @@ public class ProcessDefinitionServiceTest {
}
@Test
- public void testSubProcessViewTree() throws Exception {
+ public void testSubProcessViewTree() {
ProcessDefinition processDefinition = getProcessDefinition();
processDefinition.setProcessDefinitionJson(SHELL_JSON);
@@ -806,7 +805,7 @@ public class ProcessDefinitionServiceTest {
taskInstance.setState(ExecutionStatus.RUNNING_EXECUTION);
taskInstance.setHost("192.168.xx.xx");
taskInstance.setTaskParams("\"processDefinitionId\": \"222\",\n");
-
Mockito.when(processDefineMapper.selectById(46)).thenReturn(processDefinition);
+
Mockito.when(processDefineMapper.queryByCode(46L)).thenReturn(processDefinition);
Mockito.when(processService.genDagGraph(processDefinition)).thenReturn(new
DAG<>());
Map<String, Object> taskNotNuLLRes =
processDefinitionService.viewTree(46, 10);
Assert.assertEquals(Status.SUCCESS,
taskNotNuLLRes.get(Constants.STATUS));
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 9b2dadb..080b474 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
@@ -2428,6 +2428,7 @@ public class ProcessService {
* @param processDefinition process definition
* @return dag graph
*/
+ @Deprecated
public DAG<String, TaskNode, TaskNodeRelation>
genDagGraph(ProcessDefinition processDefinition) {
Map<String, String> locationMap =
locationToMap(processDefinition.getLocations());
List<TaskNode> taskNodeList =
genTaskNodeList(processDefinition.getCode(), processDefinition.getVersion(),
locationMap);
@@ -2539,19 +2540,19 @@ public class ProcessService {
/**
* query tasks definition list by process code and process version
*/
- public List<TaskDefinitionLog> queryTaskDefinitionList(Long processCode,
int processVersion) {
+ public List<TaskDefinitionLog> queryTaskDefinitionListByProcess(long
processCode, int processVersion) {
List<ProcessTaskRelationLog> processTaskRelationLogs =
processTaskRelationLogMapper.queryByProcessCodeAndVersion(processCode,
processVersion);
- Map<Long, TaskDefinition> postTaskDefinitionMap = new HashMap<>();
- processTaskRelationLogs.forEach(processTaskRelationLog -> {
- Long code = processTaskRelationLog.getPostTaskCode();
- int version = processTaskRelationLog.getPostTaskVersion();
- if (postTaskDefinitionMap.containsKey(code)) {
- TaskDefinition taskDefinition =
taskDefinitionLogMapper.queryByDefinitionCodeAndVersion(code, version);
- postTaskDefinitionMap.putIfAbsent(code, taskDefinition);
+ Set<TaskDefinition> taskDefinitionSet = new HashSet<>();
+ for (ProcessTaskRelationLog processTaskRelationLog :
processTaskRelationLogs) {
+ if (processTaskRelationLog.getPreTaskCode() > 0) {
+ taskDefinitionSet.add(new
TaskDefinition(processTaskRelationLog.getPreTaskCode(),
processTaskRelationLog.getPreTaskVersion()));
}
- });
- return new ArrayList(postTaskDefinitionMap.values());
+ if (processTaskRelationLog.getPostTaskCode() > 0) {
+ taskDefinitionSet.add(new
TaskDefinition(processTaskRelationLog.getPostTaskCode(),
processTaskRelationLog.getPostTaskVersion()));
+ }
+ }
+ return
taskDefinitionLogMapper.queryByTaskDefinitions(taskDefinitionSet);
}
/**