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);
}