This is an automated email from the ASF dual-hosted git repository.
zhongjiajie 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 76bbcbeb30 [Improvement][Master] Split the task dependency mode of the
dependent node into workflow dependent and task dependent (#14824)
76bbcbeb30 is described below
commit 76bbcbeb30e800e4767d5c4192b5a0f5891d6321
Author: sgw <[email protected]>
AuthorDate: Sat Sep 9 17:23:34 2023 +0800
[Improvement][Master] Split the task dependency mode of the dependent node
into workflow dependent and task dependent (#14824)
* [Improvement][Master] Split the task dependency mode of the dependent
node into workflow dependent and task dependent. (#11970)
* [Improvement][Master] Split the task dependency mode of the dependent
node into workflow dependent and task dependent. (#11970)
- delete useless output
* [Improvement][Master] Split the task dependency mode of the dependent
node into workflow dependent and task dependent. (#11970)
- add log
- run spotless check
---------
Co-authored-by: 旺阳 <[email protected]>
Co-authored-by: JinYong Li <[email protected]>
---
docs/docs/en/guide/task/dependent.md | 7 +-
docs/docs/zh/guide/task/dependent.md | 8 +-
docs/img/tasks/demo/dependent_task01.png | Bin 159694 -> 144045 bytes
docs/img/tasks/demo/dependent_task02.png | Bin 159538 -> 139834 bytes
docs/img/tasks/demo/dependent_task03.png | Bin 164972 -> 151023 bytes
.../common/constants/Constants.java | 3 +-
.../dao/mapper/TaskInstanceMapper.java | 19 +++
.../dao/repository/TaskDefinitionDao.java | 7 ++
.../dao/repository/TaskInstanceDao.java | 22 ++++
.../dao/repository/impl/TaskDefinitionDaoImpl.java | 5 +
.../dao/repository/impl/TaskInstanceDaoImpl.java | 15 +++
.../dao/mapper/TaskInstanceMapper.xml | 35 +++++-
.../DependentAsyncTaskExecuteFunction.java | 6 +
.../server/master/utils/DependentExecute.java | 138 ++++++++++++++++++++-
dolphinscheduler-ui/src/locales/en_US/project.ts | 4 +
dolphinscheduler-ui/src/locales/zh_CN/project.ts | 3 +
.../task/components/node/fields/use-dependent.ts | 42 ++++++-
.../views/projects/task/components/node/types.ts | 2 +
18 files changed, 304 insertions(+), 12 deletions(-)
diff --git a/docs/docs/en/guide/task/dependent.md
b/docs/docs/en/guide/task/dependent.md
index a0bdc6fcdf..c034e6c044 100644
--- a/docs/docs/en/guide/task/dependent.md
+++ b/docs/docs/en/guide/task/dependent.md
@@ -27,11 +27,14 @@ Dependent nodes are **dependency check nodes**. For
example, process A depends o
The Dependent node provides a logical judgment function, which can detect the
execution of the dependent node according to the logic.
-For example, process A is a weekly task, processes B and C are daily tasks,
and task A requires tasks B and C to be successfully executed every day of the
last week.
+Two dependency modes are supported, including workflow-dependent and
task-dependent. The task-dependent mode is divided into two cases: depend on
all tasks in the workflow and depend on a single task.
+The workflow-dependent mode checks the status of the dependent workflow; the
all-task-dependent mode checks the status of all tasks in the workflow; and the
single-task-dependent mode checks the status of the dependent task.
+
+For example, process A is a weekly task, processes B and C are daily tasks,
and task A requires tasks B and C to be successfully executed last week.

-And another example is that process A is a weekly report task, processes B and
C are daily tasks, and task A requires tasks B or C to be successfully executed
every day of the last week:
+And another example is that process A is a weekly report task, processes B and
C are daily tasks, and task A requires tasks B or C to be successfully executed
last week:

diff --git a/docs/docs/zh/guide/task/dependent.md
b/docs/docs/zh/guide/task/dependent.md
index c3318ae3d2..0f260f4944 100644
--- a/docs/docs/zh/guide/task/dependent.md
+++ b/docs/docs/zh/guide/task/dependent.md
@@ -27,11 +27,15 @@ Dependent 节点,就是**依赖检查节点**。比如 A 流程依赖昨天的
Dependent 节点提供了逻辑判断功能,可以按照逻辑来检测所依赖节点的执行情况。
-例如,A 流程为周报任务,B、C 流程为天任务,A 任务需要 B、C 任务在上周的每一天都执行成功,如图示:
+支持两种依赖模式,包括依赖于工作流和依赖于任务。依赖于任务的模式分依赖工作流中的所有任务和依赖单个任务两种情况。
+依赖工作流的模式会检查所依赖的工作流的状态;依赖所有任务的模式会检查工作流中所有任务的状态;
+依赖单个任务的模式会检查所依赖的任务的状态。
+
+例如,A 流程为周报任务,B、C 流程为天任务,A 任务需要 B、C 任务在上周执行成功,如图示:

-例如,A 流程为周报任务,B、C 流程为天任务,A 任务需要 B 或 C 任务在上周的每一天都执行成功,如图示:
+例如,A 流程为周报任务,B、C 流程为天任务,A 任务需要 B 或 C 任务在上周执行成功,如图示:

diff --git a/docs/img/tasks/demo/dependent_task01.png
b/docs/img/tasks/demo/dependent_task01.png
index 68a140a15d..9bdb05ba24 100644
Binary files a/docs/img/tasks/demo/dependent_task01.png and
b/docs/img/tasks/demo/dependent_task01.png differ
diff --git a/docs/img/tasks/demo/dependent_task02.png
b/docs/img/tasks/demo/dependent_task02.png
index 3f2afa4158..3a424314b3 100644
Binary files a/docs/img/tasks/demo/dependent_task02.png and
b/docs/img/tasks/demo/dependent_task02.png differ
diff --git a/docs/img/tasks/demo/dependent_task03.png
b/docs/img/tasks/demo/dependent_task03.png
index 1bc2f50aac..9dd130fbfd 100644
Binary files a/docs/img/tasks/demo/dependent_task03.png and
b/docs/img/tasks/demo/dependent_task03.png differ
diff --git
a/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/constants/Constants.java
b/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/constants/Constants.java
index e98f5f202f..5769126659 100644
---
a/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/constants/Constants.java
+++
b/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/constants/Constants.java
@@ -488,7 +488,8 @@ public final class Constants {
public static final String ALIAS = "alias";
public static final String CONTENT = "content";
public static final String DEPENDENT_SPLIT = ":||";
- public static final long DEPENDENT_ALL_TASK_CODE = 0;
+ public static final long DEPENDENT_ALL_TASK_CODE = -1;
+ public static final long DEPENDENT_WORKFLOW_CODE = 0;
/**
* preview schedule execute count
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 ae66577214..27e4eb4650 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
@@ -161,4 +161,23 @@ public interface TaskInstanceMapper extends
BaseMapper<TaskInstance> {
void deleteByWorkflowInstanceId(@Param("workflowInstanceId") int
workflowInstanceId);
List<TaskInstance> findByWorkflowInstanceId(@Param("workflowInstanceId")
Integer workflowInstanceId);
+
+ /**
+ * find last task instance list in the date interval
+ *
+ * @param taskCodes taskCodes
+ * @param startTime startTime
+ * @param endTime endTime
+ * @param testFlag testFlag
+ * @return task instance list
+ */
+ List<TaskInstance> findLastTaskInstances(@Param("taskCodes") Set<Long>
taskCodes,
+ @Param("startTime") Date
startTime,
+ @Param("endTime") Date endTime,
+ @Param("testFlag") int testFlag);
+
+ TaskInstance findLastTaskInstance(@Param("taskCode") long depTaskCode,
+ @Param("startTime") Date startTime,
+ @Param("endTime") Date endTime,
+ @Param("testFlag") int testFlag);
}
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 931a009041..652ecbfbd0 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
@@ -49,4 +49,11 @@ public interface TaskDefinitionDao extends
IDao<TaskDefinition> {
void deleteByTaskDefinitionCodes(Set<Long>
needToDeleteTaskDefinitionCodes);
List<TaskDefinition> queryByCodes(Collection<Long> taskDefinitionCodes);
+
+ /**
+ * Query task definition by code
+ * @param taskCode task code
+ * @return task definition
+ */
+ TaskDefinition queryByCode(long taskCode);
}
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 4a6d1435ff..b5d41f8783 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
@@ -19,8 +19,10 @@ package org.apache.dolphinscheduler.dao.repository;
import org.apache.dolphinscheduler.dao.entity.ProcessInstance;
import org.apache.dolphinscheduler.dao.entity.TaskInstance;
+import org.apache.dolphinscheduler.plugin.task.api.model.DateInterval;
import java.util.List;
+import java.util.Set;
/**
* Task Instance DAO
@@ -85,4 +87,24 @@ public interface TaskInstanceDao extends IDao<TaskInstance> {
void deleteByWorkflowInstanceId(int workflowInstanceId);
List<TaskInstance> queryByWorkflowInstanceId(Integer processInstanceId);
+
+ /**
+ * find last task instance list corresponding to taskCodes in the date
interval
+ *
+ * @param taskCodes taskCodes
+ * @param dateInterval dateInterval
+ * @param testFlag test flag
+ * @return task instance list
+ */
+ List<TaskInstance> queryLastTaskInstanceListIntervalByTaskCodes(Set<Long>
taskCodes, DateInterval dateInterval,
+ int
testFlag);
+
+ /**
+ * find last task instance corresponding to taskCode in the date interval
+ * @param depTaskCode taskCode
+ * @param dateInterval dateInterval
+ * @param testFlag test flag
+ * @return task instance
+ */
+ TaskInstance queryLastTaskInstanceIntervalByTaskCode(long depTaskCode,
DateInterval dateInterval, int testFlag);
}
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 1672a456d1..f92516d353 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
@@ -114,4 +114,9 @@ public class TaskDefinitionDaoImpl extends
BaseDao<TaskDefinition, TaskDefinitio
return mybatisMapper.queryByCodeList(taskDefinitionCodes);
}
+ @Override
+ public TaskDefinition queryByCode(long taskCode) {
+ return mybatisMapper.queryByCode(taskCode);
+ }
+
}
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 479247dbb9..9cbb92286c 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
@@ -27,11 +27,13 @@ import
org.apache.dolphinscheduler.dao.mapper.TaskInstanceMapper;
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.dolphinscheduler.plugin.task.api.model.DateInterval;
import org.apache.commons.lang3.StringUtils;
import java.util.Date;
import java.util.List;
+import java.util.Set;
import lombok.NonNull;
import lombok.extern.slf4j.Slf4j;
@@ -171,4 +173,17 @@ public class TaskInstanceDaoImpl extends
BaseDao<TaskInstance, TaskInstanceMappe
return mybatisMapper.findByWorkflowInstanceId(workflowInstanceId);
}
+ @Override
+ public List<TaskInstance>
queryLastTaskInstanceListIntervalByTaskCodes(Set<Long> taskCodes,
+
DateInterval dateInterval, int testFlag) {
+ return mybatisMapper.findLastTaskInstances(taskCodes,
dateInterval.getStartTime(), dateInterval.getEndTime(),
+ testFlag);
+ }
+
+ @Override
+ public TaskInstance queryLastTaskInstanceIntervalByTaskCode(long
depTaskCode, DateInterval dateInterval,
+ int testFlag) {
+ return mybatisMapper.findLastTaskInstance(depTaskCode,
dateInterval.getStartTime(), dateInterval.getEndTime(),
+ testFlag);
+ }
}
diff --git
a/dolphinscheduler-dao/src/main/resources/org/apache/dolphinscheduler/dao/mapper/TaskInstanceMapper.xml
b/dolphinscheduler-dao/src/main/resources/org/apache/dolphinscheduler/dao/mapper/TaskInstanceMapper.xml
index 6dc684ad9e..cb45014dc7 100644
---
a/dolphinscheduler-dao/src/main/resources/org/apache/dolphinscheduler/dao/mapper/TaskInstanceMapper.xml
+++
b/dolphinscheduler-dao/src/main/resources/org/apache/dolphinscheduler/dao/mapper/TaskInstanceMapper.xml
@@ -334,7 +334,40 @@
where instance.process_instance_id = #{processInstanceId}
and que.status = #{status}
</select>
-
+ <select id="findLastTaskInstances"
resultType="org.apache.dolphinscheduler.dao.entity.TaskInstance">
+ select
+ <include refid="baseSqlV2">
+ <property name="alias" value="instance"/>
+ </include>
+ from t_ds_task_instance instance
+ join (
+ select task_code, max(end_time) as max_end_time
+ from t_ds_task_instance
+ where 1=1 and test_flag = #{testFlag}
+ <if test="taskCodes != null and taskCodes.size() != 0">
+ and task_code in
+ <foreach collection="taskCodes" index="index" item="i" open="("
separator="," close=")">
+ #{i}
+ </foreach>
+ </if>
+ <if test="startTime!=null and endTime != null">
+ and start_time <![CDATA[ >= ]]> #{startTime} and start_time
<![CDATA[ <= ]]> #{endTime}
+ </if>
+ group by task_code
+ ) t_max
+ on instance.task_code = t_max.task_code and instance.end_time =
t_max.max_end_time
+ </select>
+ <select id="findLastTaskInstance"
resultType="org.apache.dolphinscheduler.dao.entity.TaskInstance">
+ select
+ <include refid="baseSql"/>
+ from t_ds_task_instance
+ where task_code = #{taskCode}
+ <if test="startTime!=null and endTime != null">
+ and start_time <![CDATA[ >= ]]> #{startTime} and start_time
<![CDATA[ <= ]]> #{endTime}
+ </if>
+ order by end_time desc limit 1
+ </select>
+
<delete id="deleteByWorkflowInstanceId">
delete
from t_ds_task_instance
diff --git
a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/runner/task/dependent/DependentAsyncTaskExecuteFunction.java
b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/runner/task/dependent/DependentAsyncTaskExecuteFunction.java
index c76f13ea4d..2a51f5d5bc 100644
---
a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/runner/task/dependent/DependentAsyncTaskExecuteFunction.java
+++
b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/runner/task/dependent/DependentAsyncTaskExecuteFunction.java
@@ -157,6 +157,12 @@ public class DependentAsyncTaskExecuteFunction implements
AsyncTaskExecuteFuncti
log.info("WorkflowName: {}",
processDefinition.getName());
log.info("TaskName: {}", "ALL");
log.info("DependentKey: {}",
dependentItem.getKey());
+ } else if (dependentItem.getDepTaskCode() ==
Constants.DEPENDENT_WORKFLOW_CODE) {
+ log.info("Add dependent task:");
+ log.info("DependentRelation: {}",
dependentTaskModel.getRelation());
+ log.info("ProjectName: {}", project.getName());
+ log.info("WorkflowName: {}",
processDefinition.getName());
+ log.info("DependentKey: {}",
dependentItem.getKey());
} else {
TaskDefinition taskDefinition =
taskDefinitionMap.get(dependentItem.getDepTaskCode());
if (taskDefinition == null) {
diff --git
a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/utils/DependentExecute.java
b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/utils/DependentExecute.java
index ddac0ddc40..b714685478 100644
---
a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/utils/DependentExecute.java
+++
b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/utils/DependentExecute.java
@@ -20,9 +20,16 @@ package org.apache.dolphinscheduler.server.master.utils;
import static
org.apache.dolphinscheduler.plugin.task.api.parameters.DependentParameters.DependentFailurePolicyEnum.DEPENDENT_FAILURE_WAITING;
import org.apache.dolphinscheduler.common.constants.Constants;
+import org.apache.dolphinscheduler.common.enums.Flag;
+import org.apache.dolphinscheduler.common.enums.TaskExecuteType;
import org.apache.dolphinscheduler.dao.entity.ProcessInstance;
+import org.apache.dolphinscheduler.dao.entity.ProcessTaskRelation;
+import org.apache.dolphinscheduler.dao.entity.TaskDefinition;
+import org.apache.dolphinscheduler.dao.entity.TaskDefinitionLog;
import org.apache.dolphinscheduler.dao.entity.TaskInstance;
import org.apache.dolphinscheduler.dao.repository.ProcessInstanceDao;
+import org.apache.dolphinscheduler.dao.repository.TaskDefinitionDao;
+import org.apache.dolphinscheduler.dao.repository.TaskDefinitionLogDao;
import org.apache.dolphinscheduler.dao.repository.TaskInstanceDao;
import org.apache.dolphinscheduler.plugin.task.api.enums.DependResult;
import org.apache.dolphinscheduler.plugin.task.api.enums.DependentRelation;
@@ -32,6 +39,7 @@ import
org.apache.dolphinscheduler.plugin.task.api.model.DependentItem;
import
org.apache.dolphinscheduler.plugin.task.api.parameters.DependentParameters;
import org.apache.dolphinscheduler.plugin.task.api.utils.DependentUtils;
import org.apache.dolphinscheduler.service.bean.SpringApplicationContext;
+import org.apache.dolphinscheduler.service.process.ProcessService;
import java.time.Duration;
import java.time.Instant;
@@ -40,6 +48,7 @@ import java.util.Date;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
+import java.util.stream.Collectors;
import lombok.extern.slf4j.Slf4j;
@@ -72,6 +81,22 @@ public class DependentExecute {
*/
private Map<String, DependResult> dependResultMap = new HashMap<>();
+ /**
+ * process service
+ */
+ private final ProcessService processService =
SpringApplicationContext.getBean(ProcessService.class);
+
+ /**
+ * task definition log dao
+ */
+ private final TaskDefinitionLogDao taskDefinitionLogDao =
+ SpringApplicationContext.getBean(TaskDefinitionLogDao.class);
+
+ /**
+ * task definition dao
+ */
+ private final TaskDefinitionDao taskDefinitionDao =
SpringApplicationContext.getBean(TaskDefinitionDao.class);
+
/**
* constructor
*
@@ -118,10 +143,13 @@ public class DependentExecute {
return DependResult.WAITING;
}
// need to check workflow for updates, so get all task and check
the task state
- if (dependentItem.getDepTaskCode() ==
Constants.DEPENDENT_ALL_TASK_CODE) {
+ if (dependentItem.getDepTaskCode() ==
Constants.DEPENDENT_WORKFLOW_CODE) {
result = dependResultByProcessInstance(processInstance);
+ } else if (dependentItem.getDepTaskCode() ==
Constants.DEPENDENT_ALL_TASK_CODE) {
+ result =
dependResultByAllTaskOfProcessInstance(processInstance, dateInterval, testFlag);
} else {
- result = getDependTaskResult(dependentItem.getDepTaskCode(),
processInstance, testFlag);
+ result = dependResultBySingleTaskInstance(processInstance,
dependentItem.getDepTaskCode(), dateInterval,
+ testFlag);
}
if (result != DependResult.SUCCESS) {
break;
@@ -131,7 +159,7 @@ public class DependentExecute {
}
/**
- * depend type = depend_all
+ * depend type = depend_work_flow
*
* @return
*/
@@ -142,6 +170,59 @@ public class DependentExecute {
if (processInstance.getState().isSuccess()) {
return DependResult.SUCCESS;
}
+ log.warn(
+ "The dependent workflow did not execute successfully, so
return depend failed. processCode: {}, processName: {}",
+ processInstance.getProcessDefinitionCode(),
processInstance.getName());
+ return DependResult.FAILED;
+ }
+
+ /**
+ * depend type = depend_all
+ *
+ * @return
+ */
+ private DependResult
dependResultByAllTaskOfProcessInstance(ProcessInstance processInstance,
+ DateInterval
dateInterval, int testFlag) {
+ if (!processInstance.getState().isFinished()) {
+ log.info("Wait for the dependent workflow to complete,
processCode: {}, processInstanceId: {}.",
+ processInstance.getProcessDefinitionCode(),
processInstance.getId());
+ return DependResult.WAITING;
+ }
+ if (processInstance.getState().isSuccess()) {
+ List<ProcessTaskRelation> processTaskRelations =
+
processService.findRelationByCode(processInstance.getProcessDefinitionCode(),
+ processInstance.getProcessDefinitionVersion());
+ List<TaskDefinitionLog> taskDefinitionLogs =
+
taskDefinitionLogDao.queryTaskDefineLogList(processTaskRelations);
+ Map<Long, String> taskDefinitionCodeMap =
+ taskDefinitionLogs.stream().filter(taskDefinitionLog ->
taskDefinitionLog.getFlag() == Flag.YES)
+
.collect(Collectors.toMap(TaskDefinitionLog::getCode,
TaskDefinitionLog::getName));
+
+ List<TaskInstance> taskInstanceList =
+
taskInstanceDao.queryLastTaskInstanceListIntervalByTaskCodes(taskDefinitionCodeMap.keySet(),
+ dateInterval, testFlag);
+ Map<Long, TaskExecutionStatus> taskExecutionStatusMap =
+ taskInstanceList.stream()
+ .filter(taskInstance ->
taskInstance.getTaskExecuteType() != TaskExecuteType.STREAM)
+
.collect(Collectors.toMap(TaskInstance::getTaskCode, TaskInstance::getState));
+
+ for (Long taskCode : taskDefinitionCodeMap.keySet()) {
+ if (!taskExecutionStatusMap.containsKey(taskCode)) {
+ log.warn(
+ "The task of the workflow is not being executed,
taskCode: {}, processInstanceId: {}, processName: {}.",
+ taskCode,
processInstance.getProcessDefinitionCode(), processInstance.getName());
+ return DependResult.FAILED;
+ } else {
+ if (!taskExecutionStatusMap.get(taskCode).isSuccess()) {
+ log.warn(
+ "The task of the workflow is not being
executed successfully, taskCode: {}, processInstanceId: {}, processName: {}.",
+ taskCode,
processInstance.getProcessDefinitionCode(), processInstance.getName());
+ return DependResult.FAILED;
+ }
+ }
+ }
+ return DependResult.SUCCESS;
+ }
return DependResult.FAILED;
}
@@ -180,6 +261,54 @@ public class DependentExecute {
return result;
}
+ /**
+ * depend type = depend_task
+ *
+ * @param processInstance last process instance in the date interval
+ * @param depTaskCode the dependent task code
+ * @param dateInterval date interval
+ * @param testFlag test flag
+ * @return depend result
+ */
+ private DependResult dependResultBySingleTaskInstance(ProcessInstance
processInstance, long depTaskCode,
+ DateInterval
dateInterval, int testFlag) {
+ TaskInstance taskInstance =
+
taskInstanceDao.queryLastTaskInstanceIntervalByTaskCode(depTaskCode,
dateInterval, testFlag);
+
+ if (taskInstance == null) {
+ TaskDefinition taskDefinition =
taskDefinitionDao.queryByCode(depTaskCode);
+
+ if (taskDefinition == null) {
+ log.error("The dependent task definition can not be find, so
return depend failed, taskCode: {}",
+ depTaskCode);
+ return DependResult.FAILED;
+ }
+
+ if (taskDefinition.getFlag() == Flag.NO) {
+ log.info(
+ "The dependent task is a forbidden task, so return
depend success. Task code: {}, task name: {}",
+ taskDefinition.getCode(), taskDefinition.getName());
+ return DependResult.SUCCESS;
+ }
+
+ if (!processInstance.getState().isFinished()) {
+ log.info("Wait for the dependent workflow to complete,
processCode: {}, processInstanceId: {}.",
+ processInstance.getProcessDefinitionCode(),
processInstance.getId());
+ return DependResult.WAITING;
+ }
+
+ return DependResult.FAILED;
+ } else {
+ if (TaskExecuteType.STREAM == taskInstance.getTaskExecuteType()) {
+ log.info(
+ "The dependent task is a streaming task, so return
depend success. Task code: {}, task name: {}.",
+ taskInstance.getTaskCode(), taskInstance.getName());
+ return DependResult.SUCCESS;
+ }
+ return getDependResultByState(taskInstance.getState());
+ }
+ }
+
/**
* find the last one process instance that :
* 1. manual run and finish between the interval
@@ -221,6 +350,9 @@ public class DependentExecute {
} else if (state.isSuccess()) {
return DependResult.SUCCESS;
} else {
+ log.warn(
+ "The dependent task were not executed successfully, so
return depend failed. Task code: {}, task name: {}.",
+ taskInstance.getTaskCode(), taskInstance.getName());
return DependResult.FAILED;
}
}
diff --git a/dolphinscheduler-ui/src/locales/en_US/project.ts
b/dolphinscheduler-ui/src/locales/en_US/project.ts
index 498ce8a538..b38462563b 100644
--- a/dolphinscheduler-ui/src/locales/en_US/project.ts
+++ b/dolphinscheduler-ui/src/locales/en_US/project.ts
@@ -876,6 +876,10 @@ export default {
child_node_instance: 'child node instance',
yarn_queue: 'Yarn Queue',
yarn_queue_tips: 'Please input yarn queue(optional)',
+ dependent_type: 'Dependency Type',
+ dependent_on_workflow: 'Dependent on workflow',
+ dependent_on_task: 'Dependent on task',
+
},
menu: {
fav: 'Favorites',
diff --git a/dolphinscheduler-ui/src/locales/zh_CN/project.ts
b/dolphinscheduler-ui/src/locales/zh_CN/project.ts
index 5b598cea17..efbd813c0a 100644
--- a/dolphinscheduler-ui/src/locales/zh_CN/project.ts
+++ b/dolphinscheduler-ui/src/locales/zh_CN/project.ts
@@ -850,6 +850,9 @@ export default {
child_node_instance: '子节点实例',
yarn_queue: 'Yarn队列',
yarn_queue_tips: '请输入Yarn队列(选填)',
+ dependent_type: '依赖类型',
+ dependent_on_workflow: '依赖于工作流',
+ dependent_on_task: '依赖于任务',
},
menu: {
fav: '收藏组件',
diff --git
a/dolphinscheduler-ui/src/views/projects/task/components/node/fields/use-dependent.ts
b/dolphinscheduler-ui/src/views/projects/task/components/node/fields/use-dependent.ts
index 433895f861..3fca7b417d 100644
---
a/dolphinscheduler-ui/src/views/projects/task/components/node/fields/use-dependent.ts
+++
b/dolphinscheduler-ui/src/views/projects/task/components/node/fields/use-dependent.ts
@@ -69,6 +69,17 @@ export function useDependent(model: { [field: string]: any
}): IJsonItem[] {
}
const selectOptions = ref([] as IDependTaskOptions[])
+ const DependentTypeOptions = [
+ {
+ value: 'DEPENDENT_ON_WORKFLOW',
+ label: t('project.node.dependent_on_workflow')
+ },
+ {
+ value: 'DEPENDENT_ON_TASK',
+ label: t('project.node.dependent_on_task')
+ }
+ ]
+
const CYCLE_LIST = [
{
value: 'month',
@@ -229,7 +240,7 @@ export function useDependent(model: { [field: string]: any
}): IJsonItem[] {
filterLabel: item.name
}))
taskList.unshift({
- value: 0,
+ value: -1,
label: 'ALL'
})
taskCache[processCode] = taskList
@@ -266,6 +277,13 @@ export function useDependent(model: { [field: string]: any
}): IJsonItem[] {
item.dependItemList?.forEach(
async (dependItem: IDependentItem, itemIndex: number) => {
itemListOptions.value[itemIndex] = {}
+
+ if (!dependItem.dependentType) {
+ if (dependItem.depTaskCode == 0)
+ dependItem.dependentType = 'DEPENDENT_ON_WORKFLOW'
+ else
+ dependItem.dependentType = 'DEPENDENT_ON_TASK'
+ }
if (dependItem.projectCode) {
itemListOptions.value[itemIndex].definitionCodeOptions =
await getProcessList(dependItem.projectCode)
@@ -298,6 +316,21 @@ export function useDependent(model: { [field: string]: any
}): IJsonItem[] {
field: 'dependItemList',
span: 18,
children: [
+ (j = 0) => ({
+ type: 'select',
+ field: 'dependentType',
+ name: t('project.node.dependent_type'),
+ span: 24,
+ props: {
+ onUpdateValue: (dependentType: string) => {
+ const item = model.dependTaskList[i].dependItemList[j]
+ if (item.definitionCode)
+ item.depTaskCode = dependentType === 'DEPENDENT_ON_WORKFLOW'
? 0 : -1
+ }
+ },
+ options: DependentTypeOptions,
+ value: 'DEPENDENT_ON_WORKFLOW'
+ }),
(j = 0) => ({
type: 'select',
field: 'projectCode',
@@ -353,7 +386,7 @@ export function useDependent(model: { [field: string]: any
}): IJsonItem[] {
const item = model.dependTaskList[i].dependItemList[j]
selectOptions.value[i].dependItemList[j].depTaskCodeOptions =
await getTaskList(item.projectCode, processCode)
- item.depTaskCode = 0
+ item.depTaskCode = item.dependentType ===
'DEPENDENT_ON_WORKFLOW' ? 0 : -1
}
},
options:
@@ -373,7 +406,10 @@ export function useDependent(model: { [field: string]: any
}): IJsonItem[] {
(j = 0) => ({
type: 'select',
field: 'depTaskCode',
- span: 24,
+ span: computed(() => {
+ const item = model.dependTaskList[i].dependItemList[j]
+ return item.dependentType === 'DEPENDENT_ON_WORKFLOW' ? 0 : 24
+ }),
name: t('project.node.task_name'),
props: {
filterable: true,
diff --git
a/dolphinscheduler-ui/src/views/projects/task/components/node/types.ts
b/dolphinscheduler-ui/src/views/projects/task/components/node/types.ts
index 84b6375a01..3f12a9e245 100644
--- a/dolphinscheduler-ui/src/views/projects/task/components/node/types.ts
+++ b/dolphinscheduler-ui/src/views/projects/task/components/node/types.ts
@@ -83,6 +83,7 @@ interface IResponseJsonItem extends Omit<IJsonItemParams,
'type'> {
}
interface IDependentItemOptions {
+ dependentTypeOptions?: IOption[]
definitionCodeOptions?: IOption[]
depTaskCodeOptions?: IOption[]
dateOptions?: IOption[]
@@ -99,6 +100,7 @@ interface IDependentItem {
definitionCode?: number
cycle?: 'month' | 'week' | 'day' | 'hour'
dateValue?: string
+ dependentType?: 'DEPENDENT_ON_WORKFLOW' | 'DEPENDENT_ON_TASK'
}
interface IDependTask {