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 ca9d7007f7 [Improvement-18568][API] Remove obsolete task
update-with-upstream API (#18569)
ca9d7007f7 is described below
commit ca9d7007f71c2fad17c37e723fd60fa6d565ba51
Author: Wenjun Ruan <[email protected]>
AuthorDate: Mon Aug 24 21:04:11 2026 +0800
[Improvement-18568][API] Remove obsolete task update-with-upstream API
(#18569)
---
docs/docs/en/guide/upgrade/incompatible.md | 1 +
docs/docs/zh/guide/upgrade/incompatible.md | 1 +
.../api/controller/TaskDefinitionController.java | 33 ---
.../apache/dolphinscheduler/api/enums/Status.java | 5 -
.../api/service/TaskDefinitionService.java | 16 --
.../service/impl/TaskDefinitionServiceImpl.java | 316 ---------------------
.../api/service/TaskDefinitionServiceImplTest.java | 83 ------
.../dao/mapper/WorkflowTaskRelationMapper.java | 17 --
.../dao/repository/WorkflowTaskRelationDao.java | 4 -
.../impl/WorkflowTaskRelationDaoImpl.java | 11 -
.../dao/mapper/WorkflowTaskRelationMapper.xml | 22 --
11 files changed, 2 insertions(+), 507 deletions(-)
diff --git a/docs/docs/en/guide/upgrade/incompatible.md
b/docs/docs/en/guide/upgrade/incompatible.md
index dd8e482802..55bc8a6041 100644
--- a/docs/docs/en/guide/upgrade/incompatible.md
+++ b/docs/docs/en/guide/upgrade/incompatible.md
@@ -48,4 +48,5 @@ This document records the incompatible updates between each
version. You need to
* Add the `missed_fire_policy` column to `t_ds_schedules`. Existing schedules
default to `FIRE_ALL_MISSED` to preserve the previous Quartz `IgnoreMisfires`
behavior. ([#18464](https://github.com/apache/dolphinscheduler/pull/18464))
* Remove the obsolete Dynamic Task query API.
([#18556](https://github.com/apache/dolphinscheduler/issues/18556))
+* Remove the obsolete task update-with-upstream API `PUT
/projects/{projectCode}/task-definition/{code}/with-upstream`.
([#18568](https://github.com/apache/dolphinscheduler/issues/18568))
diff --git a/docs/docs/zh/guide/upgrade/incompatible.md
b/docs/docs/zh/guide/upgrade/incompatible.md
index ee324952dd..ba52dd5fe8 100644
--- a/docs/docs/zh/guide/upgrade/incompatible.md
+++ b/docs/docs/zh/guide/upgrade/incompatible.md
@@ -48,4 +48,5 @@
* 为 `t_ds_schedules` 表新增 `missed_fire_policy` 字段。现有定时默认使用
`FIRE_ALL_MISSED`,以保持原有 Quartz `IgnoreMisfires`
行为。([#18464](https://github.com/apache/dolphinscheduler/pull/18464))
* 移除已废弃的 Dynamic Task
查询接口。([#18556](https://github.com/apache/dolphinscheduler/issues/18556))
+* 移除已废弃的任务及其上游关系更新接口 `PUT
/projects/{projectCode}/task-definition/{code}/with-upstream`。([#18568](https://github.com/apache/dolphinscheduler/issues/18568))
diff --git
a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/TaskDefinitionController.java
b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/TaskDefinitionController.java
index 5610412a6d..8c96bff9c8 100644
---
a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/TaskDefinitionController.java
+++
b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/TaskDefinitionController.java
@@ -23,7 +23,6 @@ import static
org.apache.dolphinscheduler.api.enums.Status.QUERY_DETAIL_OF_TASK_
import static
org.apache.dolphinscheduler.api.enums.Status.QUERY_TASK_DEFINITION_VERSIONS_ERROR;
import static
org.apache.dolphinscheduler.api.enums.Status.RELEASE_TASK_DEFINITION_ERROR;
import static
org.apache.dolphinscheduler.api.enums.Status.SWITCH_TASK_DEFINITION_VERSION_ERROR;
-import static
org.apache.dolphinscheduler.api.enums.Status.UPDATE_TASK_DEFINITION_ERROR;
import org.apache.dolphinscheduler.api.audit.OperatorLog;
import org.apache.dolphinscheduler.api.audit.enums.AuditType;
@@ -43,7 +42,6 @@ import org.springframework.web.bind.annotation.DeleteMapping;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.PostMapping;
-import org.springframework.web.bind.annotation.PutMapping;
import org.springframework.web.bind.annotation.RequestAttribute;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RequestParam;
@@ -64,37 +62,6 @@ public class TaskDefinitionController extends BaseController
{
@Autowired
private TaskDefinitionService taskDefinitionService;
- /**
- * update task definition
- *
- * @param loginUser login user
- * @param projectCode project code
- * @param code task definition code
- * @param taskDefinitionJsonObj task definition json object
- * @param upstreamCodes upstream task codes, sep comma
- * @return update result code
- */
- @Operation(summary = "updateWithUpstream", description =
"UPDATE_TASK_DEFINITION_NOTES")
- @Parameters({
- @Parameter(name = "projectCode", description = "PROJECT_CODE",
required = true, schema = @Schema(implementation = long.class)),
- @Parameter(name = "code", description = "TASK_DEFINITION_CODE",
required = true, schema = @Schema(implementation = long.class, example = "1")),
- @Parameter(name = "taskDefinitionJsonObj", description =
"TASK_DEFINITION_JSON", required = true, schema = @Schema(implementation =
String.class)),
- @Parameter(name = "upstreamCodes", description = "UPSTREAM_CODES",
required = false, schema = @Schema(implementation = String.class))
- })
- @PutMapping(value = "/{code}/with-upstream")
- @ResponseStatus(HttpStatus.OK)
- @ApiException(UPDATE_TASK_DEFINITION_ERROR)
- @OperatorLog(auditType = AuditType.TASK_UPDATE)
- public Result<Long> updateTaskWithUpstream(@Parameter(hidden = true)
@RequestAttribute(value = Constants.SESSION_USER) User loginUser,
- @Parameter(name =
"projectCode", description = "PROJECT_CODE", required = true) @PathVariable
long projectCode,
- @PathVariable(value = "code")
long code,
- @RequestParam(value =
"taskDefinitionJsonObj", required = true) String taskDefinitionJsonObj,
- @RequestParam(value =
"upstreamCodes", required = false) String upstreamCodes) {
- Long updatedTaskCode =
taskDefinitionService.updateTaskWithUpstream(loginUser, projectCode, code,
- taskDefinitionJsonObj, upstreamCodes);
- return Result.success(updatedTaskCode);
- }
-
/**
* query task definition version paging list info
*
diff --git
a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/enums/Status.java
b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/enums/Status.java
index b47a9b1089..05f951a16e 100644
---
a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/enums/Status.java
+++
b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/enums/Status.java
@@ -396,7 +396,6 @@ public enum Status {
MAIN_TABLE_USING_VERSION(50053, "the version that the master table is
using", "主表正在使用该版本"),
PROJECT_WORKFLOW_NOT_MATCH(50054, "the project and the workflow is not
match", "项目和工作流不匹配"),
DELETE_EDGE_ERROR(50055, "delete edge error", "删除工作流任务连接线错误"),
- NOT_SUPPORT_UPDATE_TASK_DEFINITION(50056, "task state does not support
modification", "当前任务不支持修改"),
TASK_DEFINITION_NOT_MODIFY_ERROR(50057, "task [{0}] definition not modify
error", "任务[{0}]定义未修改错误"),
BATCH_EXECUTE_WORKFLOW_INSTANCE_ERROR(50058, "change workflow instance
status error: {0}", "修改工作实例状态错误: {0}"),
START_TASK_INSTANCE_ERROR(50059, "start task instance error", "运行任务流实例错误"),
@@ -406,16 +405,12 @@ public enum Status {
TASK_DEFINITION_NOT_CHANGE(50063, "task definition {0} do not change",
"任务定义 {0} 没有变化"),
TASK_DEFINITION_NOT_EXISTS(50064, "task definition {0} do not exists",
"任务定义 {0} 不存在"),
UPDATE_UPSTREAM_TASK_WORKFLOW_RELATION_ERROR(50065, "update task upstream
relation error", "更新任务上游关系错误"),
- CREATE_WORKFLOW_TASK_RELATION_LOG_ERROR(50066, "create workflow task
relation log {0}-{1} error",
- "创建任务关系日志 {0}-{1} 错误"),
WORKFLOW_TASK_RELATION_NOT_EXPECT(50067, "workflow task relation number
not expect, expect {0} but get {1}",
"工作流任务关系数量不符合预期,预期 {0} 但是实际 {1}"),
WORKFLOW_TASK_RELATION_BATCH_DELETE_ERROR(50068, "batch delete workflow
task relation {0} error",
"批量删除工作流任务关系 {0} 错误"),
WORKFLOW_TASK_RELATION_BATCH_CREATE_ERROR(50069, "batch create workflow
task relation {0} error",
"批量创建工作流任务关系 {0} 错误"),
- WORKFLOW_TASK_RELATION_BATCH_UPDATE_ERROR(50070, "batch update workflow
task relation error",
- "批量修改工作流任务关系错误"),
UPSTREAM_TASK_NOT_EXISTS(50071, "upstream task want to set dependence do
not exists {0}", "指定的上游任务 {0} 不存在"),
WORKFLOW_INSTANCE_IS_NOT_FINISHED(50071, "the workflow instance is not
finished, can not do this operation",
diff --git
a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/TaskDefinitionService.java
b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/TaskDefinitionService.java
index f2750efdca..b4e150be0e 100644
---
a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/TaskDefinitionService.java
+++
b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/TaskDefinitionService.java
@@ -51,22 +51,6 @@ public interface TaskDefinitionService {
TaskDefinition getTaskDefinition(User loginUser,
long taskCode);
- /**
- * update task definition and upstream
- *
- * @param loginUser login user
- * @param projectCode project code
- * @param taskCode task definition code
- * @param taskDefinitionJsonObj task definition json object
- * @param upstreamCodes upstream task codes, sep comma
- * @return updated task code
- */
- Long updateTaskWithUpstream(User loginUser,
- long projectCode,
- long taskCode,
- String taskDefinitionJsonObj,
- String upstreamCodes);
-
/**
* update task definition
*
diff --git
a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/TaskDefinitionServiceImpl.java
b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/TaskDefinitionServiceImpl.java
index b39b0de17c..cb1a8afc51 100644
---
a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/TaskDefinitionServiceImpl.java
+++
b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/TaskDefinitionServiceImpl.java
@@ -19,7 +19,6 @@ package org.apache.dolphinscheduler.api.service.impl;
import static
org.apache.dolphinscheduler.api.constants.ApiFuncIdentificationConstant.TASK_DEFINITION;
import static
org.apache.dolphinscheduler.api.constants.ApiFuncIdentificationConstant.TASK_VERSION_VIEW;
-import static
org.apache.dolphinscheduler.plugin.task.api.TaskPluginManager.checkTaskParameters;
import org.apache.dolphinscheduler.api.enums.Status;
import org.apache.dolphinscheduler.api.exceptions.ServiceException;
@@ -33,44 +32,33 @@ import org.apache.dolphinscheduler.api.utils.Result;
import org.apache.dolphinscheduler.api.vo.TaskDefinitionVO;
import org.apache.dolphinscheduler.common.constants.Constants;
import org.apache.dolphinscheduler.common.enums.AuthorizationType;
-import org.apache.dolphinscheduler.common.enums.ConditionType;
import org.apache.dolphinscheduler.common.enums.Flag;
import org.apache.dolphinscheduler.common.enums.ReleaseState;
-import org.apache.dolphinscheduler.common.enums.TaskExecuteType;
-import org.apache.dolphinscheduler.common.enums.TimeoutFlag;
import org.apache.dolphinscheduler.common.utils.CodeGenerateUtils;
-import org.apache.dolphinscheduler.common.utils.JSONUtils;
import org.apache.dolphinscheduler.dao.entity.Project;
import org.apache.dolphinscheduler.dao.entity.TaskDefinition;
import org.apache.dolphinscheduler.dao.entity.TaskDefinitionLog;
import org.apache.dolphinscheduler.dao.entity.User;
import org.apache.dolphinscheduler.dao.entity.WorkflowDefinition;
-import org.apache.dolphinscheduler.dao.entity.WorkflowDefinitionLog;
import org.apache.dolphinscheduler.dao.entity.WorkflowTaskRelation;
import org.apache.dolphinscheduler.dao.entity.WorkflowTaskRelationLog;
import org.apache.dolphinscheduler.dao.mapper.TaskDefinitionLogMapper;
-import org.apache.dolphinscheduler.dao.mapper.WorkflowDefinitionLogMapper;
import org.apache.dolphinscheduler.dao.repository.ProjectDao;
import org.apache.dolphinscheduler.dao.repository.TaskDefinitionDao;
import org.apache.dolphinscheduler.dao.repository.WorkflowDefinitionDao;
import org.apache.dolphinscheduler.dao.repository.WorkflowTaskRelationDao;
-import org.apache.dolphinscheduler.dao.repository.WorkflowTaskRelationLogDao;
import org.apache.dolphinscheduler.service.process.ProcessService;
import org.apache.commons.collections4.CollectionUtils;
-import org.apache.commons.collections4.MapUtils;
import org.apache.commons.lang3.StringUtils;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
import java.util.Date;
-import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
-import java.util.Map;
import java.util.Set;
-import java.util.function.Function;
import java.util.stream.Collectors;
import lombok.extern.slf4j.Slf4j;
@@ -102,9 +90,6 @@ public class TaskDefinitionServiceImpl extends
BaseServiceImpl implements TaskDe
@Autowired
private WorkflowTaskRelationDao workflowTaskRelationDao;
- @Autowired
- private WorkflowTaskRelationLogDao workflowTaskRelationLogDao;
-
@Autowired
private WorkflowTaskRelationService workflowTaskRelationService;
@@ -114,12 +99,8 @@ public class TaskDefinitionServiceImpl extends
BaseServiceImpl implements TaskDe
@Autowired
private ProcessService processService;
- @Autowired
- private WorkflowDefinitionLogMapper workflowDefinitionLogMapper;
-
@Autowired
private TaskDatasourcePermissionChecker taskDatasourcePermissionChecker;
-
/**
* query task definition
*
@@ -197,303 +178,6 @@ public class TaskDefinitionServiceImpl extends
BaseServiceImpl implements TaskDe
return taskDefinition;
}
- /**
- * Update the task body and cascade through workflow task relations.
- *
- * @return the persisted {@link TaskDefinitionLog}, or {@code null} when
the
- * task body is unchanged (callers may still need to apply upstream
- * changes). All other failure modes throw {@link
ServiceException}.
- */
- private TaskDefinitionLog updateTask(User loginUser, long projectCode,
long taskCode,
- String taskDefinitionJsonObj) {
- Project project = projectDao.queryByCode(projectCode);
-
- // check if user have write perm for project
- projectService.checkHasProjectWritePermissionThrowException(loginUser,
project);
-
- TaskDefinition taskDefinition =
taskDefinitionDao.queryByCode(taskCode);
- if (taskDefinition == null) {
- log.error("Task definition does not exist,
taskDefinitionCode:{}.", taskCode);
- throw new ServiceException(Status.TASK_DEFINE_NOT_EXIST,
String.valueOf(taskCode));
- }
- if (processService.isTaskOnline(taskCode) && taskDefinition.getFlag()
== Flag.YES) {
- // if stream, can update task definition without online check
- if (taskDefinition.getTaskExecuteType() != TaskExecuteType.STREAM)
{
- log.warn("Only {} type task can be updated without online
check, taskDefinitionCode:{}.",
- TaskExecuteType.STREAM, taskCode);
- throw new
ServiceException(Status.NOT_SUPPORT_UPDATE_TASK_DEFINITION);
- }
- }
- TaskDefinitionLog taskDefinitionToUpdate =
- JSONUtils.parseObject(taskDefinitionJsonObj,
TaskDefinitionLog.class);
- if (taskDefinitionToUpdate == null) {
- log.warn("Parameter taskDefinitionJson is invalid.");
- throw new ServiceException(Status.DATA_IS_NOT_VALID,
taskDefinitionJsonObj);
- }
- if (!checkTaskParameters(taskDefinitionToUpdate.getTaskType(),
taskDefinitionToUpdate.getTaskParams())) {
- throw new
ServiceException(Status.WORKFLOW_NODE_S_PARAMETER_INVALID,
taskDefinitionToUpdate.getName());
- }
- taskDatasourcePermissionChecker.checkPermission(loginUser,
Collections.singletonList(taskDefinitionToUpdate));
- if (TimeoutFlag.CLOSE == taskDefinition.getTimeoutFlag()) {
- taskDefinition.setTimeoutNotifyStrategy(null);
- }
- if (taskDefinition.equals(taskDefinitionToUpdate)) {
- log.warn("Task definition does not need update because no change,
taskDefinitionCode:{}.", taskCode);
- return null;
- }
- Integer version =
taskDefinitionLogMapper.queryMaxVersionForDefinition(taskCode);
- if (version == null || version == 0) {
- log.error("Max version task definitionLog can not be found in
database, taskDefinitionCode:{}.",
- taskCode);
- throw new ServiceException(Status.DATA_IS_NOT_VALID, taskCode);
- }
- Date now = new Date();
- taskDefinitionToUpdate.setCode(taskCode);
- taskDefinitionToUpdate.setId(taskDefinition.getId());
- taskDefinitionToUpdate.setProjectCode(projectCode);
- taskDefinitionToUpdate.setUserId(taskDefinition.getUserId());
- taskDefinitionToUpdate.setVersion(++version);
-
taskDefinitionToUpdate.setTaskType(taskDefinitionToUpdate.getTaskType().toUpperCase());
- taskDefinitionToUpdate.setUpdateTime(now);
- boolean updateSuccess =
taskDefinitionDao.updateById(taskDefinitionToUpdate);
- taskDefinitionToUpdate.setOperator(loginUser.getId());
- taskDefinitionToUpdate.setOperateTime(now);
- taskDefinitionToUpdate.setCreateTime(now);
- taskDefinitionToUpdate.setId(null);
- int insert = taskDefinitionLogMapper.insert(taskDefinitionToUpdate);
- if (!updateSuccess || insert != 1) {
- log.error("Update task definition or definitionLog error,
projectCode:{}, taskDefinitionCode:{}.",
- projectCode, taskCode);
- throw new ServiceException(Status.UPDATE_TASK_DEFINITION_ERROR);
- }
- log.info(
- "Update task definition and definitionLog complete,
projectCode:{}, taskDefinitionCode:{}, newTaskVersion:{}.",
- projectCode, taskCode, taskDefinitionToUpdate.getVersion());
- // update workflow task relation
- List<WorkflowTaskRelation> workflowTaskRelations =
workflowTaskRelationDao
-
.queryWorkflowTaskRelationByTaskCodeAndTaskVersion(taskDefinitionToUpdate.getCode(),
- taskDefinition.getVersion());
- if (CollectionUtils.isNotEmpty(workflowTaskRelations)) {
- Map<Long, List<WorkflowTaskRelation>>
workflowTaskRelationGroupList = workflowTaskRelations.stream()
-
.collect(Collectors.groupingBy(WorkflowTaskRelation::getWorkflowDefinitionCode));
- for (Map.Entry<Long, List<WorkflowTaskRelation>>
workflowTaskRelationMap : workflowTaskRelationGroupList
- .entrySet()) {
- Long workflowDefinitionCode = workflowTaskRelationMap.getKey();
- int workflowDefinitionVersion =
-
workflowDefinitionLogMapper.queryMaxVersionForDefinition(workflowDefinitionCode)
- + 1;
- List<WorkflowTaskRelation> workflowTaskRelationList =
workflowTaskRelationMap.getValue();
- for (WorkflowTaskRelation workflowTaskRelation :
workflowTaskRelationList) {
- if (taskCode == workflowTaskRelation.getPreTaskCode()) {
- workflowTaskRelation.setPreTaskVersion(version);
- } else if (taskCode ==
workflowTaskRelation.getPostTaskCode()) {
- workflowTaskRelation.setPostTaskVersion(version);
- }
-
workflowTaskRelation.setWorkflowDefinitionVersion(workflowDefinitionVersion);
- if
(!workflowTaskRelationDao.updateWorkflowTaskRelationTaskVersion(workflowTaskRelation))
{
- log.error("batch update workflow task relation error,
projectCode:{}, taskDefinitionCode:{}.",
- projectCode, taskCode);
- throw new
ServiceException(Status.WORKFLOW_TASK_RELATION_BATCH_UPDATE_ERROR);
- }
- WorkflowTaskRelationLog workflowTaskRelationLog = new
WorkflowTaskRelationLog(workflowTaskRelation);
- workflowTaskRelationLog.setOperator(loginUser.getId());
- workflowTaskRelationLog.setId(null);
- workflowTaskRelationLog.setOperateTime(now);
- int insertWorkflowTaskRelationLogCount =
workflowTaskRelationLogDao.insert(workflowTaskRelationLog);
- if (insertWorkflowTaskRelationLogCount != 1) {
- log.error("batch update workflow task relation error,
projectCode:{}, taskDefinitionCode:{}.",
- projectCode, taskCode);
- throw new
ServiceException(Status.CREATE_WORKFLOW_TASK_RELATION_LOG_ERROR);
- }
- }
- WorkflowDefinition workflowDefinition =
-
workflowDefinitionDao.queryByCode(workflowDefinitionCode).orElse(null);
- workflowDefinition.setVersion(workflowDefinitionVersion);
- workflowDefinition.setUpdateTime(now);
- workflowDefinition.setUserId(loginUser.getId());
- // update workflow definition
- boolean updateWorkflowDefinitionSuccess =
workflowDefinitionDao.updateById(workflowDefinition);
- WorkflowDefinitionLog workflowDefinitionLog = new
WorkflowDefinitionLog(workflowDefinition);
- workflowDefinitionLog.setOperateTime(now);
- workflowDefinitionLog.setId(null);
- workflowDefinitionLog.setOperator(loginUser.getId());
- int insertWorkflowDefinitionLogCount =
workflowDefinitionLogMapper.insert(workflowDefinitionLog);
- if (!updateWorkflowDefinitionSuccess ||
insertWorkflowDefinitionLogCount != 1) {
- throw new
ServiceException(Status.UPDATE_WORKFLOW_DEFINITION_ERROR);
- }
- }
- }
- return taskDefinitionToUpdate;
- }
-
- /**
- * update task definition and upstream
- *
- * @param loginUser login user
- * @param projectCode project code
- * @param taskCode task definition code
- * @param taskDefinitionJsonObj task definition json object
- * @param upstreamCodes upstream task codes, sep comma
- * @return updated task code
- */
- @Override
- public Long updateTaskWithUpstream(User loginUser, long projectCode, long
taskCode,
- String taskDefinitionJsonObj, String
upstreamCodes) {
- TaskDefinitionLog taskDefinitionToUpdate =
- updateTask(loginUser, projectCode, taskCode,
taskDefinitionJsonObj);
- List<WorkflowTaskRelation> upstreamTaskRelations =
- workflowTaskRelationDao.queryUpstreamByCode(projectCode,
taskCode);
- Set<Long> upstreamCodeSet =
-
upstreamTaskRelations.stream().map(WorkflowTaskRelation::getPreTaskCode).collect(Collectors.toSet());
- Set<Long> upstreamTaskCodes = Collections.emptySet();
- if (StringUtils.isNotEmpty(upstreamCodes)) {
- upstreamTaskCodes =
Arrays.stream(upstreamCodes.split(Constants.COMMA)).map(Long::parseLong)
- .collect(Collectors.toSet());
- }
- if (CollectionUtils.isEqualCollection(upstreamCodeSet,
upstreamTaskCodes) && taskDefinitionToUpdate == null) {
- return taskCode;
- }
- Map<Long, TaskDefinition> queryUpStreamTaskCodeMap;
- if (CollectionUtils.isNotEmpty(upstreamTaskCodes)) {
- List<TaskDefinition> upstreamTaskDefinitionList =
taskDefinitionDao.queryByCodes(upstreamTaskCodes);
- queryUpStreamTaskCodeMap = upstreamTaskDefinitionList.stream()
- .collect(Collectors.toMap(TaskDefinition::getCode,
taskDefinition -> taskDefinition));
- // upstreamTaskCodes - queryUpStreamTaskCodeMap.keySet
- upstreamTaskCodes.removeAll(queryUpStreamTaskCodeMap.keySet());
- if (CollectionUtils.isNotEmpty(upstreamTaskCodes)) {
- String notExistTaskCodes = StringUtils.join(upstreamTaskCodes,
Constants.COMMA);
- log.error("Some task definitions in parameter
upstreamTaskCodes do not exist, notExistTaskCodes:{}.",
- notExistTaskCodes);
- throw new ServiceException(Status.TASK_DEFINE_NOT_EXIST,
notExistTaskCodes);
- }
- } else {
- queryUpStreamTaskCodeMap = new HashMap<>();
- }
- if (MapUtils.isNotEmpty(queryUpStreamTaskCodeMap)) {
- WorkflowTaskRelation taskRelation = upstreamTaskRelations.get(0);
- List<WorkflowTaskRelation> workflowTaskRelations =
-
workflowTaskRelationDao.queryByWorkflowDefinitionCode(taskRelation.getWorkflowDefinitionCode());
-
- // set upstream code list
- updateUpstreamTask(new
HashSet<>(queryUpStreamTaskCodeMap.keySet()),
- taskCode, projectCode,
taskRelation.getWorkflowDefinitionCode(), loginUser);
-
- List<WorkflowTaskRelation> workflowTaskRelationList =
Lists.newArrayList(workflowTaskRelations);
- List<WorkflowTaskRelation> relationList = Lists.newArrayList();
- for (WorkflowTaskRelation workflowTaskRelation :
workflowTaskRelationList) {
- if (workflowTaskRelation.getPostTaskCode() == taskCode) {
- if
(queryUpStreamTaskCodeMap.containsKey(workflowTaskRelation.getPreTaskCode())
- && workflowTaskRelation.getPreTaskCode() != 0L) {
-
queryUpStreamTaskCodeMap.remove(workflowTaskRelation.getPreTaskCode());
- } else {
- workflowTaskRelation.setPreTaskCode(0L);
- workflowTaskRelation.setPreTaskVersion(0);
- relationList.add(workflowTaskRelation);
- }
- }
- }
- workflowTaskRelationList.removeAll(relationList);
- for (Map.Entry<Long, TaskDefinition> queryUpStreamTask :
queryUpStreamTaskCodeMap.entrySet()) {
- taskRelation.setPreTaskCode(queryUpStreamTask.getKey());
-
taskRelation.setPreTaskVersion(queryUpStreamTask.getValue().getVersion());
- workflowTaskRelationList.add(taskRelation);
- }
- if (MapUtils.isEmpty(queryUpStreamTaskCodeMap) &&
CollectionUtils.isNotEmpty(workflowTaskRelationList)) {
- workflowTaskRelationList.add(workflowTaskRelationList.get(0));
- }
- }
- log.info(
- "Update task with upstream tasks complete, projectCode:{},
taskDefinitionCode:{}, upstreamTaskCodes:{}.",
- projectCode, taskCode, upstreamTaskCodes);
- return taskCode;
- }
-
- private void updateUpstreamTask(Set<Long> allPreTaskCodeSet, long
taskCode, long projectCode,
- long workflowDefinitionCode, User
loginUser) {
- // query all workflow task relation
- List<WorkflowTaskRelation> hadWorkflowTaskRelationList =
workflowTaskRelationDao
- .queryUpstreamByCode(projectCode, taskCode);
- // remove pre
- Set<Long> removePreTaskSet = new HashSet<>();
- List<WorkflowTaskRelation> removePreTaskList = new ArrayList<>();
- // add pre
- Set<Long> addPreTaskSet = new HashSet<>();
- List<WorkflowTaskRelation> addPreTaskList = new ArrayList<>();
-
- List<WorkflowTaskRelationLog> workflowTaskRelationLogList = new
ArrayList<>();
-
- // filter all workflow task relation
- if (CollectionUtils.isNotEmpty(hadWorkflowTaskRelationList)) {
- for (WorkflowTaskRelation workflowTaskRelation :
hadWorkflowTaskRelationList) {
- if (workflowTaskRelation.getPreTaskCode() == 0) {
- continue;
- }
- // had
- if
(allPreTaskCodeSet.contains(workflowTaskRelation.getPreTaskCode())) {
-
allPreTaskCodeSet.remove(workflowTaskRelation.getPreTaskCode());
- } else {
- // remove
-
removePreTaskSet.add(workflowTaskRelation.getPreTaskCode());
- workflowTaskRelation.setPreTaskCode(0);
- workflowTaskRelation.setPreTaskVersion(0);
- removePreTaskList.add(workflowTaskRelation);
-
workflowTaskRelationLogList.add(createWorkflowTaskRelationLog(loginUser,
workflowTaskRelation));
- }
- }
- }
- // add
- if (allPreTaskCodeSet.size() != 0) {
- addPreTaskSet.addAll(allPreTaskCodeSet);
- }
- // get add task code map
- allPreTaskCodeSet.add(Long.valueOf(taskCode));
- List<TaskDefinition> taskDefinitionList =
taskDefinitionDao.queryByCodes(allPreTaskCodeSet);
- Map<Long, TaskDefinition> taskCodeMap =
taskDefinitionList.stream().collect(Collectors
- .toMap(TaskDefinition::getCode, Function.identity(), (a, b) ->
a));
-
- WorkflowDefinition workflowDefinition =
workflowDefinitionDao.queryByCode(workflowDefinitionCode).orElse(null);
- TaskDefinition taskDefinition = taskCodeMap.get(taskCode);
-
- for (Long preTaskCode : addPreTaskSet) {
- TaskDefinition preTaskRelation = taskCodeMap.get(preTaskCode);
- WorkflowTaskRelation workflowTaskRelation = new
WorkflowTaskRelation(
- null, workflowDefinition.getVersion(), projectCode,
workflowDefinition.getCode(),
- preTaskRelation.getCode(), preTaskRelation.getVersion(),
- taskDefinition.getCode(), taskDefinition.getVersion(),
ConditionType.NONE, "{}");
- addPreTaskList.add(workflowTaskRelation);
-
workflowTaskRelationLogList.add(createWorkflowTaskRelationLog(loginUser,
workflowTaskRelation));
- }
- int insert = 0;
- int remove = 0;
- int log = 0;
- // insert workflow task relation table data
- if (CollectionUtils.isNotEmpty(addPreTaskList)) {
- insert = workflowTaskRelationDao.batchInsert(addPreTaskList);
- }
- if (CollectionUtils.isNotEmpty(removePreTaskList)) {
- for (WorkflowTaskRelation workflowTaskRelation :
removePreTaskList) {
- remove +=
workflowTaskRelationDao.updateById(workflowTaskRelation) ? 1 : 0;
- }
- }
- if (CollectionUtils.isNotEmpty(workflowTaskRelationLogList)) {
- log =
workflowTaskRelationLogDao.batchInsert(workflowTaskRelationLogList);
- }
- if (insert + remove != log) {
- throw new RuntimeException("updateUpstreamTask error");
- }
- }
-
- private WorkflowTaskRelationLog createWorkflowTaskRelationLog(User
loginUser,
-
WorkflowTaskRelation workflowTaskRelation) {
- Date now = new Date();
- WorkflowTaskRelationLog workflowTaskRelationLog = new
WorkflowTaskRelationLog(workflowTaskRelation);
- workflowTaskRelationLog.setOperator(loginUser.getId());
- workflowTaskRelationLog.setOperateTime(now);
- workflowTaskRelationLog.setCreateTime(now);
- workflowTaskRelationLog.setUpdateTime(now);
- return workflowTaskRelationLog;
- }
-
/**
* switch task definition
*
diff --git
a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/TaskDefinitionServiceImplTest.java
b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/TaskDefinitionServiceImplTest.java
index f34c1c7e81..364f89d722 100644
---
a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/TaskDefinitionServiceImplTest.java
+++
b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/TaskDefinitionServiceImplTest.java
@@ -56,13 +56,10 @@ import
org.apache.dolphinscheduler.dao.repository.ProjectDao;
import org.apache.dolphinscheduler.dao.repository.TaskDefinitionDao;
import org.apache.dolphinscheduler.dao.repository.WorkflowDefinitionDao;
import org.apache.dolphinscheduler.dao.repository.WorkflowTaskRelationDao;
-import org.apache.dolphinscheduler.dao.repository.WorkflowTaskRelationLogDao;
-import org.apache.dolphinscheduler.plugin.task.api.TaskPluginManager;
import org.apache.dolphinscheduler.service.process.ProcessService;
import org.apache.dolphinscheduler.service.process.ProcessServiceImpl;
import java.util.ArrayList;
-import java.util.Arrays;
import java.util.List;
import java.util.Map;
import java.util.Optional;
@@ -73,7 +70,6 @@ import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.extension.ExtendWith;
import org.mockito.InjectMocks;
import org.mockito.Mock;
-import org.mockito.MockedStatic;
import org.mockito.Mockito;
import org.mockito.junit.jupiter.MockitoExtension;
import org.mockito.junit.jupiter.MockitoSettings;
@@ -128,15 +124,11 @@ public class TaskDefinitionServiceImplTest {
@Mock
private WorkflowDefinitionDao workflowDefinitionDao;
- @Mock
- private WorkflowTaskRelationLogDao workflowTaskRelationLogDao;
-
private static final String TASK_PARAMETER =
"{\"resourceList\":[],\"localParams\":[],\"rawScript\":\"echo
1\",\"conditionResult\":{\"successNode\":[\"\"],\"failedNode\":[\"\"]},\"dependence\":{}}";;
private static final long PROJECT_CODE = 1L;
private static final long PROCESS_DEFINITION_CODE = 2L;
private static final long TASK_CODE = 3L;
- private static final String UPSTREAM_CODE = "3,5";
private static final int VERSION = 1;
private static final int RESOURCE_RATE = -1;
protected User user;
@@ -385,60 +377,6 @@ public class TaskDefinitionServiceImplTest {
Assertions.assertDoesNotThrow(() ->
taskDefinitionService.getTaskDefinition(user, TASK_CODE));
}
- @Test
- public void testUpdateTaskWithUpstream() {
- try (
- MockedStatic<TaskPluginManager> taskPluginManagerMockedStatic =
- Mockito.mockStatic(TaskPluginManager.class)) {
- taskPluginManagerMockedStatic
- .when(() ->
TaskPluginManager.checkTaskParameters(Mockito.any(), Mockito.any()))
- .thenReturn(true);
- String taskDefinitionJson = getTaskDefinitionJson();
- TaskDefinition taskDefinition = getTaskDefinition();
- taskDefinition.setFlag(Flag.NO);
- TaskDefinition taskDefinitionSecond = getTaskDefinition();
- taskDefinitionSecond.setCode(5);
-
- user.setUserType(UserType.ADMIN_USER);
-
when(projectDao.queryByCode(PROJECT_CODE)).thenReturn(getProject());
-
Mockito.doNothing().when(projectService).checkHasProjectWritePermissionThrowException(eq(user),
- eq(getProject()));
-
when(taskDefinitionDao.queryByCode(TASK_CODE)).thenReturn(taskDefinition);
-
when(taskDefinitionLogMapper.queryMaxVersionForDefinition(TASK_CODE)).thenReturn(1);
- when(taskDefinitionDao.updateById(Mockito.any())).thenReturn(true);
- when(taskDefinitionLogMapper.insert(Mockito.any())).thenReturn(1);
-
- when(taskDefinitionDao.queryByCodes(Mockito.anySet()))
- .thenReturn(Arrays.asList(taskDefinition,
taskDefinitionSecond));
-
- when(workflowTaskRelationDao.queryUpstreamByCode(PROJECT_CODE,
TASK_CODE))
- .thenReturn(getProcessTaskRelationListV2());
- when(workflowDefinitionDao.queryByCode(PROCESS_DEFINITION_CODE))
- .thenReturn(Optional.of(getProcessDefinition()));
-
when(workflowTaskRelationDao.batchInsert(Mockito.anyList())).thenReturn(1);
-
when(workflowTaskRelationDao.updateById(Mockito.any())).thenReturn(true);
-
when(workflowTaskRelationLogDao.batchInsert(Mockito.anyList())).thenReturn(2);
- // success
- Long updatedTaskCode =
taskDefinitionService.updateTaskWithUpstream(user, PROJECT_CODE, TASK_CODE,
- taskDefinitionJson, UPSTREAM_CODE);
- assertEquals(TASK_CODE, updatedTaskCode);
-
Mockito.verify(taskDatasourcePermissionChecker).checkPermission(eq(user),
anyList());
- user.setUserType(UserType.GENERAL_USER);
- }
- }
-
- private String getTaskDefinitionJson() {
- return
"{\"name\":\"detail_up\",\"description\":\"\",\"taskType\":\"SHELL\",\"taskParams\":"
- +
"\"{\\\"resourceList\\\":[],\\\"localParams\\\":[{\\\"prop\\\":\\\"datetime\\\","
- + "\\\"direct\\\":\\\"IN\\\",\\\"type\\\":\\\"VARCHAR\\\","
- +
"\\\"value\\\":\\\"${system.datetime}\\\"}],\\\"rawScript\\\":\\\"echo
${datetime}\\\","
- +
"\\\"conditionResult\\\":\\\"{\\\\\\\"successNode\\\\\\\":[\\\\\\\"\\\\\\\"],"
- +
"\\\\\\\"failedNode\\\\\\\":[\\\\\\\"\\\\\\\"]}\\\",\\\"dependence\\\":{}}\","
- +
"\"flag\":0,\"taskPriority\":0,\"workerGroup\":\"default\",\"failRetryTimes\":0,"
- +
"\"failRetryInterval\":0,\"timeoutFlag\":0,\"timeoutNotifyStrategy\":0,\"timeout\":0,"
- + "\"delayTime\":0,\"resourceIds\":\"\"}";
- }
-
/**
* create admin user
*/
@@ -511,27 +449,6 @@ public class TaskDefinitionServiceImplTest {
return workflowTaskRelationList;
}
- private List<WorkflowTaskRelation> getProcessTaskRelationListV2() {
- List<WorkflowTaskRelation> workflowTaskRelationList = new
ArrayList<>();
-
- WorkflowTaskRelation workflowTaskRelation = new WorkflowTaskRelation();
- fillProcessTaskRelation(workflowTaskRelation);
-
- workflowTaskRelationList.add(workflowTaskRelation);
- workflowTaskRelation = new WorkflowTaskRelation();
- fillProcessTaskRelation(workflowTaskRelation);
- workflowTaskRelation.setPreTaskCode(4L);
- workflowTaskRelationList.add(workflowTaskRelation);
- return workflowTaskRelationList;
- }
-
- private void fillProcessTaskRelation(WorkflowTaskRelation
workflowTaskRelation) {
- workflowTaskRelation.setProjectCode(PROJECT_CODE);
-
workflowTaskRelation.setWorkflowDefinitionCode(PROCESS_DEFINITION_CODE);
- workflowTaskRelation.setPreTaskCode(TASK_CODE);
- workflowTaskRelation.setPostTaskCode(TASK_CODE + 1L);
- }
-
private List<WorkflowTaskRelationLog> getProcessTaskRelationLogList() {
List<WorkflowTaskRelationLog> processTaskRelationLogList = new
ArrayList<>();
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 a733aeaff9..917a9031fe 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
@@ -165,27 +165,10 @@ public interface WorkflowTaskRelationMapper extends
BaseMapper<WorkflowTaskRelat
IPage<WorkflowTaskRelation>
filterWorkflowTaskRelation(IPage<WorkflowTaskRelation> page,
@Param("relation")
WorkflowTaskRelation workflowTaskRelation);
- /**
- * batch update workflow task relation version
- *
- * @param workflowTaskRelation workflow task relation list
- * @return update num
- */
- 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);
- /**
- * workflow task relation by taskCode and postTaskVersion
- *
- * @param taskCode taskCode
- * @param postTaskVersion postTaskVersion
- * @return ProcessTaskRelation
- */
- List<WorkflowTaskRelation>
queryWorkflowTaskRelationByTaskCodeAndTaskVersion(@Param("taskCode") long
taskCode,
-
@Param("postTaskVersion") long postTaskVersion);
}
diff --git
a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/WorkflowTaskRelationDao.java
b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/WorkflowTaskRelationDao.java
index d84dd0c796..d0b161d839 100644
---
a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/WorkflowTaskRelationDao.java
+++
b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/WorkflowTaskRelationDao.java
@@ -36,12 +36,8 @@ public interface WorkflowTaskRelationDao extends
IDao<WorkflowTaskRelation> {
List<WorkflowTaskRelation> queryDownstreamByWorkflowDefinitionCode(long
workflowDefinitionCode);
- boolean updateWorkflowTaskRelationTaskVersion(WorkflowTaskRelation
workflowTaskRelation);
-
void deleteByWorkflowDefinitionCodeAndVersion(long workflowDefinitionCode,
int workflowDefinitionVersion);
- List<WorkflowTaskRelation>
queryWorkflowTaskRelationByTaskCodeAndTaskVersion(long taskCode, long
postTaskVersion);
-
List<WorkflowTaskRelation> queryByCode(long projectCode, long
workflowDefinitionCode, long preTaskCode,
long postTaskCode);
}
diff --git
a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/WorkflowTaskRelationDaoImpl.java
b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/WorkflowTaskRelationDaoImpl.java
index 79d84f84c1..c8a9281ca5 100644
---
a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/WorkflowTaskRelationDaoImpl.java
+++
b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/repository/impl/WorkflowTaskRelationDaoImpl.java
@@ -69,22 +69,11 @@ public class WorkflowTaskRelationDaoImpl extends
BaseDao<WorkflowTaskRelation, W
return
mybatisMapper.queryDownstreamByWorkflowDefinitionCode(workflowDefinitionCode);
}
- @Override
- public boolean updateWorkflowTaskRelationTaskVersion(WorkflowTaskRelation
workflowTaskRelation) {
- return
mybatisMapper.updateWorkflowTaskRelationTaskVersion(workflowTaskRelation) > 0;
- }
-
@Override
public void deleteByWorkflowDefinitionCodeAndVersion(long
workflowDefinitionCode, int workflowDefinitionVersion) {
mybatisMapper.deleteByWorkflowDefinitionCodeAndVersion(workflowDefinitionCode,
workflowDefinitionVersion);
}
- @Override
- public List<WorkflowTaskRelation>
queryWorkflowTaskRelationByTaskCodeAndTaskVersion(long taskCode,
-
long postTaskVersion) {
- return
mybatisMapper.queryWorkflowTaskRelationByTaskCodeAndTaskVersion(taskCode,
postTaskVersion);
- }
-
@Override
public List<WorkflowTaskRelation> queryByCode(long projectCode, long
workflowDefinitionCode, long preTaskCode,
long postTaskCode) {
diff --git
a/dolphinscheduler-dao/src/main/resources/org/apache/dolphinscheduler/dao/mapper/WorkflowTaskRelationMapper.xml
b/dolphinscheduler-dao/src/main/resources/org/apache/dolphinscheduler/dao/mapper/WorkflowTaskRelationMapper.xml
index 7fa82d3892..a46d221242 100644
---
a/dolphinscheduler-dao/src/main/resources/org/apache/dolphinscheduler/dao/mapper/WorkflowTaskRelationMapper.xml
+++
b/dolphinscheduler-dao/src/main/resources/org/apache/dolphinscheduler/dao/mapper/WorkflowTaskRelationMapper.xml
@@ -183,14 +183,6 @@
where r.workflow_definition_code = #{workflowCode}
and d.name = #{taskName}
</select>
- <update id="updateWorkflowTaskRelationTaskVersion">
- update t_ds_workflow_task_relation
- set pre_task_version=#{workflowTaskRelation.preTaskVersion},
- post_task_version=#{workflowTaskRelation.postTaskVersion},
-
workflow_definition_version=#{workflowTaskRelation.workflowDefinitionVersion}
- where id = #{workflowTaskRelation.id}
- </update>
-
<delete id="deleteByWorkflowDefinitionCodeAndVersion">
delete
from t_ds_workflow_task_relation
@@ -198,18 +190,4 @@
and workflow_definition_version = #{workflowDefinitionVersion}
</delete>
- <select id="queryWorkflowTaskRelationByTaskCodeAndTaskVersion"
resultType="org.apache.dolphinscheduler.dao.entity.WorkflowTaskRelation">
- select
- <include refid="baseSql"/>
- from t_ds_workflow_task_relation
- WHERE workflow_definition_code in (
- SELECT
- workflow_definition_code
- FROM
- t_ds_workflow_task_relation
- WHERE
- post_task_code = #{taskCode}
- and post_task_version =
#{postTaskVersion}
- )
- </select>
</mapper>