This is an automated email from the ASF dual-hosted git repository. ruanwenjun pushed a commit to branch revert-18153-Improvement-18056-dao-syc in repository https://gitbox.apache.org/repos/asf/dolphinscheduler.git
commit 98ea98b570d9a734c3f55dfd9d4267117e2420eb Author: Wenjun Ruan <[email protected]> AuthorDate: Fri May 1 12:17:49 2026 +0800 Revert "[Improvement-18056] Clean up unused methods and classes in the dolphi…" This reverts commit f80a62e98ce1638998439aceea779c3f26bdd341. --- .../dao/entity/DependentWorkflowDefinition.java | 50 +++++++++++++ .../dao/mapper/AlertGroupMapper.java | 9 +++ .../dao/mapper/DataSourceMapper.java | 9 +++ .../dao/mapper/K8sNamespaceMapper.java | 15 ++++ .../dao/mapper/K8sNamespaceUserMapper.java | 9 +++ .../dolphinscheduler/dao/mapper/ProjectMapper.java | 21 ++++++ .../dao/mapper/RelationSubWorkflowMapper.java | 2 + .../dao/mapper/ScheduleMapper.java | 10 +++ .../dao/mapper/TaskDefinitionMapper.java | 21 ++++++ .../dao/mapper/TaskGroupQueueMapper.java | 35 +++++++++ .../dao/mapper/TaskInstanceMapper.java | 37 ++++++++++ .../dolphinscheduler/dao/mapper/TenantMapper.java | 11 +++ .../dao/mapper/WorkerGroupMapper.java | 5 ++ .../dao/mapper/WorkflowDefinitionLogMapper.java | 4 +- .../dao/mapper/WorkflowDefinitionMapper.java | 24 +++++++ .../dao/mapper/WorkflowInstanceMapper.java | 69 ++++++++++++++++++ .../dao/mapper/WorkflowTaskRelationLogMapper.java | 17 +++++ .../dao/mapper/WorkflowTaskRelationMapper.java | 53 ++++++++++++++ .../dolphinscheduler/dao/repository/BaseDao.java | 25 +++++++ .../dolphinscheduler/dao/repository/IDao.java | 15 ++++ .../dao/repository/TaskDefinitionDao.java | 10 +++ .../dao/repository/TaskInstanceContextDao.java | 3 + .../dao/repository/TaskInstanceDao.java | 30 ++++++++ .../dao/repository/WorkerGroupDao.java | 2 + .../dao/repository/WorkflowInstanceDao.java | 2 + .../dao/repository/WorkflowTaskLineageDao.java | 1 + .../dao/repository/impl/TaskDefinitionDaoImpl.java | 10 +++ .../impl/TaskInstanceContextDaoImpl.java | 9 +++ .../dao/repository/impl/TaskInstanceDaoImpl.java | 84 ++++++++++++++++++++++ .../dao/repository/impl/WorkerGroupDaoImpl.java | 6 ++ .../repository/impl/WorkflowInstanceDaoImpl.java | 13 ++++ .../impl/WorkflowTaskLineageDaoImpl.java | 11 +++ .../dao/mapper/AlertMapperTest.java | 20 ++++++ .../dao/mapper/AuditLogMapperTest.java | 15 ++++ .../dao/mapper/CommandMapperTest.java | 20 ++++++ .../dao/mapper/ProjectWorkerGroupMapperTest.java | 1 + .../dao/mapper/TaskDefinitionLogMapperTest.java | 1 + .../dao/mapper/TaskGroupQueueMapperTest.java | 2 + .../dao/mapper/TenantMapperTest.java | 1 - .../dao/mapper/UserMapperTest.java | 18 +++++ .../dao/mapper/WorkflowInstanceMapMapperTest.java | 1 + .../dao/mapper/WorkflowInstanceMapperTest.java | 3 + .../mapper/WorkflowTaskRelationLogMapperTest.java | 2 +- .../dao/mapper/WorkflowTaskRelationMapperTest.java | 8 +-- 44 files changed, 705 insertions(+), 9 deletions(-) diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/entity/DependentWorkflowDefinition.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/entity/DependentWorkflowDefinition.java index 671e065312..01e48e4669 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/entity/DependentWorkflowDefinition.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/entity/DependentWorkflowDefinition.java @@ -17,6 +17,14 @@ package org.apache.dolphinscheduler.dao.entity; +import org.apache.dolphinscheduler.common.enums.CycleEnum; +import org.apache.dolphinscheduler.common.utils.JSONUtils; +import org.apache.dolphinscheduler.plugin.task.api.model.DependentItem; +import org.apache.dolphinscheduler.plugin.task.api.model.DependentTaskModel; +import org.apache.dolphinscheduler.plugin.task.api.parameters.DependentParameters; + +import java.util.List; + import lombok.Data; @Data @@ -32,4 +40,46 @@ public class DependentWorkflowDefinition { private String workerGroup; + public CycleEnum getDependentCycle(long upstreamProcessDefinitionCode) { + DependentParameters dependentParameters = this.getDependentParameters(); + List<DependentTaskModel> dependentTaskModelList = dependentParameters.getDependence().getDependTaskList(); + + for (DependentTaskModel dependentTaskModel : dependentTaskModelList) { + List<DependentItem> dependentItemList = dependentTaskModel.getDependItemList(); + for (DependentItem dependentItem : dependentItemList) { + if (upstreamProcessDefinitionCode == dependentItem.getDefinitionCode()) { + return cycle2CycleEnum(dependentItem.getCycle()); + } + } + } + + return CycleEnum.DAY; + } + + public CycleEnum cycle2CycleEnum(String cycle) { + CycleEnum cycleEnum = null; + + switch (cycle) { + case "day": + cycleEnum = CycleEnum.DAY; + break; + case "hour": + cycleEnum = CycleEnum.HOUR; + break; + case "week": + cycleEnum = CycleEnum.WEEK; + break; + case "month": + cycleEnum = CycleEnum.MONTH; + break; + default: + break; + } + return cycleEnum; + } + + public DependentParameters getDependentParameters() { + return JSONUtils.parseObject(taskParams, DependentParameters.class); + } + } diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/AlertGroupMapper.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/AlertGroupMapper.java index d1d1b57c4c..712b67e5f2 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/AlertGroupMapper.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/AlertGroupMapper.java @@ -81,6 +81,15 @@ public interface AlertGroupMapper extends BaseMapper<AlertGroup> { */ String queryAlertGroupInstanceIdsById(@Param("alertGroupId") int alertGroupId); + /** + * list authorized AlertGroup + * @param userId + * @param alertGroupsIds + * @return + */ + <T> List<AlertGroup> listAuthorizedAlertGroupList(@Param("userId") int userId, + @Param("alertGroupsIds") List<Integer> alertGroupsIds); + /** * queryAlertGroupPageByIds * @param page diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/DataSourceMapper.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/DataSourceMapper.java index e06374c32a..0eb87a9dcb 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/DataSourceMapper.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/DataSourceMapper.java @@ -90,6 +90,15 @@ public interface DataSourceMapper extends BaseMapper<DataSource> { <T> List<DataSource> listAuthorizedDataSource(@Param("userId") int userId, @Param("dataSourceIds") T[] dataSourceIds); + /** + * query datasource by name and user id + * + * @param userId userId + * @param name datasource name + * @return If the name does not exist or the user does not have permission, it will return null + */ + DataSource queryDataSourceByNameAndUserId(@Param("userId") int userId, @Param("name") String name); + /** * selectPagingByIds * @param dataSourcePage diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/K8sNamespaceMapper.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/K8sNamespaceMapper.java index f9f92c7493..6c773c3e1d 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/K8sNamespaceMapper.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/K8sNamespaceMapper.java @@ -50,6 +50,14 @@ public interface K8sNamespaceMapper extends BaseMapper<K8sNamespace> { */ Boolean existNamespace(@Param("namespace") String namespace, @Param("clusterCode") Long clusterCode); + /** + * query namespace except userId + * + * @param userId userId + * @return namespace list + */ + List<K8sNamespace> queryNamespaceExceptUserId(@Param("userId") int userId); + /** * query authed namespace list by userId * @@ -58,4 +66,11 @@ public interface K8sNamespaceMapper extends BaseMapper<K8sNamespace> { */ List<K8sNamespace> queryAuthedNamespaceListByUserId(@Param("userId") Integer userId); + /** + * check the target namespace + * + * @param namespaceCode namespaceCode + * @return true if exist else return null + */ + K8sNamespace queryByNamespaceCode(@Param("clusterCode") Long namespaceCode); } diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/K8sNamespaceUserMapper.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/K8sNamespaceUserMapper.java index 013fe1746d..46723cc41d 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/K8sNamespaceUserMapper.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/K8sNamespaceUserMapper.java @@ -38,4 +38,13 @@ public interface K8sNamespaceUserMapper extends BaseMapper<K8sNamespaceUser> { int deleteNamespaceRelation(@Param("namespaceId") int namespaceId, @Param("userId") int userId); + /** + * query namespace relation + * + * @param namespaceId namespaceId + * @param userId userId + * @return namespace user relation + */ + K8sNamespaceUser queryNamespaceRelation(@Param("namespaceId") int namespaceId, + @Param("userId") int userId); } diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/ProjectMapper.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/ProjectMapper.java index cef148f15c..566e1beb7d 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/ProjectMapper.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/ProjectMapper.java @@ -52,6 +52,13 @@ public interface ProjectMapper extends BaseMapper<Project> { */ Project queryDetailById(@Param("projectId") int projectId); + /** + * query project detail by code + * @param projectCode projectCode + * @return project + */ + Project queryDetailByCode(@Param("projectCode") long projectCode); + /** * query project by name * @param projectName projectName @@ -84,6 +91,13 @@ public interface ProjectMapper extends BaseMapper<Project> { */ List<Project> queryAuthedProjectListByUserId(@Param("userId") int userId); + /** + * query relation project list by userId + * @param userId userId + * @return project list + */ + List<Project> queryRelationProjectListByUserId(@Param("userId") int userId); + /** * query project except userId * @param userId userId @@ -116,6 +130,7 @@ public interface ProjectMapper extends BaseMapper<Project> { * list authorized Projects * @param userId * @param projectsIds + * @param <T> * @return */ List<Project> listAuthorizedProjects(@Param("userId") int userId, @Param("projectsIds") List<Integer> projectsIds); @@ -133,4 +148,10 @@ public interface ProjectMapper extends BaseMapper<Project> { */ Project queryProjectByTaskInstanceId(@Param("taskInstanceId") int taskInstanceId); + /** + * query all workflow count + * @param projectsCodes projectsCodes + * @return workflow count + */ + int queryAllWorkflowCounts(@Param("projectsCodes") List<Long> projectsCodes); } diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/RelationSubWorkflowMapper.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/RelationSubWorkflowMapper.java index a1ef713417..78fcdd89ea 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/RelationSubWorkflowMapper.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/RelationSubWorkflowMapper.java @@ -32,4 +32,6 @@ public interface RelationSubWorkflowMapper extends BaseMapper<RelationSubWorkflo List<RelationSubWorkflow> queryAllSubWorkflowInstance(@Param("parentWorkflowInstanceId") Long parentWorkflowInstanceId, @Param("parentTaskCode") Long parentTaskCode); + RelationSubWorkflow queryParentWorkflowInstance(@Param("subWorkflowInstanceId") Long subWorkflowInstanceId); + } diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/ScheduleMapper.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/ScheduleMapper.java index a3cc7a9b12..82afa56553 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/ScheduleMapper.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/ScheduleMapper.java @@ -69,6 +69,16 @@ public interface ScheduleMapper extends BaseMapper<Schedule> { @Param("workflowDefinitionCode") long workflowDefinitionCode, @Param("searchVal") String searchVal); + /** + * Filter schedule + * + * @param page page + * @param schedule schedule + * @return schedule IPage + */ + IPage<Schedule> filterSchedules(IPage<Schedule> page, + @Param("schedule") Schedule schedule); + /** * query schedule list by project name * diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/TaskDefinitionMapper.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/TaskDefinitionMapper.java index b43e519df4..d9e322192c 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/TaskDefinitionMapper.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/TaskDefinitionMapper.java @@ -19,6 +19,7 @@ package org.apache.dolphinscheduler.dao.mapper; import org.apache.dolphinscheduler.dao.entity.TaskDefinition; import org.apache.dolphinscheduler.dao.entity.TaskDefinitionLog; +import org.apache.dolphinscheduler.dao.entity.TaskMainInfo; import org.apache.dolphinscheduler.dao.model.WorkflowDefinitionCountDto; import org.apache.ibatis.annotations.Param; @@ -27,6 +28,7 @@ import java.util.Collection; import java.util.List; import com.baomidou.mybatisplus.core.mapper.BaseMapper; +import com.baomidou.mybatisplus.core.metadata.IPage; public interface TaskDefinitionMapper extends BaseMapper<TaskDefinition> { @@ -82,6 +84,15 @@ public interface TaskDefinitionMapper extends BaseMapper<TaskDefinition> { */ int batchInsert(@Param("taskDefinitions") List<TaskDefinitionLog> taskDefinitions); + /** + * task main info + * @param projectCode project code + * @param codeList code list + * @return task main info + */ + List<TaskMainInfo> queryDefineListByCodeList(@Param("projectCode") long projectCode, + @Param("codeList") List<Long> codeList); + /** * query task definition by code list * @@ -90,6 +101,16 @@ public interface TaskDefinitionMapper extends BaseMapper<TaskDefinition> { */ List<TaskDefinition> queryByCodeList(@Param("codes") Collection<Long> codes); + /** + * Filter task definition + * + * @param page page + * @param taskDefinition task definition + * @return task definition IPage + */ + IPage<TaskDefinition> filterTaskDefinition(IPage<TaskDefinition> page, + @Param("task") TaskDefinition taskDefinition); + /** * batch delete task by task code * diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/TaskGroupQueueMapper.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/TaskGroupQueueMapper.java index 20c9bc5e64..c4e6495e90 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/TaskGroupQueueMapper.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/TaskGroupQueueMapper.java @@ -36,6 +36,18 @@ import com.baomidou.mybatisplus.extension.plugins.pagination.Page; */ public interface TaskGroupQueueMapper extends BaseMapper<TaskGroupQueue> { + /** + * select task group queues by some conditions + * + * @param page page + * @param groupId group id + * @return task group queue list + */ + IPage<TaskGroupQueue> queryTaskGroupQueuePaging(IPage<TaskGroupQueue> page, + @Param("groupId") int groupId); + + TaskGroupQueue queryByTaskId(@Param("taskId") int taskId); + /** * query by status * @@ -61,6 +73,24 @@ public interface TaskGroupQueueMapper extends BaseMapper<TaskGroupQueue> { */ int updateStatusByTaskId(@Param("taskId") int taskId, @Param("status") int status); + /** + * Query the {@link TaskGroupQueue}, who's priority > the given <code>priority</code> + */ + List<TaskGroupQueue> queryHighPriorityTasks(@Param("groupId") int groupId, @Param("priority") int priority, + @Param("status") int status); + + TaskGroupQueue queryTheHighestPriorityTasks(@Param("groupId") int groupId, @Param("status") int status, + @Param("forceStart") int forceStart, @Param("inQueue") int inQueue); + + void updateInQueue(@Param("inQueue") int inQueue, @Param("id") int id); + + void updateForceStart(@Param("queueId") int queueId, @Param("forceStart") int forceStart); + + int updateInQueueLimit1(@Param("oldValue") int oldValue, @Param("newValue") int newValue, @Param("groupId") int id, + @Param("status") int status); + + int updateInQueueCAS(@Param("oldValue") int oldValue, @Param("newValue") int newValue, @Param("id") int id); + void modifyPriority(@Param("queueId") int queueId, @Param("priority") int priority); IPage<TaskGroupQueue> queryTaskGroupQueueByTaskGroupIdPaging(Page<TaskGroupQueue> page, @@ -70,12 +100,17 @@ public interface TaskGroupQueueMapper extends BaseMapper<TaskGroupQueue> { @Param("groupId") int groupId, @Param("projects") List<Project> projects); + void deleteByTaskInstanceIds(@Param("taskInstanceIds") List<Integer> taskInstanceIds); + void deleteByWorkflowInstanceId(@Param("workflowInstanceId") Integer workflowInstanceId); void deleteByWorkflowInstanceIds(@Param("workflowInstanceIds") List<Integer> workflowInstanceIds); void deleteByTaskGroupIds(@Param("taskGroupIds") List<Integer> taskGroupIds); + void updateTaskGroupPriorityByTaskInstanceId(@Param("taskInstanceId") Integer taskInstanceId, + @Param("priority") int taskGroupPriority); + List<TaskGroupQueue> queryAllInQueueTaskGroupQueueByGroupId(@Param("taskGroupId") Integer taskGroupId, @Param("inQueue") int inQueue); diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/TaskInstanceMapper.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/TaskInstanceMapper.java index 02811da19a..32a05a21ec 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/TaskInstanceMapper.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/TaskInstanceMapper.java @@ -19,8 +19,10 @@ package org.apache.dolphinscheduler.dao.mapper; import org.apache.dolphinscheduler.common.enums.Flag; import org.apache.dolphinscheduler.common.enums.TaskExecuteType; +import org.apache.dolphinscheduler.dao.entity.ExecuteStatusCount; import org.apache.dolphinscheduler.dao.entity.TaskInstance; import org.apache.dolphinscheduler.dao.model.TaskInstanceStatusCountDto; +import org.apache.dolphinscheduler.plugin.task.api.enums.TaskExecutionStatus; import org.apache.ibatis.annotations.Param; @@ -57,6 +59,41 @@ public interface TaskInstanceMapper extends BaseMapper<TaskInstance> { @Param("endTime") Date endTime, @Param("projectCodes") Collection<Long> projectCodes); + /** + * Statistics task instance group by given project ids list by start time + * <p> + * We only need project ids to determine whether the task instance belongs to the user or not. + * + * @param startTime Statistics start time + * @param endTime Statistics end time + * @param projectIds Project ids list to filter + * @return List of ExecuteStatusCount + */ + List<ExecuteStatusCount> countTaskInstanceStateByProjectIdsV2(@Param("startTime") Date startTime, + @Param("endTime") Date endTime, + @Param("projectIds") Set<Integer> projectIds); + + /** + * Statistics task instance group by given project codes list by submit time + * <p> + * We only need project codes to determine whether the task instance belongs to the user or not. + * + * @param startTime Statistics start time + * @param endTime Statistics end time + * @param projectCode projectCode + * @param model model + * @param projectIds projectIds + * @return List of ExecuteStatusCount + */ + List<ExecuteStatusCount> countTaskInstanceStateByProjectCodesAndStatesBySubmitTimeV2(@Param("startTime") Date startTime, + @Param("endTime") Date endTime, + @Param("projectCode") Long projectCode, + @Param("workflowCode") Long workflowCode, + @Param("taskCode") Long taskCode, + @Param("model") Integer model, + @Param("projectIds") Set<Integer> projectIds, + @Param("states") List<TaskExecutionStatus> states); + IPage<TaskInstance> queryTaskInstanceListPaging(IPage<TaskInstance> page, @Param("projectCode") Long projectCode, @Param("workflowInstanceId") Integer workflowInstanceId, diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/TenantMapper.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/TenantMapper.java index 2b21ea6743..69b904f57d 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/TenantMapper.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/TenantMapper.java @@ -25,6 +25,7 @@ import java.util.List; import com.baomidou.mybatisplus.core.mapper.BaseMapper; import com.baomidou.mybatisplus.core.metadata.IPage; +import com.baomidou.mybatisplus.extension.plugins.pagination.Page; public interface TenantMapper extends BaseMapper<Tenant> { @@ -79,6 +80,16 @@ public interface TenantMapper extends BaseMapper<Tenant> { */ Boolean existTenant(@Param("tenantCode") String tenantCode); + /** + * queryTenantPagingByIds + * @param page + * @param ids + * @param searchVal + * @return + */ + IPage<Tenant> queryTenantPagingByIds(Page<Tenant> page, @Param("ids") List<Integer> ids, + @Param("searchVal") String searchVal); + /** * queryAll * @return diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/WorkerGroupMapper.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/WorkerGroupMapper.java index 5679921608..06a39223ef 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/WorkerGroupMapper.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/WorkerGroupMapper.java @@ -17,6 +17,7 @@ package org.apache.dolphinscheduler.dao.mapper; +import org.apache.dolphinscheduler.common.enums.WorkerGroupSource; import org.apache.dolphinscheduler.dao.entity.WorkerGroup; import org.apache.ibatis.annotations.Param; @@ -48,5 +49,9 @@ public interface WorkerGroupMapper extends BaseMapper<WorkerGroup> { */ List<WorkerGroup> queryWorkerGroupByName(@Param("name") String name); + int updateAddrListByWorkerGroupName(@Param("name") String name, + @Param("addrList") String addrList, + @Param("source") WorkerGroupSource source); + int deleteByWorkerGroupName(@Param("name") String name); } diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/WorkflowDefinitionLogMapper.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/WorkflowDefinitionLogMapper.java index d8c1719366..58b4c29d7f 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/WorkflowDefinitionLogMapper.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/WorkflowDefinitionLogMapper.java @@ -27,9 +27,6 @@ import com.baomidou.mybatisplus.core.mapper.BaseMapper; import com.baomidou.mybatisplus.core.metadata.IPage; import com.baomidou.mybatisplus.extension.plugins.pagination.Page; -/** - * workflow definition log mapper interface - */ public interface WorkflowDefinitionLogMapper extends BaseMapper<WorkflowDefinitionLog> { /** @@ -39,6 +36,7 @@ public interface WorkflowDefinitionLogMapper extends BaseMapper<WorkflowDefiniti * @param version version number * @return the workflow definition version info */ + WorkflowDefinitionLog queryByDefinitionCodeAndVersion(@Param("code") long code, @Param("version") int version); /** diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/WorkflowDefinitionMapper.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/WorkflowDefinitionMapper.java index 36b42d4b01..6fa7e6e50f 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/WorkflowDefinitionMapper.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/WorkflowDefinitionMapper.java @@ -103,6 +103,16 @@ public interface WorkflowDefinitionMapper extends BaseMapper<WorkflowDefinition> @Param("userId") int userId, @Param("projectCode") long projectCode); + /** + * Filter workflow definitions + * + * @param page page + * @param workflowDefinition workflow definition object + * @return workflow definition IPage + */ + IPage<WorkflowDefinition> filterWorkflowDefinition(IPage<WorkflowDefinition> page, + @Param("pd") WorkflowDefinition workflowDefinition); + /** * query all workflow definition list * @@ -138,6 +148,20 @@ public interface WorkflowDefinitionMapper extends BaseMapper<WorkflowDefinition> */ List<WorkflowDefinitionCountDto> countDefinitionByProjectCodes(@Param("projectCodes") Collection<Long> projectCodes); + /** + * Statistics workflow definition group by project codes list + * <p> + * We only need project codes to determine whether the definition belongs to the user or not. + * + * @param projectCodes projectCodes + * @param userId userId + * @param releaseState releaseState + * @return definition group by user + */ + List<WorkflowDefinitionCountDto> countDefinitionByProjectCodesV2(@Param("projectCodes") List<Long> projectCodes, + @Param("userId") Integer userId, + @Param("releaseState") Integer releaseState); + /** * list all project ids * diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/WorkflowInstanceMapper.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/WorkflowInstanceMapper.java index 5ffe392917..aa88680ae7 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/WorkflowInstanceMapper.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/WorkflowInstanceMapper.java @@ -18,6 +18,7 @@ package org.apache.dolphinscheduler.dao.mapper; import org.apache.dolphinscheduler.common.enums.WorkflowExecutionStatus; +import org.apache.dolphinscheduler.dao.entity.ExecuteStatusCount; import org.apache.dolphinscheduler.dao.entity.WorkflowInstance; import org.apache.dolphinscheduler.dao.model.WorkflowInstanceStatusCountDto; @@ -26,6 +27,7 @@ import org.apache.ibatis.annotations.Param; import java.util.Collection; import java.util.Date; import java.util.List; +import java.util.Set; import com.baomidou.mybatisplus.core.mapper.BaseMapper; import com.baomidou.mybatisplus.core.metadata.IPage; @@ -77,6 +79,20 @@ public interface WorkflowInstanceMapper extends BaseMapper<WorkflowInstance> { List<WorkflowInstance> queryByWorkerGroupNameAndStatus(@Param("workerGroupName") String workerGroupName, @Param("states") int[] states); + /** + * workflow instance page + * @param page page + * @param projectId projectId + * @param processDefinitionId processDefinitionId + * @param searchVal searchVal + * @param executorId executorId + * @param statusArray statusArray + * @param host host + * @param startTime startTime + * @param endTime endTime + * @return workflow instance IPage + */ + /** * workflow instance page * @@ -101,6 +117,16 @@ public interface WorkflowInstanceMapper extends BaseMapper<WorkflowInstance> { @Param("startTime") Date startTime, @Param("endTime") Date endTime); + /** + * set failover by host and state array + * + * @param host host + * @param stateArray stateArray + * @return set result + */ + int setFailoverByHostAndStateArray(@Param("host") String host, + @Param("states") int[] stateArray); + /** * Update the workflow instance state from originState to destState */ @@ -215,6 +241,7 @@ public interface WorkflowInstanceMapper extends BaseMapper<WorkflowInstance> { * @param projectCode project code * @return ProcessInstance list */ + List<WorkflowInstance> queryTopNWorkflowInstance(@Param("size") int size, @Param("startTime") Date startTime, @Param("endTime") Date endTime, @@ -228,6 +255,7 @@ public interface WorkflowInstanceMapper extends BaseMapper<WorkflowInstance> { * @param states states array * @return workflow instance list */ + List<WorkflowInstance> queryByWorkflowDefinitionCodeAndStatus(@Param("workflowDefinitionCode") Long workflowDefinitionCode, @Param("states") int[] states); @@ -235,6 +263,47 @@ public interface WorkflowInstanceMapper extends BaseMapper<WorkflowInstance> { @Param("workflowDefinitionVersion") int workflowDefinitionVersion, @Param("states") int[] states); + /** + * Filter workflow instance + * + * @param page page + * @param workflowDefinitionCode workflowDefinitionCode + * @param name name + * @param host host + * @param startTime startTime + * @param endTime endTime + * @return workflow instance IPage + */ + IPage<WorkflowInstance> queryWorkflowInstanceListV2Paging(Page<WorkflowInstance> page, + @Param("projectCode") Long projectCode, + @Param("workflowDefinitionCode") Long workflowDefinitionCode, + @Param("name") String name, + @Param("startTime") String startTime, + @Param("endTime") String endTime, + @Param("state") Integer state, + @Param("host") String host); + + /** + * Statistics workflow instance state v2 + * <p> + * We only need project codes to determine whether the workflow instance belongs to the user or not. + * + * @param startTime startTime + * @param endTime endTime + * @param projectCode projectCode + * @param workflowCode workflowCode + * @param model model + * @param projectIds projectIds + * @return ExecuteStatusCount list + */ + List<ExecuteStatusCount> countInstanceStateV2( + @Param("startTime") Date startTime, + @Param("endTime") Date endTime, + @Param("projectCode") Long projectCode, + @Param("workflowCode") Long workflowCode, + @Param("model") Integer model, + @Param("projectIds") Set<Integer> projectIds); + /** * query process list by triggerCode * diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/WorkflowTaskRelationLogMapper.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/WorkflowTaskRelationLogMapper.java index 8c439d5536..82ea32f4b2 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/WorkflowTaskRelationLogMapper.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/WorkflowTaskRelationLogMapper.java @@ -17,6 +17,7 @@ package org.apache.dolphinscheduler.dao.mapper; +import org.apache.dolphinscheduler.dao.entity.WorkflowTaskRelation; import org.apache.dolphinscheduler.dao.entity.WorkflowTaskRelationLog; import org.apache.ibatis.annotations.Param; @@ -55,6 +56,22 @@ public interface WorkflowTaskRelationLogMapper extends BaseMapper<WorkflowTaskRe int deleteByCode(@Param("workflowDefinitionCode") long workflowDefinitionCode, @Param("workflowDefinitionVersion") int workflowDefinitionVersion); + /** + * delete workflow task relation + * + * @param workflowTaskRelationLog workflowTaskRelationLog + * @return int + */ + int deleteRelation(@Param("workflowTaskRelationLog") WorkflowTaskRelationLog workflowTaskRelationLog); + + /** + * query workflow task relation log + * + * @param workflowTaskRelation workflowTaskRelation + * @return workflow task relation log + */ + WorkflowTaskRelationLog queryRelationLogByRelation(@Param("workflowTaskRelation") WorkflowTaskRelation workflowTaskRelation); + List<WorkflowTaskRelationLog> queryByWorkflowDefinitionCode(@Param("workflowDefinitionCode") long workflowDefinitionCode); void deleteByWorkflowDefinitionCode(@Param("workflowDefinitionCode") long workflowDefinitionCode); diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/WorkflowTaskRelationMapper.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/WorkflowTaskRelationMapper.java index bba992615c..a733aeaff9 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/WorkflowTaskRelationMapper.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/mapper/WorkflowTaskRelationMapper.java @@ -18,12 +18,14 @@ package org.apache.dolphinscheduler.dao.mapper; import org.apache.dolphinscheduler.dao.entity.WorkflowTaskRelation; +import org.apache.dolphinscheduler.dao.entity.WorkflowTaskRelationLog; import org.apache.ibatis.annotations.Param; import java.util.List; import com.baomidou.mybatisplus.core.mapper.BaseMapper; +import com.baomidou.mybatisplus.core.metadata.IPage; public interface WorkflowTaskRelationMapper extends BaseMapper<WorkflowTaskRelation> { @@ -74,6 +76,14 @@ public interface WorkflowTaskRelationMapper extends BaseMapper<WorkflowTaskRelat */ int batchInsert(@Param("taskRelationList") List<WorkflowTaskRelation> taskRelationList); + /** + * query downstream workflow task relation by taskCode + * + * @param taskCode taskCode + * @return ProcessTaskRelation + */ + List<WorkflowTaskRelation> queryDownstreamByTaskCode(@Param("taskCode") long taskCode); + /** * query upstream workflow task relation by taskCode * @@ -84,6 +94,28 @@ public interface WorkflowTaskRelationMapper extends BaseMapper<WorkflowTaskRelat List<WorkflowTaskRelation> queryUpstreamByCode(@Param("projectCode") long projectCode, @Param("taskCode") long taskCode); + /** + * query downstream workflow task relation by taskCode + * + * @param projectCode projectCode + * @param taskCode taskCode + * @return ProcessTaskRelation + */ + List<WorkflowTaskRelation> queryDownstreamByCode(@Param("projectCode") long projectCode, + @Param("taskCode") long taskCode); + + /** + * query task relation by codes + * + * @param projectCode projectCode + * @param taskCode taskCode + * @param preTaskCodes preTaskCode list + * @return ProcessTaskRelation + */ + List<WorkflowTaskRelation> queryUpstreamByCodes(@Param("projectCode") long projectCode, + @Param("taskCode") long taskCode, + @Param("preTaskCodes") Long[] preTaskCodes); + /** * query workflow task relation by process definition code * @@ -108,6 +140,14 @@ public interface WorkflowTaskRelationMapper extends BaseMapper<WorkflowTaskRelat @Param("preTaskCode") long preTaskCode, @Param("postTaskCode") long postTaskCode); + /** + * delete workflow task relation + * + * @param workflowTaskRelationLog workflowTaskRelationLog + * @return int + */ + int deleteRelation(@Param("workflowTaskRelationLog") WorkflowTaskRelationLog workflowTaskRelationLog); + /** * query downstream workflow task relation by workflowDefinitionCode * @param workflowDefinitionCode @@ -115,6 +155,16 @@ public interface WorkflowTaskRelationMapper extends BaseMapper<WorkflowTaskRelat */ List<WorkflowTaskRelation> queryDownstreamByWorkflowDefinitionCode(@Param("workflowDefinitionCode") long workflowDefinitionCode); + /** + * Filter workflow task relation + * + * @param page page + * @param workflowTaskRelation process definition object + * @return workflow task relation IPage + */ + IPage<WorkflowTaskRelation> filterWorkflowTaskRelation(IPage<WorkflowTaskRelation> page, + @Param("relation") WorkflowTaskRelation workflowTaskRelation); + /** * batch update workflow task relation version * @@ -123,6 +173,9 @@ public interface WorkflowTaskRelationMapper extends BaseMapper<WorkflowTaskRelat */ int updateWorkflowTaskRelationTaskVersion(@Param("workflowTaskRelation") WorkflowTaskRelation workflowTaskRelation); + Long queryTaskCodeByTaskName(@Param("workflowCode") Long workflowCode, + @Param("taskName") String taskName); + void deleteByWorkflowDefinitionCodeAndVersion(@Param("workflowDefinitionCode") long workflowDefinitionCode, @Param("workflowDefinitionVersion") int workflowDefinitionVersion); diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/BaseDao.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/BaseDao.java index e2a53e75fe..664b56ee47 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/BaseDao.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/BaseDao.java @@ -27,6 +27,7 @@ import java.util.Optional; import lombok.NonNull; +import com.baomidou.mybatisplus.core.conditions.query.QueryWrapper; import com.baomidou.mybatisplus.core.mapper.BaseMapper; public abstract class BaseDao<ENTITY, MYBATIS_MAPPER extends BaseMapper<ENTITY>> implements IDao<ENTITY> { @@ -60,6 +61,14 @@ public abstract class BaseDao<ENTITY, MYBATIS_MAPPER extends BaseMapper<ENTITY>> return mybatisMapper.selectList(null); } + @Override + public List<ENTITY> queryByCondition(ENTITY queryCondition) { + if (queryCondition == null) { + throw new IllegalArgumentException("queryCondition can not be null"); + } + return mybatisMapper.selectList(new QueryWrapper<>(queryCondition)); + } + @Override public int insert(@NonNull ENTITY model) { return mybatisMapper.insert(model); @@ -85,4 +94,20 @@ public abstract class BaseDao<ENTITY, MYBATIS_MAPPER extends BaseMapper<ENTITY>> return mybatisMapper.deleteById(id) > 0; } + @Override + public boolean deleteByIds(Collection<? extends Serializable> ids) { + if (CollectionUtils.isEmpty(ids)) { + return true; + } + return mybatisMapper.deleteBatchIds(ids) > 0; + } + + @Override + public boolean deleteByCondition(ENTITY queryCondition) { + if (queryCondition == null) { + throw new IllegalArgumentException("queryCondition can not be null"); + } + return mybatisMapper.delete(new QueryWrapper<>(queryCondition)) > 0; + } + } diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/IDao.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/IDao.java index 98393210cd..ab77419600 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/IDao.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/IDao.java @@ -46,6 +46,11 @@ public interface IDao<Entity> { */ List<Entity> queryAll(); + /** + * Query the entity by condition. + */ + List<Entity> queryByCondition(Entity queryCondition); + /** * Insert the entity. */ @@ -66,4 +71,14 @@ public interface IDao<Entity> { */ boolean deleteById(@NonNull Serializable id); + /** + * Delete the entities by primary keys. + */ + boolean deleteByIds(Collection<? extends Serializable> ids); + + /** + * Delete the entities by condition. + */ + boolean deleteByCondition(Entity queryCondition); + } diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/TaskDefinitionDao.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/TaskDefinitionDao.java index 2aa08a40c8..8eb5856c43 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/TaskDefinitionDao.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/TaskDefinitionDao.java @@ -36,6 +36,16 @@ public interface TaskDefinitionDao extends IDao<TaskDefinition> { */ List<TaskDefinition> getTaskDefinitionListByDefinition(long workflowDefinitionCode); + /** + * Query task definition by code and version + * @param taskCode task code + * @param taskDefinitionVersion task definition version + * @return task definition + */ + TaskDefinition findTaskDefinition(long taskCode, int taskDefinitionVersion); + + void deleteByWorkflowDefinitionCodeAndVersion(long workflowDefinitionCode, int workflowDefinitionVersion); + void deleteByTaskDefinitionCodes(Set<Long> needToDeleteTaskDefinitionCodes); List<TaskDefinition> queryByCodes(Collection<Long> taskDefinitionCodes); diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/TaskInstanceContextDao.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/TaskInstanceContextDao.java index 44a8eb788f..9fcdec80f3 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/TaskInstanceContextDao.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/TaskInstanceContextDao.java @@ -24,6 +24,9 @@ import java.util.List; public interface TaskInstanceContextDao extends IDao<TaskInstanceContext> { + List<TaskInstanceContext> queryListByTaskInstanceIdAndContextType(Integer taskInstanceId, + ContextType contextType); + int deleteByTaskInstanceIdAndContextType(Integer taskInstanceId, ContextType contextType); diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/TaskInstanceDao.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/TaskInstanceDao.java index 8aa25949cb..39b2a073d7 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/TaskInstanceDao.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/TaskInstanceDao.java @@ -18,6 +18,8 @@ package org.apache.dolphinscheduler.dao.repository; import org.apache.dolphinscheduler.dao.entity.TaskInstance; +import org.apache.dolphinscheduler.dao.entity.WorkflowInstance; +import org.apache.dolphinscheduler.plugin.task.api.enums.TaskExecutionStatus; import java.util.List; import java.util.Set; @@ -34,6 +36,15 @@ public interface TaskInstanceDao extends IDao<TaskInstance> { */ boolean upsertTaskInstance(TaskInstance taskInstance); + /** + * Submit a task instance to DB. + * + * @param taskInstance task instance + * @param workflowInstance workflow instance + * @return task instance + */ + boolean submitTaskInstanceToDB(TaskInstance taskInstance, WorkflowInstance workflowInstance); + /** * Mark the task instance as invalid */ @@ -47,6 +58,23 @@ public interface TaskInstanceDao extends IDao<TaskInstance> { */ List<TaskInstance> queryValidTaskListByWorkflowInstanceId(Integer workflowInstanceId); + /** + * Query list of task instance by workflow instance id and task code + * + * @param workflowInstanceId workflowInstanceId + * @param taskCode task code + * @return list of valid task instance + */ + TaskInstance queryByWorkflowInstanceIdAndTaskCode(Integer workflowInstanceId, Long taskCode); + + /** + * find previous task list by work workflow id + * + * @param workflowInstanceId workflowInstanceId + * @return task instance list + */ + List<TaskInstance> queryPreviousTaskListByWorkflowInstanceId(Integer workflowInstanceId); + void deleteByWorkflowInstanceId(int workflowInstanceId); List<TaskInstance> queryByWorkflowInstanceId(Integer workflowInstanceId); @@ -71,4 +99,6 @@ public interface TaskInstanceDao extends IDao<TaskInstance> { TaskInstance queryLastTaskInstanceIntervalInWorkflowInstance(Integer workflowInstanceId, long depTaskCode); + void updateTaskInstanceState(Integer taskInstanceId, TaskExecutionStatus originState, + TaskExecutionStatus targetState); } diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/WorkerGroupDao.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/WorkerGroupDao.java index bf5d00f87e..7db0c7d10f 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/WorkerGroupDao.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/WorkerGroupDao.java @@ -23,6 +23,8 @@ import java.util.List; public interface WorkerGroupDao extends IDao<WorkerGroup> { + boolean deleteByWorkerGroupName(String workerGroupName); + List<String> queryAllWorkerGroupNames(); List<WorkerGroup> queryAllWorkerGroup(); diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/WorkflowInstanceDao.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/WorkflowInstanceDao.java index 861586d52d..52b0770109 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/WorkflowInstanceDao.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/WorkflowInstanceDao.java @@ -80,6 +80,8 @@ public interface WorkflowInstanceDao extends IDao<WorkflowInstance> { */ WorkflowInstance queryFirstStartWorkflowInstance(Long definitionCode); + WorkflowInstance querySubWorkflowInstanceByParentId(Integer workflowInstanceId, Integer taskInstanceId); + List<WorkflowInstance> queryByWorkflowCodeVersionStatus(Long workflowDefinitionCode, int workflowDefinitionVersion, int[] states); diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/WorkflowTaskLineageDao.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/WorkflowTaskLineageDao.java index 09fa59b49a..ac7f3ea56d 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/WorkflowTaskLineageDao.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/WorkflowTaskLineageDao.java @@ -41,4 +41,5 @@ public interface WorkflowTaskLineageDao extends IDao<WorkflowTaskLineage> { List<WorkflowTaskLineage> queryByWorkflowDefinitionCode(long workflowDefinitionCode); + int updateWorkflowTaskLineage(List<WorkflowTaskLineage> workflowTaskLineages); } diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/TaskDefinitionDaoImpl.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/TaskDefinitionDaoImpl.java index 228e1dcc7d..ca3c9b2e6f 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/TaskDefinitionDaoImpl.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/TaskDefinitionDaoImpl.java @@ -86,6 +86,16 @@ public class TaskDefinitionDaoImpl extends BaseDao<TaskDefinition, TaskDefinitio return Lists.newArrayList(taskDefinitionLogs); } + @Override + public TaskDefinition findTaskDefinition(long taskCode, int taskDefinitionVersion) { + return taskDefinitionLogMapper.queryByDefinitionCodeAndVersion(taskCode, taskDefinitionVersion); + } + + @Override + public void deleteByWorkflowDefinitionCodeAndVersion(long workflowDefinitionCode, int workflowDefinitionVersion) { + mybatisMapper.deleteByWorkflowDefinitionCodeAndVersion(workflowDefinitionCode, workflowDefinitionVersion); + } + @Override public void deleteByTaskDefinitionCodes(Set<Long> needToDeleteTaskDefinitionCodes) { if (CollectionUtils.isEmpty(needToDeleteTaskDefinitionCodes)) { diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/TaskInstanceContextDaoImpl.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/TaskInstanceContextDaoImpl.java index 1e796f459a..5a147d1c17 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/TaskInstanceContextDaoImpl.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/TaskInstanceContextDaoImpl.java @@ -46,6 +46,15 @@ public class TaskInstanceContextDaoImpl extends BaseDao<TaskInstanceContext, Tas super(taskInstanceContextMapper); } + @Override + public List<TaskInstanceContext> queryListByTaskInstanceIdAndContextType(Integer taskInstanceId, + ContextType contextType) { + if (taskInstanceId == null) { + return Collections.emptyList(); + } + return mybatisMapper.queryListByTaskInstanceIdAndContextType(taskInstanceId, contextType); + } + @Override public int deleteByTaskInstanceIdAndContextType(Integer taskInstanceId, ContextType contextType) { if (taskInstanceId == null) { diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/TaskInstanceDaoImpl.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/TaskInstanceDaoImpl.java index 3448bb6eaf..c675f502aa 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/TaskInstanceDaoImpl.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/TaskInstanceDaoImpl.java @@ -17,15 +17,20 @@ package org.apache.dolphinscheduler.dao.repository.impl; +import org.apache.dolphinscheduler.common.enums.FailureStrategy; import org.apache.dolphinscheduler.common.enums.Flag; +import org.apache.dolphinscheduler.common.enums.WorkflowExecutionStatus; import org.apache.dolphinscheduler.dao.entity.TaskInstance; +import org.apache.dolphinscheduler.dao.entity.WorkflowInstance; import org.apache.dolphinscheduler.dao.mapper.TaskInstanceMapper; import org.apache.dolphinscheduler.dao.mapper.WorkflowInstanceMapper; import org.apache.dolphinscheduler.dao.repository.BaseDao; import org.apache.dolphinscheduler.dao.repository.TaskInstanceDao; +import org.apache.dolphinscheduler.plugin.task.api.enums.TaskExecutionStatus; import org.apache.commons.collections4.CollectionUtils; +import java.util.Date; import java.util.List; import java.util.Set; @@ -58,6 +63,31 @@ public class TaskInstanceDaoImpl extends BaseDao<TaskInstance, TaskInstanceMappe } } + @Override + public boolean submitTaskInstanceToDB(TaskInstance taskInstance, WorkflowInstance workflowInstance) { + WorkflowExecutionStatus processInstanceState = workflowInstance.getState(); + if (processInstanceState.isFinalState() || processInstanceState == WorkflowExecutionStatus.READY_STOP) { + log.warn("processInstance: {} state was: {}, skip submit this task, taskCode: {}", + workflowInstance.getId(), + processInstanceState, + taskInstance.getTaskCode()); + return false; + } + if (processInstanceState == WorkflowExecutionStatus.READY_PAUSE) { + taskInstance.setState(TaskExecutionStatus.PAUSE); + } + taskInstance.setExecutorId(workflowInstance.getExecutorId()); + taskInstance.setExecutorName(workflowInstance.getExecutorName()); + taskInstance.setState(getSubmitTaskState(taskInstance, workflowInstance)); + if (taskInstance.getSubmitTime() == null) { + taskInstance.setSubmitTime(new Date()); + } + if (taskInstance.getFirstSubmitTime() == null) { + taskInstance.setFirstSubmitTime(taskInstance.getSubmitTime()); + } + return upsertTaskInstance(taskInstance); + } + @Override public void markTaskInstanceInvalid(List<TaskInstance> taskInstances) { if (CollectionUtils.isEmpty(taskInstances)) { @@ -69,11 +99,59 @@ public class TaskInstanceDaoImpl extends BaseDao<TaskInstance, TaskInstanceMappe } } + private TaskExecutionStatus getSubmitTaskState(TaskInstance taskInstance, WorkflowInstance workflowInstance) { + TaskExecutionStatus state = taskInstance.getState(); + if (state == TaskExecutionStatus.RUNNING_EXECUTION + || state == TaskExecutionStatus.DELAY_EXECUTION + || state == TaskExecutionStatus.KILL + || state == TaskExecutionStatus.DISPATCH) { + return state; + } + + if (workflowInstance.getState() == WorkflowExecutionStatus.READY_PAUSE) { + state = TaskExecutionStatus.PAUSE; + } else if (workflowInstance.getState() == WorkflowExecutionStatus.READY_STOP + || !checkProcessStrategy(taskInstance, workflowInstance)) { + state = TaskExecutionStatus.KILL; + } else { + state = TaskExecutionStatus.SUBMITTED_SUCCESS; + } + return state; + } + + private boolean checkProcessStrategy(TaskInstance taskInstance, WorkflowInstance workflowInstance) { + FailureStrategy failureStrategy = workflowInstance.getFailureStrategy(); + if (failureStrategy == FailureStrategy.CONTINUE) { + return true; + } + List<TaskInstance> taskInstances = + this.queryValidTaskListByWorkflowInstanceId(taskInstance.getWorkflowInstanceId()); + + for (TaskInstance task : taskInstances) { + if (task.getState() == TaskExecutionStatus.FAILURE + && task.getRetryTimes() >= task.getMaxRetryTimes()) { + return false; + } + } + return true; + } + @Override public List<TaskInstance> queryValidTaskListByWorkflowInstanceId(Integer processInstanceId) { return mybatisMapper.findValidTaskListByWorkflowInstanceId(processInstanceId, Flag.YES); } + @Override + public TaskInstance queryByWorkflowInstanceIdAndTaskCode(Integer workflowInstanceId, Long taskCode) { + return mybatisMapper.queryByInstanceIdAndCode(workflowInstanceId, taskCode); + } + + @Override + public List<TaskInstance> queryPreviousTaskListByWorkflowInstanceId(Integer workflowInstanceId) { + WorkflowInstance workflowInstance = workflowInstanceMapper.selectById(workflowInstanceId); + return mybatisMapper.findValidTaskListByWorkflowInstanceId(workflowInstanceId, Flag.NO); + } + @Override public void deleteByWorkflowInstanceId(int workflowInstanceId) { mybatisMapper.deleteByWorkflowInstanceId(workflowInstanceId); @@ -95,4 +173,10 @@ public class TaskInstanceDaoImpl extends BaseDao<TaskInstance, TaskInstanceMappe return mybatisMapper.findLastTaskInstance(workflowInstanceId, depTaskCode); } + @Override + public void updateTaskInstanceState(Integer taskInstanceId, + TaskExecutionStatus originState, + TaskExecutionStatus targetState) { + mybatisMapper.updateTaskInstanceState(taskInstanceId, originState.getCode(), targetState.getCode()); + } } diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/WorkerGroupDaoImpl.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/WorkerGroupDaoImpl.java index f51a253d53..b5203bb661 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/WorkerGroupDaoImpl.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/WorkerGroupDaoImpl.java @@ -36,6 +36,12 @@ public class WorkerGroupDaoImpl extends BaseDao<WorkerGroup, WorkerGroupMapper> super(workerGroupMapper); } + @Override + public boolean deleteByWorkerGroupName(String workerGroupName) { + int deleted = mybatisMapper.deleteByWorkerGroupName(workerGroupName); + return deleted > 0; + } + @Override public List<String> queryAllWorkerGroupNames() { return mybatisMapper.queryAllWorkerGroup().stream() diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/WorkflowInstanceDaoImpl.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/WorkflowInstanceDaoImpl.java index c1b93972b6..6f35057d44 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/WorkflowInstanceDaoImpl.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/WorkflowInstanceDaoImpl.java @@ -19,6 +19,7 @@ package org.apache.dolphinscheduler.dao.repository.impl; import org.apache.dolphinscheduler.common.enums.WorkflowExecutionStatus; import org.apache.dolphinscheduler.dao.entity.WorkflowInstance; +import org.apache.dolphinscheduler.dao.entity.WorkflowInstanceRelation; import org.apache.dolphinscheduler.dao.mapper.WorkflowInstanceMapper; import org.apache.dolphinscheduler.dao.mapper.WorkflowInstanceRelationMapper; import org.apache.dolphinscheduler.dao.repository.BaseDao; @@ -143,6 +144,18 @@ public class WorkflowInstanceDaoImpl extends BaseDao<WorkflowInstance, WorkflowI return mybatisMapper.queryFirstStartWorkflowInstance(definitionCode); } + @Override + public WorkflowInstance querySubWorkflowInstanceByParentId(Integer workflowInstanceId, Integer taskInstanceId) { + WorkflowInstance workflowInstance = null; + WorkflowInstanceRelation workflowInstanceRelation = + workflowInstanceRelationMapper.queryByParentId(workflowInstanceId, taskInstanceId); + if (workflowInstanceRelation == null || workflowInstanceRelation.getWorkflowInstanceId() == 0) { + return workflowInstance; + } + workflowInstance = queryById(workflowInstanceRelation.getWorkflowInstanceId()); + return workflowInstance; + } + @Override public List<WorkflowInstance> queryByWorkflowCodeVersionStatus(Long workflowDefinitionCode, int workflowDefinitionVersion, diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/WorkflowTaskLineageDaoImpl.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/WorkflowTaskLineageDaoImpl.java index 6d5ed08527..0539b34b95 100644 --- a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/WorkflowTaskLineageDaoImpl.java +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/WorkflowTaskLineageDaoImpl.java @@ -26,6 +26,7 @@ import org.apache.dolphinscheduler.dao.repository.WorkflowTaskLineageDao; import org.apache.commons.collections4.CollectionUtils; import java.util.List; +import java.util.stream.Collectors; import lombok.NonNull; @@ -83,4 +84,14 @@ public class WorkflowTaskLineageDaoImpl extends BaseDao<WorkflowTaskLineage, Wor return mybatisMapper.queryByWorkflowDefinitionCode(workflowDefinitionCode); } + @Override + public int updateWorkflowTaskLineage(List<WorkflowTaskLineage> workflowTaskLineages) { + if (CollectionUtils.isEmpty(workflowTaskLineages)) { + return 0; + } + this.batchDeleteByWorkflowDefinitionCode( + workflowTaskLineages.stream().map(WorkflowTaskLineage::getWorkflowDefinitionCode) + .distinct().collect(Collectors.toList())); + return mybatisMapper.batchInsert(workflowTaskLineages); + } } diff --git a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/AlertMapperTest.java b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/AlertMapperTest.java index 48ffb2e295..346004453f 100644 --- a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/AlertMapperTest.java +++ b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/AlertMapperTest.java @@ -25,6 +25,9 @@ import org.apache.dolphinscheduler.dao.entity.Alert; import org.apache.commons.codec.digest.DigestUtils; +import java.util.HashMap; +import java.util.Map; + import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; import org.springframework.beans.factory.annotation.Autowired; @@ -90,6 +93,23 @@ public class AlertMapperTest extends BaseDaoTest { Assertions.assertNull(actualAlert); } + /** + * create alert map + * + * @param count alert count + * @param alertStatus alert status + * @return alert map + */ + private Map<Integer, Alert> createAlertMap(Integer count, AlertStatus alertStatus) { + Map<Integer, Alert> alertMap = new HashMap<>(); + + for (int i = 0; i < count; i++) { + Alert alert = createAlert(alertStatus); + alertMap.put(alert.getId(), alert); + } + return alertMap; + } + /** * create alert * diff --git a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/AuditLogMapperTest.java b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/AuditLogMapperTest.java index eae7a1da0b..92418193b8 100644 --- a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/AuditLogMapperTest.java +++ b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/AuditLogMapperTest.java @@ -21,6 +21,7 @@ import org.apache.dolphinscheduler.common.enums.AuditModelType; import org.apache.dolphinscheduler.common.enums.AuditOperationType; import org.apache.dolphinscheduler.dao.BaseDaoTest; import org.apache.dolphinscheduler.dao.entity.AuditLog; +import org.apache.dolphinscheduler.dao.entity.Project; import java.util.ArrayList; import java.util.Date; @@ -39,6 +40,9 @@ public class AuditLogMapperTest extends BaseDaoTest { @Autowired private AuditLogMapper logMapper; + @Autowired + private ProjectMapper projectMapper; + private void insertOne(AuditModelType objectType) { AuditLog auditLog = new AuditLog(); auditLog.setUserId(1); @@ -53,6 +57,17 @@ public class AuditLogMapperTest extends BaseDaoTest { logMapper.insert(auditLog); } + private Project insertProject() { + Project project = new Project(); + project.setName("ut project"); + project.setUserId(111); + project.setCode(1L); + project.setCreateTime(new Date()); + project.setUpdateTime(new Date()); + projectMapper.insert(project); + return project; + } + /** * test page query */ diff --git a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/CommandMapperTest.java b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/CommandMapperTest.java index da0226f3c2..59c69c981d 100644 --- a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/CommandMapperTest.java +++ b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/CommandMapperTest.java @@ -34,7 +34,9 @@ import org.apache.dolphinscheduler.dao.entity.WorkflowDefinition; import org.apache.dolphinscheduler.dao.utils.WorkerGroupUtils; import java.util.Date; +import java.util.HashMap; import java.util.List; +import java.util.Map; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; @@ -128,6 +130,8 @@ public class CommandMapperTest extends BaseDaoTest { public void testGetAll() { Integer count = 10; + Map<Integer, Command> commandMap = createCommandMap(count); + List<Command> actualCommands = commandMapper.selectList(null); Assertions.assertTrue(actualCommands.size() >= count); @@ -251,6 +255,22 @@ public class CommandMapperTest extends BaseDaoTest { return workflowDefinition; } + /** + * create command map + * + * @param count map count + * @return command map + */ + private Map<Integer, Command> createCommandMap(Integer count) { + Map<Integer, Command> commandMap = new HashMap<>(); + + for (int i = 0; i < count; i++) { + Command command = createCommand(); + commandMap.put(command.getId(), command); + } + return commandMap; + } + /** * create command * diff --git a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/ProjectWorkerGroupMapperTest.java b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/ProjectWorkerGroupMapperTest.java index c438363399..bddcd10e89 100644 --- a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/ProjectWorkerGroupMapperTest.java +++ b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/ProjectWorkerGroupMapperTest.java @@ -77,6 +77,7 @@ public class ProjectWorkerGroupMapperTest extends BaseDaoTest { */ @Test public void testQuery() { + ProjectWorkerGroup projectWorkerGroup = insertOne(); // query List<ProjectWorkerGroup> projectUsers = projectWorkerGroupMapper.selectList(null); Assertions.assertNotEquals(0, projectUsers.size()); diff --git a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/TaskDefinitionLogMapperTest.java b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/TaskDefinitionLogMapperTest.java index 133d3487f5..4d815f4b03 100644 --- a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/TaskDefinitionLogMapperTest.java +++ b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/TaskDefinitionLogMapperTest.java @@ -89,6 +89,7 @@ public class TaskDefinitionLogMapperTest extends BaseDaoTest { ArrayList<TaskDefinition> taskDefinitions = new ArrayList<>(); taskDefinitions.add(taskDefinition); + TaskDefinitionLog taskDefinitionLog = insertOne(); List<TaskDefinitionLog> taskDefinitionLogs = taskDefinitionLogMapper.queryByTaskDefinitions(taskDefinitions); Assertions.assertNotEquals(0, taskDefinitionLogs.size()); } diff --git a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/TaskGroupQueueMapperTest.java b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/TaskGroupQueueMapperTest.java index db97433cce..b46bed5953 100644 --- a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/TaskGroupQueueMapperTest.java +++ b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/TaskGroupQueueMapperTest.java @@ -33,6 +33,8 @@ public class TaskGroupQueueMapperTest extends BaseDaoTest { @Autowired TaskGroupQueueMapper taskGroupQueueMapper; + int userId = 1; + public TaskGroupQueue insertOne() { TaskGroupQueue taskGroupQueue = new TaskGroupQueue(); taskGroupQueue.setTaskName("task1"); diff --git a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/TenantMapperTest.java b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/TenantMapperTest.java index f8b3ae6f3f..22465b62a5 100644 --- a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/TenantMapperTest.java +++ b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/TenantMapperTest.java @@ -143,7 +143,6 @@ public class TenantMapperTest extends BaseDaoTest { Assertions.assertNotEquals(0, tenantIPage.getTotal()); } - @Test public void testExistTenant() { String tenantCode = "test_code"; Assertions.assertNull(tenantMapper.existTenant(tenantCode)); diff --git a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/UserMapperTest.java b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/UserMapperTest.java index b2a0fc5dff..5716b7072b 100644 --- a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/UserMapperTest.java +++ b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/UserMapperTest.java @@ -21,6 +21,7 @@ import org.apache.dolphinscheduler.common.enums.UserType; import org.apache.dolphinscheduler.common.utils.DateUtils; import org.apache.dolphinscheduler.dao.BaseDaoTest; import org.apache.dolphinscheduler.dao.entity.AccessToken; +import org.apache.dolphinscheduler.dao.entity.AlertGroup; import org.apache.dolphinscheduler.dao.entity.Queue; import org.apache.dolphinscheduler.dao.entity.Tenant; import org.apache.dolphinscheduler.dao.entity.User; @@ -117,6 +118,23 @@ public class UserMapperTest extends BaseDaoTest { return user; } + /** + * insert one AlertGroup + * + * @return AlertGroup + */ + private AlertGroup insertOneAlertGroup() { + // insertOne + AlertGroup alertGroup = new AlertGroup(); + alertGroup.setGroupName("alert group 1"); + alertGroup.setDescription("alert test1"); + + alertGroup.setCreateTime(new Date()); + alertGroup.setUpdateTime(new Date()); + alertGroupMapper.insert(alertGroup); + return alertGroup; + } + /** * insert one AccessToken * diff --git a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/WorkflowInstanceMapMapperTest.java b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/WorkflowInstanceMapMapperTest.java index 76a42d0565..527fce4aec 100644 --- a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/WorkflowInstanceMapMapperTest.java +++ b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/WorkflowInstanceMapMapperTest.java @@ -74,6 +74,7 @@ public class WorkflowInstanceMapMapperTest extends BaseDaoTest { */ @Test public void testQuery() { + WorkflowInstanceRelation workflowInstanceRelation = insertOne(); // query List<WorkflowInstanceRelation> dataSources = workflowInstanceRelationMapper.selectList(null); Assertions.assertNotEquals(0, dataSources.size()); diff --git a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/WorkflowInstanceMapperTest.java b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/WorkflowInstanceMapperTest.java index b847db8134..952c3480e5 100644 --- a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/WorkflowInstanceMapperTest.java +++ b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/WorkflowInstanceMapperTest.java @@ -45,6 +45,9 @@ public class WorkflowInstanceMapperTest extends BaseDaoTest { @Autowired private WorkflowDefinitionMapper workflowDefinitionMapper; + @Autowired + private ProjectMapper projectMapper; + /** * insert process instance with specified start time and end time,set state to SUCCESS */ diff --git a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/WorkflowTaskRelationLogMapperTest.java b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/WorkflowTaskRelationLogMapperTest.java index 1279e49bea..9c00aead01 100644 --- a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/WorkflowTaskRelationLogMapperTest.java +++ b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/WorkflowTaskRelationLogMapperTest.java @@ -54,7 +54,7 @@ public class WorkflowTaskRelationLogMapperTest extends BaseDaoTest { @Test public void testQueryByWorkflowCodeAndVersion() { - insertOne(); + WorkflowTaskRelationLog processTaskRelationLog = insertOne(); List<WorkflowTaskRelationLog> processTaskRelationLogs = workflowTaskRelationLogMapper .queryByWorkflowCodeAndVersion(1L, 1); Assertions.assertNotEquals(0, processTaskRelationLogs.size()); diff --git a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/WorkflowTaskRelationMapperTest.java b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/WorkflowTaskRelationMapperTest.java index cb180d7803..e227980963 100644 --- a/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/WorkflowTaskRelationMapperTest.java +++ b/dolphinscheduler-dao/src/test/java/org/apache/dolphinscheduler/dao/mapper/WorkflowTaskRelationMapperTest.java @@ -55,21 +55,21 @@ public class WorkflowTaskRelationMapperTest extends BaseDaoTest { @Test public void testQueryByWorkflowDefinitionCode() { - insertOne(); + WorkflowTaskRelation workflowTaskRelation = insertOne(); List<WorkflowTaskRelation> workflowTaskRelations = workflowTaskRelationMapper.queryByWorkflowDefinitionCode(1L); Assertions.assertNotEquals(0, workflowTaskRelations.size()); } @Test public void testQueryByTaskCode() { - insertOne(); + WorkflowTaskRelation workflowTaskRelation = insertOne(); List<WorkflowTaskRelation> workflowTaskRelations = workflowTaskRelationMapper.queryByTaskCode(2L); Assertions.assertNotEquals(0, workflowTaskRelations.size()); } @Test public void testQueryByTaskCodes() { - insertOne(); + WorkflowTaskRelation workflowTaskRelation = insertOne(); Long[] codes = Arrays.array(1L, 2L); List<WorkflowTaskRelation> workflowTaskRelations = workflowTaskRelationMapper.queryByTaskCodes(codes); @@ -78,7 +78,7 @@ public class WorkflowTaskRelationMapperTest extends BaseDaoTest { @Test public void testDeleteByWorkflowDefinitionCode() { - insertOne(); + WorkflowTaskRelation workflowTaskRelation = insertOne(); int i = workflowTaskRelationMapper.deleteByWorkflowDefinitionCode(1L, 1L); Assertions.assertNotEquals(0, i); }
