This is an automated email from the ASF dual-hosted git repository. leonbao pushed a commit to branch json_split in repository https://gitbox.apache.org/repos/asf/incubator-dolphinscheduler.git
commit f05ba849682c7d79e49ce50bddc45ab149b20e8f Merge: 010b49c 5d264c9 Author: lenboo <[email protected]> AuthorDate: Thu Mar 25 23:54:01 2021 +0800 Merge remote-tracking branch 'upstream/dev' into spilit # Conflicts: # dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/ProcessInstanceMapper.java # dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/TaskInstanceMapperTest.java .../conf/dolphinscheduler/master.properties.tpl | 2 +- .../conf/dolphinscheduler/worker.properties.tpl | 10 +- docker/kubernetes/dolphinscheduler/README.md | 6 +- .../api/controller/WorkerGroupController.java | 74 ++++++-- .../apache/dolphinscheduler/api/enums/Status.java | 4 + .../dolphinscheduler/api/service/UsersService.java | 8 + .../api/service/WorkerGroupService.java | 18 ++ .../api/service/impl/MonitorServiceImpl.java | 2 +- .../service/impl/ProcessInstanceServiceImpl.java | 4 +- .../api/service/impl/UsersServiceImpl.java | 8 + .../api/service/impl/WorkerGroupServiceImpl.java | 199 ++++++++++++++++++--- .../api/utils/ZookeeperMonitor.java | 143 +++++++-------- .../api/service/UsersServiceTest.java | 13 ++ .../api/service/WorkerGroupServiceTest.java | 91 +++++++++- .../apache/dolphinscheduler/common/Constants.java | 4 +- .../dolphinscheduler/common/enums/ZKNodeType.java | 3 +- .../common/utils/CollectionUtils.java | 34 ++++ .../common/utils/CollectionUtilsTest.java | 16 ++ .../dolphinscheduler/dao/entity/WorkerGroup.java | 69 +++++-- .../dao/mapper/ProcessInstanceMapper.java | 15 +- .../dolphinscheduler/dao/mapper/UserMapper.java | 8 + .../dao/mapper/WorkerGroupMapper.java | 33 ++-- .../dao/mapper/ProcessInstanceMapper.xml | 12 +- .../dolphinscheduler/dao/mapper/UserMapper.xml | 8 + .../dao/mapper/WorkerGroupMapper.xml | 31 ++++ .../dao/mapper/TaskInstanceMapperTest.java | 26 ++- .../dao/mapper/UserMapperTest.java | 11 ++ .../master/dispatch/host/CommonHostManager.java | 92 +++++++--- .../master/dispatch/host/HostManagerConfig.java | 2 +- .../dispatch/host/LowerWeightHostManager.java | 64 +++++-- .../server/master/registry/MasterRegistry.java | 4 +- .../server/registry/HeartBeatTask.java | 2 +- .../server/registry/ZookeeperRegistryCenter.java | 6 +- .../server/worker/config/WorkerConfig.java | 10 +- .../server/worker/registry/WorkerRegistry.java | 4 +- .../dolphinscheduler/server/zk/ZKMasterClient.java | 8 +- .../src/main/resources/master.properties | 2 +- .../src/main/resources/worker.properties | 10 +- .../consumer/TaskPriorityQueueConsumerTest.java | 10 +- .../service/zk/AbstractZKClient.java | 52 ++++-- .../service/zk/RegisterOperator.java | 8 +- .../service/zk/RegisterOperatorTest.java | 8 +- .../home/pages/dag/_source/formModel/tasks/sql.vue | 4 +- .../conf/home/pages/security/pages/queue/index.vue | 2 - .../home/pages/security/pages/tenement/index.vue | 2 - .../pages/security/pages/users/_source/list.vue | 2 +- .../conf/home/pages/security/pages/users/index.vue | 3 - .../security/pages/warningGroups/_source/list.vue | 2 +- .../pages/security/pages/warningGroups/index.vue | 3 - .../pages/security/pages/warningInstance/index.vue | 3 - .../pages/workerGroups/_source/createWorker.vue | 67 ++++--- .../security/pages/workerGroups/_source/list.vue | 23 ++- .../pages/security/pages/workerGroups/index.vue | 31 +++- .../src/js/conf/home/store/security/actions.js | 2 +- .../src/js/module/i18n/locale/en_US.js | 16 +- .../src/js/module/i18n/locale/zh_CN.js | 14 +- sql/dolphinscheduler_mysql.sql | 16 ++ sql/dolphinscheduler_postgre.sql | 5 +- .../1.3.6_schema/mysql/dolphinscheduler_ddl.sql | 35 ++-- .../1.3.6_schema/mysql/dolphinscheduler_dml.sql | 17 +- .../postgresql/dolphinscheduler_ddl.sql | 40 +++++ .../postgresql/dolphinscheduler_dml.sql | 17 +- 62 files changed, 1040 insertions(+), 398 deletions(-) diff --cc dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/ProcessInstanceServiceImpl.java index 95e3f5c,f8a9250..54d07f7 --- a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/ProcessInstanceServiceImpl.java +++ b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/ProcessInstanceServiceImpl.java @@@ -257,13 -247,13 +257,15 @@@ public class ProcessInstanceServiceImp PageInfo<ProcessInstance> pageInfo = new PageInfo<>(pageNo, pageSize); int executorId = usersService.getUserIdByName(executorName); + ProcessDefinition processDefinition = processDefineMapper.queryByDefineId(processDefineId); + IPage<ProcessInstance> processInstanceList = processInstanceMapper.queryProcessInstanceListPaging(page, - project.getId(), processDefineId, searchVal, executorId, statusArray, host, start, end); + project.getCode(), processDefinition.getCode(), searchVal, executorId, statusArray, host, start, end); List<ProcessInstance> processInstances = processInstanceList.getRecords(); + List<Integer> userIds = CollectionUtils.transformToList(processInstances, ProcessInstance::getExecutorId); + Map<Integer, User> idToUserMap = CollectionUtils.collectionToMap(usersService.queryUser(userIds), User::getId); for (ProcessInstance processInstance : processInstances) { processInstance.setDuration(DateUtils.format2Duration(processInstance.getStartTime(), processInstance.getEndTime())); diff --cc dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/ProcessInstanceMapper.java index 6e53b45,08d1740..7be58a7 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/ProcessInstanceMapper.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/ProcessInstanceMapper.java @@@ -61,12 -58,11 +61,10 @@@ public interface ProcessInstanceMapper * @return process instance list */ List<ProcessInstance> queryByTenantIdAndStatus(@Param("tenantId") int tenantId, - @Param("states") int[] states); + @Param("states") int[] states); /** -- * query process instance by worker group and stateArray - * - * @param workerGroupId workerGroupId + * @param workerGroupName workerGroupName * @param states states array * @return process instance list */ @@@ -142,13 -134,13 +140,14 @@@ @Param("destTenantId") int destTenantId); /** - * update process instance by worker group name + * update process instance by worker groupId + * - * @param originWorkerGroupId originWorkerGroupId - * @param destWorkerGroupId destWorkerGroupId + * @param originWorkerGroupName originWorkerGroupName + * @param destWorkerGroupName destWorkerGroupName * @return update result */ - int updateProcessInstanceByWorkerGroupId(@Param("originWorkerGroupId") int originWorkerGroupId, @Param("destWorkerGroupId") int destWorkerGroupId); + int updateProcessInstanceByWorkerGroupName(@Param("originWorkerGroupName") String originWorkerGroupName, + @Param("destWorkerGroupName") String destWorkerGroupName); /** * count process instance state by user diff --cc dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/TaskInstanceMapperTest.java index a0930fb,9ad8677..6dc348c --- a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/TaskInstanceMapperTest.java +++ b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/TaskInstanceMapperTest.java @@@ -17,7 -17,10 +17,9 @@@ package org.apache.dolphinscheduler.dao.mapper; - import org.apache.dolphinscheduler.common.enums.CommandType; + import com.baomidou.mybatisplus.core.metadata.IPage; + import com.baomidou.mybatisplus.extension.plugins.pagination.Page; + -import org.apache.dolphinscheduler.common.enums.CommandType; import org.apache.dolphinscheduler.common.enums.ExecutionStatus; import org.apache.dolphinscheduler.common.enums.Flag; import org.apache.dolphinscheduler.common.enums.TaskType; @@@ -41,12 -45,10 +44,9 @@@ import org.springframework.transaction. @RunWith(SpringRunner.class) @SpringBootTest @Transactional - @Rollback(true) + @Rollback public class TaskInstanceMapperTest { - @Autowired TaskInstanceMapper taskInstanceMapper; @@@ -99,7 -95,7 +110,8 @@@ taskInstance.setTaskJson("{}"); taskInstance.setProcessInstanceId(processInstanceId); taskInstance.setTaskType(taskType); - taskInstance.setProcessDefinitionId(processDefinitionId); + taskInstance.setProcessDefinitionCode(1L); ++// taskInstance.setProcessDefinitionId(processDefinitionId); taskInstanceMapper.insert(taskInstance); return taskInstance; } @@@ -337,24 -287,22 +349,20 @@@ */ @Test public void testQueryTaskInstanceListPaging() { - ProcessDefinition definition = new ProcessDefinition(); + definition.setCode(1L); definition.setProjectId(1111); + definition.setProjectCode(1111L); + definition.setCreateTime(new Date()); + definition.setUpdateTime(new Date()); processDefinitionMapper.insert(definition); - ProcessInstance processInstance = new ProcessInstance(); - processInstance.setProcessDefinitionId(definition.getId()); - processInstance.setState(ExecutionStatus.RUNNING_EXECUTION); - processInstance.setName("ut process"); - processInstance.setStartTime(new Date()); - processInstance.setEndTime(new Date()); - processInstance.setCommandType(CommandType.START_PROCESS); - processInstanceMapper.insert(processInstance); + // insert ProcessInstance + ProcessInstance processInstance = insertProcessInstance(); - TaskInstance task = insertOne("us task", processInstance.getId(), ExecutionStatus.RUNNING_EXECUTION, TaskType.SHELL.toString(),definition.getId()); + // insert taskInstance + TaskInstance task = insertTaskInstance(processInstance.getId()); - task.setProcessDefinitionId(definition.getId()); - task.setProcessInstanceId(processInstance.getId()); - taskInstanceMapper.updateById(task); - Page<TaskInstance> page = new Page(1, 3); IPage<TaskInstance> taskInstanceIPage = taskInstanceMapper.queryTaskInstanceListPaging( page,
