This is an automated email from the ASF dual-hosted git repository.

SbloodyS pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/dolphinscheduler.git


The following commit(s) were added to refs/heads/dev by this push:
     new 0f975bac39 Revert "[Improvement-18056] Clean up unused methods and 
classes in the dolphi…" (#18206)
0f975bac39 is described below

commit 0f975bac396da3604aa5b28757420ac8112f2777
Author: Wenjun Ruan <[email protected]>
AuthorDate: Fri May 1 21:28:38 2026 +0800

    Revert "[Improvement-18056] Clean up unused methods and classes in the 
dolphi…" (#18206)
---
 .../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);
     }

Reply via email to