diff --git a/docs/docs/en/guide/upgrade/incompatible.md b/docs/docs/en/guide/upgrade/incompatible.md index ad118a1d6088..cadcc249dcc2 100644 --- a/docs/docs/en/guide/upgrade/incompatible.md +++ b/docs/docs/en/guide/upgrade/incompatible.md @@ -55,4 +55,5 @@ This document records the incompatible updates between each version. You need to * **Removed transient fields**: `stateDescList`, `workflowDefinition`, `dagData`, `queue`, `locations`, `dependenceScheduleTimes` * **Removed derived properties**: `cmdTypeIfComplement`, `complementData` (related to complement-data executions; use the detail API to obtain them) * To obtain any of these fields, use the detail API `GET /projects/{projectCode}/workflow-instances/{id}` instead, which continues to return the full `WorkflowInstance` object. ([#18444](https://github.com/apache/dolphinscheduler/pull/18444)) +* The task instance list APIs (`GET /projects/{projectCode}/task-instances`, `GET /projects/{projectCode}/task-instances/visual`) no longer return the heavy fields `taskParams`, `varPool`, and `logPath` in the response body. To obtain these fields, use the task detail obtained from the workflow instance detail API, which continues to return the full `TaskInstance` object. ([#18595](https://github.com/apache/dolphinscheduler/pull/18595)) diff --git a/docs/docs/zh/guide/upgrade/incompatible.md b/docs/docs/zh/guide/upgrade/incompatible.md index a0f945501989..b7dbea6e30bd 100644 --- a/docs/docs/zh/guide/upgrade/incompatible.md +++ b/docs/docs/zh/guide/upgrade/incompatible.md @@ -55,4 +55,5 @@ * **移除的非数据库字段**:`stateDescList`、`workflowDefinition`、`dagData`、`queue`、`locations`、`dependenceScheduleTimes` * **移除的派生属性**:`cmdTypeIfComplement`、`complementData`(补数执行相关,如需获取请使用详情接口) * 如需获取这些字段,请使用详情接口 `GET /projects/{projectCode}/workflow-instances/{id}`,该接口仍返回完整的 `WorkflowInstance` 对象 ([#18444](https://github.com/apache/dolphinscheduler/pull/18444)) +* 任务实例列表接口(`GET /projects/{projectCode}/task-instances`、`GET /projects/{projectCode}/task-instances/visual`)的响应体不再返回大字段 `taskParams`、`varPool` 和 `logPath`。如需获取这些字段,请通过工作流实例详情接口获取任务详情,该接口仍返回完整的 `TaskInstance` 对象 ([#18595](https://github.com/apache/dolphinscheduler/pull/18595)) diff --git a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/TaskInstanceController.java b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/TaskInstanceController.java index 915b848187f7..276dd934cfb0 100644 --- a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/TaskInstanceController.java +++ b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/controller/TaskInstanceController.java @@ -28,16 +28,13 @@ import org.apache.dolphinscheduler.api.service.TaskInstanceService; import org.apache.dolphinscheduler.api.utils.PageInfo; import org.apache.dolphinscheduler.api.utils.Result; -import org.apache.dolphinscheduler.api.utils.SensitivePropertyUtils; +import org.apache.dolphinscheduler.api.vo.TaskInstanceSummaryVO; import org.apache.dolphinscheduler.common.constants.Constants; import org.apache.dolphinscheduler.common.enums.TaskExecuteType; -import org.apache.dolphinscheduler.dao.entity.TaskInstance; import org.apache.dolphinscheduler.dao.entity.User; import org.apache.dolphinscheduler.plugin.task.api.enums.TaskExecutionStatus; import org.apache.dolphinscheduler.plugin.task.api.utils.ParameterUtils; -import java.util.stream.Collectors; - import org.springframework.beans.factory.annotation.Autowired; import org.springframework.http.HttpStatus; import org.springframework.web.bind.annotation.GetMapping; @@ -99,25 +96,25 @@ public class TaskInstanceController extends BaseController { @GetMapping() @ResponseStatus(HttpStatus.OK) @ApiException(QUERY_TASK_LIST_PAGING_ERROR) - public Result> queryTaskListPaging(@Parameter(hidden = true) @RequestAttribute(value = Constants.SESSION_USER) User loginUser, - @Parameter(name = "projectCode", description = "PROJECT_CODE", required = true) @PathVariable long projectCode, - @RequestParam(value = "workflowInstanceId", required = false, defaultValue = "0") Integer workflowInstanceId, - @RequestParam(value = "workflowInstanceName", required = false) String workflowInstanceName, - @RequestParam(value = "workflowDefinitionName", required = false) String workflowDefinitionName, - @RequestParam(value = "searchVal", required = false) String searchVal, - @RequestParam(value = "taskName", required = false) String taskName, - @RequestParam(value = "taskCode", required = false) Long taskCode, - @RequestParam(value = "executorName", required = false) String executorName, - @RequestParam(value = "stateType", required = false) TaskExecutionStatus stateType, - @RequestParam(value = "host", required = false) String host, - @RequestParam(value = "startDate", required = false) String startTime, - @RequestParam(value = "endDate", required = false) String endTime, - @RequestParam(value = "taskExecuteType", required = false, defaultValue = "BATCH") TaskExecuteType taskExecuteType, - @RequestParam("pageNo") Integer pageNo, - @RequestParam("pageSize") Integer pageSize) { + public Result> queryTaskListPaging(@Parameter(hidden = true) @RequestAttribute(value = Constants.SESSION_USER) User loginUser, + @Parameter(name = "projectCode", description = "PROJECT_CODE", required = true) @PathVariable long projectCode, + @RequestParam(value = "workflowInstanceId", required = false, defaultValue = "0") Integer workflowInstanceId, + @RequestParam(value = "workflowInstanceName", required = false) String workflowInstanceName, + @RequestParam(value = "workflowDefinitionName", required = false) String workflowDefinitionName, + @RequestParam(value = "searchVal", required = false) String searchVal, + @RequestParam(value = "taskName", required = false) String taskName, + @RequestParam(value = "taskCode", required = false) Long taskCode, + @RequestParam(value = "executorName", required = false) String executorName, + @RequestParam(value = "stateType", required = false) TaskExecutionStatus stateType, + @RequestParam(value = "host", required = false) String host, + @RequestParam(value = "startDate", required = false) String startTime, + @RequestParam(value = "endDate", required = false) String endTime, + @RequestParam(value = "taskExecuteType", required = false, defaultValue = "BATCH") TaskExecuteType taskExecuteType, + @RequestParam("pageNo") Integer pageNo, + @RequestParam("pageSize") Integer pageSize) { checkPageParams(pageNo, pageSize); searchVal = ParameterUtils.handleEscapes(searchVal); - Result> result = taskInstanceService.queryTaskListPaging( + return taskInstanceService.queryTaskListPaging( loginUser, projectCode, workflowInstanceId, @@ -134,13 +131,6 @@ public Result> queryTaskListPaging(@Parameter(hidden = tr taskExecuteType, pageNo, pageSize); - PageInfo pageInfo = result.getData(); - if (pageInfo != null && pageInfo.getTotalList() != null) { - pageInfo.setTotalList(pageInfo.getTotalList().stream() - .map(SensitivePropertyUtils::mask) - .collect(Collectors.toList())); - } - return result; } /** diff --git a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/TaskInstanceService.java b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/TaskInstanceService.java index 50c735c3300c..2e9d60799f6b 100644 --- a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/TaskInstanceService.java +++ b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/TaskInstanceService.java @@ -19,8 +19,8 @@ import org.apache.dolphinscheduler.api.utils.PageInfo; import org.apache.dolphinscheduler.api.utils.Result; +import org.apache.dolphinscheduler.api.vo.TaskInstanceSummaryVO; import org.apache.dolphinscheduler.common.enums.TaskExecuteType; -import org.apache.dolphinscheduler.dao.entity.TaskInstance; import org.apache.dolphinscheduler.dao.entity.User; import org.apache.dolphinscheduler.plugin.task.api.enums.TaskExecutionStatus; @@ -44,22 +44,22 @@ public interface TaskInstanceService { * @param pageSize page size * @return task list page */ - Result> queryTaskListPaging(User loginUser, - long projectCode, - Integer workflowInstanceId, - String workflowInstanceName, - String workflowDefinitionName, - String taskName, - Long taskCode, - String executorName, - String startDate, - String endDate, - String searchVal, - TaskExecutionStatus stateType, - String host, - TaskExecuteType taskExecuteType, - Integer pageNo, - Integer pageSize); + Result> queryTaskListPaging(User loginUser, + long projectCode, + Integer workflowInstanceId, + String workflowInstanceName, + String workflowDefinitionName, + String taskName, + Long taskCode, + String executorName, + String startDate, + String endDate, + String searchVal, + TaskExecutionStatus stateType, + String host, + TaskExecuteType taskExecuteType, + Integer pageNo, + Integer pageSize); /** * change one task instance's state from failure to forced success diff --git a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/TaskInstanceServiceImpl.java b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/TaskInstanceServiceImpl.java index d060661980cb..caac8d2a5ec9 100644 --- a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/TaskInstanceServiceImpl.java +++ b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/TaskInstanceServiceImpl.java @@ -27,12 +27,14 @@ import org.apache.dolphinscheduler.api.service.UsersService; import org.apache.dolphinscheduler.api.utils.PageInfo; import org.apache.dolphinscheduler.api.utils.Result; +import org.apache.dolphinscheduler.api.vo.TaskInstanceSummaryVO; import org.apache.dolphinscheduler.common.enums.TaskExecuteType; import org.apache.dolphinscheduler.common.utils.DateUtils; import org.apache.dolphinscheduler.dao.entity.Project; import org.apache.dolphinscheduler.dao.entity.TaskInstance; import org.apache.dolphinscheduler.dao.entity.User; import org.apache.dolphinscheduler.dao.entity.WorkflowInstance; +import org.apache.dolphinscheduler.dao.model.TaskInstanceSummaryDto; import org.apache.dolphinscheduler.dao.repository.ProjectDao; import org.apache.dolphinscheduler.dao.repository.TaskInstanceDao; import org.apache.dolphinscheduler.dao.repository.WorkflowInstanceDao; @@ -106,23 +108,23 @@ public class TaskInstanceServiceImpl extends BaseServiceImpl implements TaskInst * @return task list page */ @Override - public Result> queryTaskListPaging(User loginUser, - long projectCode, - Integer workflowInstanceId, - String workflowInstanceName, - String workflowDefinitionName, - String taskName, - Long taskCode, - String executorName, - String startDate, - String endDate, - String searchVal, - TaskExecutionStatus stateType, - String host, - TaskExecuteType taskExecuteType, - Integer pageNo, - Integer pageSize) { - Result> result = new Result<>(); + public Result> queryTaskListPaging(User loginUser, + long projectCode, + Integer workflowInstanceId, + String workflowInstanceName, + String workflowDefinitionName, + String taskName, + Long taskCode, + String executorName, + String startDate, + String endDate, + String searchVal, + TaskExecutionStatus stateType, + String host, + TaskExecuteType taskExecuteType, + Integer pageNo, + Integer pageSize) { + Result> result = new Result<>(); // check user access for project projectService.checkProjectAndAuthThrowException(loginUser, projectCode, TASK_INSTANCE); int[] statusArray = null; @@ -131,9 +133,9 @@ public Result> queryTaskListPaging(User loginUser, } Date start = checkAndParseDateParameters(startDate); Date end = checkAndParseDateParameters(endDate); - Page page = new Page<>(pageNo, pageSize); - PageInfo pageInfo = new PageInfo<>(pageNo, pageSize); - IPage taskInstanceIPage; + Page page = new Page<>(pageNo, pageSize); + PageInfo pageInfo = new PageInfo<>(pageNo, pageSize); + IPage taskInstanceIPage; if (taskExecuteType == TaskExecuteType.STREAM) { // stream task without workflow instance taskInstanceIPage = taskInstanceDao.queryStreamTaskInstanceListPaging( @@ -165,12 +167,13 @@ public Result> queryTaskListPaging(User loginUser, start, end); } - List taskInstanceList = taskInstanceIPage.getRecords(); + List taskInstanceList = taskInstanceIPage.getRecords(); List executorIds = - taskInstanceList.stream().map(TaskInstance::getExecutorId).distinct().collect(Collectors.toList()); + taskInstanceList.stream().map(TaskInstanceSummaryDto::getExecutorId).distinct() + .collect(Collectors.toList()); List users = usersService.queryUser(executorIds); Map userMap = users.stream().collect(Collectors.toMap(User::getId, v -> v)); - for (TaskInstance taskInstance : taskInstanceList) { + for (TaskInstanceSummaryDto taskInstance : taskInstanceList) { taskInstance.setDuration(DateUtils.format2Duration(taskInstance.getStartTime(), taskInstance.getEndTime())); User user = userMap.get(taskInstance.getExecutorId()); if (user != null) { @@ -178,7 +181,9 @@ public Result> queryTaskListPaging(User loginUser, } } pageInfo.setTotal((int) taskInstanceIPage.getTotal()); - pageInfo.setTotalList(taskInstanceList); + pageInfo.setTotalList(taskInstanceList.stream() + .map(TaskInstanceSummaryVO::fromSummaryDto) + .collect(Collectors.toList())); result.setData(pageInfo); putMsg(result, Status.SUCCESS); return result; diff --git a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/WorkflowInstanceServiceImpl.java b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/WorkflowInstanceServiceImpl.java index bf4356d1f46e..30bdc234c2cf 100644 --- a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/WorkflowInstanceServiceImpl.java +++ b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/service/impl/WorkflowInstanceServiceImpl.java @@ -66,6 +66,7 @@ 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.model.TaskInstanceSummaryDto; import org.apache.dolphinscheduler.dao.model.WorkflowInstanceSummaryDto; import org.apache.dolphinscheduler.dao.repository.ProjectDao; import org.apache.dolphinscheduler.dao.repository.TaskDefinitionDao; @@ -707,11 +708,11 @@ public GanttDto viewGantt(User loginUser, long projectCode, Integer workflowInst List taskList = new ArrayList<>(); if (CollectionUtils.isNotEmpty(nodeList)) { - List taskInstances = taskInstanceDao.queryByWorkflowInstanceIdsAndTaskCodes( + List taskInstances = taskInstanceDao.queryByWorkflowInstanceIdsAndTaskCodes( Collections.singletonList(workflowInstanceId), nodeList); for (Long node : nodeList) { - TaskInstance taskInstance = null; - for (TaskInstance instance : taskInstances) { + TaskInstanceSummaryDto taskInstance = null; + for (TaskInstanceSummaryDto instance : taskInstances) { if (instance.getWorkflowInstanceId() == workflowInstanceId && instance.getTaskCode() == node) { taskInstance = instance; break; diff --git a/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/vo/TaskInstanceSummaryVO.java b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/vo/TaskInstanceSummaryVO.java new file mode 100644 index 000000000000..09cae80cbf2c --- /dev/null +++ b/dolphinscheduler-api/src/main/java/org/apache/dolphinscheduler/api/vo/TaskInstanceSummaryVO.java @@ -0,0 +1,217 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.dolphinscheduler.api.vo; + +import org.apache.dolphinscheduler.common.enums.Flag; +import org.apache.dolphinscheduler.common.enums.Priority; +import org.apache.dolphinscheduler.common.enums.TaskExecuteType; +import org.apache.dolphinscheduler.dao.entity.TaskInstance; +import org.apache.dolphinscheduler.dao.model.TaskInstanceSummaryDto; +import org.apache.dolphinscheduler.plugin.task.api.enums.TaskExecutionStatus; + +import java.util.Date; + +import lombok.AllArgsConstructor; +import lombok.Data; +import lombok.NoArgsConstructor; +import io.swagger.v3.oas.annotations.media.Schema; + +/** + * Lightweight response VO for task instance list / paging queries. + * + *

Unlike {@link TaskInstance}, this VO intentionally omits heavy columns + * that are only needed for task execution or detail views. This allows the + * corresponding DAO queries to use the optimized {@code listSql} projection + * instead of the full {@code baseSql}. + * + *

Incompatible API change (documented): The following properties that + * were previously present in task instance list API responses are no longer + * returned: + * + *

Heavy DB-backed fields removed from the SQL projection: + *

    + *
  • {@code taskParams}
  • + *
  • {@code varPool}
  • + *
  • {@code logPath}
  • + *
+ * + *

Transient (non-DB) fields removed from the entity that were always + * {@code null} in list responses: + *

    + *
  • {@code processDefinitionName}
  • + *
  • {@code taskGroupPriority}
  • + *
  • {@code workflowInstance}
  • + *
  • {@code workflowDefinition}
  • + *
  • {@code taskDefine}
  • + *
  • {@code workflowInstancePriority}
  • + *
+ * + *

Consumers that require {@code taskParams}, {@code varPool}, or + * {@code logPath} should use the task instance detail query path (e.g. + * {@code queryById}) instead, which continues to return the full + * {@link TaskInstance}. + */ +@Data +@NoArgsConstructor +@AllArgsConstructor +@Schema(name = "TASK_INSTANCE_QUERY_RESPONSE") +public class TaskInstanceSummaryVO { + + @Schema(description = "task instance id") + private Integer id; + + @Schema(description = "task instance name") + private String name; + + @Schema(description = "task type") + private String taskType; + + @Schema(description = "workflow instance id") + private int workflowInstanceId; + + @Schema(description = "workflow instance name") + private String workflowInstanceName; + + @Schema(description = "project code") + private Long projectCode; + + @Schema(description = "task code") + private long taskCode; + + @Schema(description = "task definition version") + private int taskDefinitionVersion; + + @Schema(description = "task execution status") + private TaskExecutionStatus state; + + @Schema(description = "first submit time") + private Date firstSubmitTime; + + @Schema(description = "submit time") + private Date submitTime; + + @Schema(description = "start time") + private Date startTime; + + @Schema(description = "end time") + private Date endTime; + + @Schema(description = "host") + private String host; + + @Schema(description = "execute path") + private String executePath; + + @Schema(description = "alert flag") + private Flag alertFlag; + + @Schema(description = "retry times") + private int retryTimes; + + @Schema(description = "pid") + private int pid; + + @Schema(description = "app link") + private String appLink; + + @Schema(description = "flag") + private Flag flag; + + @Schema(description = "max retry times") + private int maxRetryTimes; + + @Schema(description = "retry interval") + private int retryInterval; + + @Schema(description = "task instance priority") + private Priority taskInstancePriority; + + @Schema(description = "worker group") + private String workerGroup; + + @Schema(description = "environment code") + private Long environmentCode; + + @Schema(description = "executor id") + private int executorId; + + @Schema(description = "executor name") + private String executorName; + + @Schema(description = "delay time") + private int delayTime; + + @Schema(description = "dry run") + private int dryRun; + + @Schema(description = "task group id") + private int taskGroupId; + + @Schema(description = "cpu quota") + private Integer cpuQuota; + + @Schema(description = "memory max") + private Integer memoryMax; + + @Schema(description = "task execute type") + private TaskExecuteType taskExecuteType; + + @Schema(description = "duration string, e.g. 1h 2m 3s") + private String duration; + + /** + * Create a {@link TaskInstanceSummaryVO} from a {@link TaskInstanceSummaryDto} DAO DTO. + */ + public static TaskInstanceSummaryVO fromSummaryDto(TaskInstanceSummaryDto dto) { + return new TaskInstanceSummaryVO( + dto.getId(), + dto.getName(), + dto.getTaskType(), + dto.getWorkflowInstanceId(), + dto.getWorkflowInstanceName(), + dto.getProjectCode(), + dto.getTaskCode(), + dto.getTaskDefinitionVersion(), + dto.getState(), + dto.getFirstSubmitTime(), + dto.getSubmitTime(), + dto.getStartTime(), + dto.getEndTime(), + dto.getHost(), + dto.getExecutePath(), + dto.getAlertFlag(), + dto.getRetryTimes(), + dto.getPid(), + dto.getAppLink(), + dto.getFlag(), + dto.getMaxRetryTimes(), + dto.getRetryInterval(), + dto.getTaskInstancePriority(), + dto.getWorkerGroup(), + dto.getEnvironmentCode(), + dto.getExecutorId(), + dto.getExecutorName(), + dto.getDelayTime(), + dto.getDryRun(), + dto.getTaskGroupId(), + dto.getCpuQuota(), + dto.getMemoryMax(), + dto.getTaskExecuteType(), + dto.getDuration()); + } +} diff --git a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/controller/TaskInstanceControllerTest.java b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/controller/TaskInstanceControllerTest.java index 53e618bbe791..e9475d6f98cb 100644 --- a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/controller/TaskInstanceControllerTest.java +++ b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/controller/TaskInstanceControllerTest.java @@ -30,15 +30,12 @@ import org.apache.dolphinscheduler.api.service.TaskInstanceService; import org.apache.dolphinscheduler.api.utils.PageInfo; import org.apache.dolphinscheduler.api.utils.Result; +import org.apache.dolphinscheduler.api.vo.TaskInstanceSummaryVO; import org.apache.dolphinscheduler.common.enums.TaskExecuteType; import org.apache.dolphinscheduler.common.utils.JSONUtils; -import org.apache.dolphinscheduler.dao.entity.TaskInstance; import org.apache.dolphinscheduler.dao.entity.User; -import org.apache.dolphinscheduler.plugin.task.api.TaskConstants; import org.apache.dolphinscheduler.plugin.task.api.enums.TaskExecutionStatus; -import java.util.Collections; - import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Disabled; import org.junit.jupiter.api.Test; @@ -61,31 +58,23 @@ public class TaskInstanceControllerTest extends AbstractControllerTest { @Test public void testQueryTaskListPaging() { - Result> result = new Result<>(); Integer pageNo = 1; Integer pageSize = 20; - TaskInstance taskInstance = new TaskInstance(); - taskInstance.setTaskParams("{\"localParams\":[{\"prop\":\"token\",\"direct\":\"IN\",\"type\":\"VARCHAR\"," - + "\"value\":\"abc\",\"sensitive\":true}]}"); - PageInfo pageInfo = new PageInfo<>(pageNo, pageSize); - pageInfo.setTotalList(Collections.singletonList(taskInstance)); - result.setData(pageInfo); - result.setCode(Status.SUCCESS.getCode()); - result.setMsg(Status.SUCCESS.getMsg()); + PageInfo pageInfo = new PageInfo<>(pageNo, pageSize); + Result> mockResult = new Result<>(); + mockResult.setData(pageInfo); + mockResult.setCode(Status.SUCCESS.getCode()); + mockResult.setMsg(Status.SUCCESS.getMsg()); when(taskInstanceService.queryTaskListPaging(any(), eq(1L), eq(1), eq(""), eq(""), eq(""), any(), eq(""), any(), any(), eq(""), Mockito.any(), eq("192.168.xx.xx"), eq(TaskExecuteType.BATCH), any(), any())) - .thenReturn(result); - Result> taskResult = taskInstanceController.queryTaskListPaging(null, 1L, 1, "", "", "", + .thenReturn(mockResult); + Result> taskResult = taskInstanceController.queryTaskListPaging(null, 1L, 1, "", + "", "", "", 1L, "", TaskExecutionStatus.SUCCESS, "192.168.xx.xx", "2020-01-01 00:00:00", "2020-01-02 00:00:00", TaskExecuteType.BATCH, pageNo, pageSize); Assertions.assertEquals(Integer.valueOf(Status.SUCCESS.getCode()), taskResult.getCode()); - PageInfo maskedPage = taskResult.getData(); - Assertions.assertTrue( - maskedPage.getTotalList().get(0).getTaskParams().contains(TaskConstants.SENSITIVE_DATA_MASK)); - Assertions.assertFalse(maskedPage.getTotalList().get(0).getTaskParams().contains("abc")); - Assertions.assertTrue(taskInstance.getTaskParams().contains("abc")); } @Disabled diff --git a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/TaskInstanceServiceTest.java b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/TaskInstanceServiceTest.java index d7f5771121a7..86654661b2cf 100644 --- a/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/TaskInstanceServiceTest.java +++ b/dolphinscheduler-api/src/test/java/org/apache/dolphinscheduler/api/service/TaskInstanceServiceTest.java @@ -39,6 +39,7 @@ import org.apache.dolphinscheduler.dao.entity.TaskInstance; import org.apache.dolphinscheduler.dao.entity.User; import org.apache.dolphinscheduler.dao.entity.WorkflowInstance; +import org.apache.dolphinscheduler.dao.model.TaskInstanceSummaryDto; import org.apache.dolphinscheduler.dao.repository.ProjectDao; import org.apache.dolphinscheduler.dao.repository.TaskInstanceDao; import org.apache.dolphinscheduler.dao.repository.WorkflowInstanceDao; @@ -144,9 +145,15 @@ public void queryTaskListPaging() { Date end = DateUtils.stringToDate("2020-01-02 00:00:00"); WorkflowInstance workflowInstance = getProcessInstance(); TaskInstance taskInstance = getTaskInstance(); - List taskInstanceList = new ArrayList<>(); - Page pageReturn = new Page<>(1, 10); - taskInstanceList.add(taskInstance); + List taskInstanceList = new ArrayList<>(); + Page pageReturn = new Page<>(1, 10); + TaskInstanceSummaryDto summaryDto = new TaskInstanceSummaryDto(); + summaryDto.setId(taskInstance.getId()); + summaryDto.setName(taskInstance.getName()); + summaryDto.setExecutorId(taskInstance.getExecutorId()); + summaryDto.setStartTime(taskInstance.getStartTime()); + summaryDto.setEndTime(taskInstance.getEndTime()); + taskInstanceList.add(summaryDto); pageReturn.setRecords(taskInstanceList); doNothing().when(projectService).checkProjectAndAuthThrowException(loginUser, projectCode, TASK_INSTANCE); when(usersService.queryUser(loginUser.getId())).thenReturn(loginUser); 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 32a05a21ecc7..4da68ec647ba 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 @@ -22,6 +22,7 @@ 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.dao.model.TaskInstanceSummaryDto; import org.apache.dolphinscheduler.plugin.task.api.enums.TaskExecutionStatus; import org.apache.ibatis.annotations.Param; @@ -42,8 +43,8 @@ List findValidTaskListByWorkflowInstanceId(@Param("workflowInstanc TaskInstance queryByInstanceIdAndCode(@Param("workflowInstanceId") int workflowInstanceId, @Param("taskCode") Long taskCode); - List queryByWorkflowInstanceIdsAndTaskCodes(@Param("workflowInstanceIds") List workflowInstanceIds, - @Param("taskCodes") List taskCodes); + List queryByWorkflowInstanceIdsAndTaskCodes(@Param("workflowInstanceIds") List workflowInstanceIds, + @Param("taskCodes") List taskCodes); /** * Statistics task instance group by given project codes list by start time @@ -94,32 +95,32 @@ List countTaskInstanceStateByProjectCodesAndStatesBySubmitTi @Param("projectIds") Set projectIds, @Param("states") List states); - IPage queryTaskInstanceListPaging(IPage page, - @Param("projectCode") Long projectCode, - @Param("workflowInstanceId") Integer workflowInstanceId, - @Param("workflowInstanceName") String workflowInstanceName, - @Param("searchVal") String searchVal, - @Param("taskName") String taskName, - @Param("taskCode") Long taskCode, - @Param("executorName") String executorName, - @Param("states") int[] statusArray, - @Param("host") String host, - @Param("taskExecuteType") TaskExecuteType taskExecuteType, - @Param("startTime") Date startTime, - @Param("endTime") Date endTime); - - IPage queryStreamTaskInstanceListPaging(IPage page, - @Param("projectCode") Long projectCode, - @Param("workflowDefinitionName") String workflowDefinitionName, - @Param("searchVal") String searchVal, - @Param("taskName") String taskName, - @Param("taskCode") Long taskCode, - @Param("executorName") String executorName, - @Param("states") int[] statusArray, - @Param("host") String host, - @Param("taskExecuteType") TaskExecuteType taskExecuteType, - @Param("startTime") Date startTime, - @Param("endTime") Date endTime); + IPage queryTaskInstanceListPaging(IPage page, + @Param("projectCode") Long projectCode, + @Param("workflowInstanceId") Integer workflowInstanceId, + @Param("workflowInstanceName") String workflowInstanceName, + @Param("searchVal") String searchVal, + @Param("taskName") String taskName, + @Param("taskCode") Long taskCode, + @Param("executorName") String executorName, + @Param("states") int[] statusArray, + @Param("host") String host, + @Param("taskExecuteType") TaskExecuteType taskExecuteType, + @Param("startTime") Date startTime, + @Param("endTime") Date endTime); + + IPage queryStreamTaskInstanceListPaging(IPage page, + @Param("projectCode") Long projectCode, + @Param("workflowDefinitionName") String workflowDefinitionName, + @Param("searchVal") String searchVal, + @Param("taskName") String taskName, + @Param("taskCode") Long taskCode, + @Param("executorName") String executorName, + @Param("states") int[] statusArray, + @Param("host") String host, + @Param("taskExecuteType") TaskExecuteType taskExecuteType, + @Param("startTime") Date startTime, + @Param("endTime") Date endTime); void deleteByWorkflowInstanceId(@Param("workflowInstanceId") int workflowInstanceId); diff --git a/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/model/TaskInstanceSummaryDto.java b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/model/TaskInstanceSummaryDto.java new file mode 100644 index 000000000000..d261dfb1e68b --- /dev/null +++ b/dolphinscheduler-dao/src/main/java/org/apache/dolphinscheduler/dao/model/TaskInstanceSummaryDto.java @@ -0,0 +1,122 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.dolphinscheduler.dao.model; + +import org.apache.dolphinscheduler.common.enums.Flag; +import org.apache.dolphinscheduler.common.enums.Priority; +import org.apache.dolphinscheduler.common.enums.TaskExecuteType; +import org.apache.dolphinscheduler.plugin.task.api.enums.TaskExecutionStatus; + +import java.util.Date; + +import lombok.AllArgsConstructor; +import lombok.Data; +import lombok.NoArgsConstructor; + +/** + * Lightweight DTO for task instance list / paging queries. + * + *

Unlike {@link org.apache.dolphinscheduler.dao.entity.TaskInstance}, this DTO + * intentionally omits the heavy text columns {@code task_params}, + * {@code var_pool} and {@code log_path} that are only needed for task + * execution, detail views, or log retrieval. + * This allows the corresponding DAO queries to use the optimized + * {@code listSql} projection instead of the full {@code baseSql}. + * + *

Fields correspond 1:1 to the columns in the {@code listSql} / + * {@code listSqlV2} SQL projections in {@code TaskInstanceMapper.xml}. + * The {@code duration} field is not a DB column; it is populated by the + * application layer. + */ +@Data +@NoArgsConstructor +@AllArgsConstructor +public class TaskInstanceSummaryDto { + + private Integer id; + + private String name; + + private String taskType; + + private int workflowInstanceId; + + private String workflowInstanceName; + + private Long projectCode; + + private long taskCode; + + private int taskDefinitionVersion; + + private TaskExecutionStatus state; + + private Date firstSubmitTime; + + private Date submitTime; + + private Date startTime; + + private Date endTime; + + private String host; + + private String executePath; + + private Flag alertFlag; + + private int retryTimes; + + private int pid; + + private String appLink; + + private Flag flag; + + private int maxRetryTimes; + + private int retryInterval; + + private Priority taskInstancePriority; + + private String workerGroup; + + private Long environmentCode; + + private int executorId; + + private String executorName; + + private int delayTime; + + private int dryRun; + + private int taskGroupId; + + private Integer cpuQuota; + + private Integer memoryMax; + + private TaskExecuteType taskExecuteType; + + /** + * Task execution duration, e.g. "1h 2m 3s", populated by the application layer. + */ + private String duration; + +} 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 978dc0103b03..4792c27ac131 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 @@ -21,6 +21,7 @@ import org.apache.dolphinscheduler.dao.entity.TaskInstance; import org.apache.dolphinscheduler.dao.entity.WorkflowInstance; import org.apache.dolphinscheduler.dao.model.TaskInstanceStatusCountDto; +import org.apache.dolphinscheduler.dao.model.TaskInstanceSummaryDto; import org.apache.dolphinscheduler.plugin.task.api.enums.TaskExecutionStatus; import java.util.Collection; @@ -112,33 +113,33 @@ List countTaskInstanceStateByProjectCodes(Date start Date endTime, Collection projectCodes); - List queryByWorkflowInstanceIdsAndTaskCodes(List workflowInstanceIds, - List taskCodes); - - IPage queryTaskInstanceListPaging(IPage page, - Long projectCode, - Integer workflowInstanceId, - String workflowInstanceName, - String searchVal, - String taskName, - Long taskCode, - String executorName, - int[] statusArray, - String host, - TaskExecuteType taskExecuteType, - Date startTime, - Date endTime); - - IPage queryStreamTaskInstanceListPaging(IPage page, - Long projectCode, - String workflowDefinitionName, - String searchVal, - String taskName, - Long taskCode, - String executorName, - int[] statusArray, - String host, - TaskExecuteType taskExecuteType, - Date startTime, - Date endTime); + List queryByWorkflowInstanceIdsAndTaskCodes(List workflowInstanceIds, + List taskCodes); + + IPage queryTaskInstanceListPaging(IPage page, + Long projectCode, + Integer workflowInstanceId, + String workflowInstanceName, + String searchVal, + String taskName, + Long taskCode, + String executorName, + int[] statusArray, + String host, + TaskExecuteType taskExecuteType, + Date startTime, + Date endTime); + + IPage queryStreamTaskInstanceListPaging(IPage page, + Long projectCode, + String workflowDefinitionName, + String searchVal, + String taskName, + Long taskCode, + String executorName, + int[] statusArray, + String host, + TaskExecuteType taskExecuteType, + Date startTime, + Date endTime); } 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 893ab06e5fdc..a30fb61418f7 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 @@ -26,6 +26,7 @@ import org.apache.dolphinscheduler.dao.mapper.TaskInstanceMapper; import org.apache.dolphinscheduler.dao.mapper.WorkflowInstanceMapper; import org.apache.dolphinscheduler.dao.model.TaskInstanceStatusCountDto; +import org.apache.dolphinscheduler.dao.model.TaskInstanceSummaryDto; import org.apache.dolphinscheduler.dao.repository.BaseDao; import org.apache.dolphinscheduler.dao.repository.TaskInstanceDao; import org.apache.dolphinscheduler.plugin.task.api.enums.TaskExecutionStatus; @@ -193,42 +194,42 @@ public List countTaskInstanceStateByProjectCodes(Dat } @Override - public List queryByWorkflowInstanceIdsAndTaskCodes(List workflowInstanceIds, - List taskCodes) { + public List queryByWorkflowInstanceIdsAndTaskCodes(List workflowInstanceIds, + List taskCodes) { return mybatisMapper.queryByWorkflowInstanceIdsAndTaskCodes(workflowInstanceIds, taskCodes); } @Override - public IPage queryTaskInstanceListPaging(IPage page, - Long projectCode, - Integer workflowInstanceId, - String workflowInstanceName, - String searchVal, - String taskName, - Long taskCode, - String executorName, - int[] statusArray, - String host, - TaskExecuteType taskExecuteType, - Date startTime, - Date endTime) { + public IPage queryTaskInstanceListPaging(IPage page, + Long projectCode, + Integer workflowInstanceId, + String workflowInstanceName, + String searchVal, + String taskName, + Long taskCode, + String executorName, + int[] statusArray, + String host, + TaskExecuteType taskExecuteType, + Date startTime, + Date endTime) { return mybatisMapper.queryTaskInstanceListPaging(page, projectCode, workflowInstanceId, workflowInstanceName, searchVal, taskName, taskCode, executorName, statusArray, host, taskExecuteType, startTime, endTime); } @Override - public IPage queryStreamTaskInstanceListPaging(IPage page, - Long projectCode, - String workflowDefinitionName, - String searchVal, - String taskName, - Long taskCode, - String executorName, - int[] statusArray, - String host, - TaskExecuteType taskExecuteType, - Date startTime, - Date endTime) { + public IPage queryStreamTaskInstanceListPaging(IPage page, + Long projectCode, + String workflowDefinitionName, + String searchVal, + String taskName, + Long taskCode, + String executorName, + int[] statusArray, + String host, + TaskExecuteType taskExecuteType, + Date startTime, + Date endTime) { return mybatisMapper.queryStreamTaskInstanceListPaging(page, projectCode, workflowDefinitionName, searchVal, taskName, taskCode, executorName, statusArray, host, taskExecuteType, startTime, endTime); } 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 de1883346aa0..be6cf033e8c3 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 @@ -32,14 +32,14 @@ id, name, task_type, workflow_instance_id, workflow_instance_name, project_code, task_code, task_definition_version, state, submit_time, - start_time, end_time, host, alert_flag, retry_times, pid, app_link, + start_time, end_time, host, execute_path, alert_flag, retry_times, pid, app_link, flag, retry_interval, max_retry_times, task_instance_priority, worker_group,environment_code , executor_id, executor_name, first_submit_time, delay_time, dry_run, task_group_id, cpu_quota, memory_max, task_execute_type ${alias}.id, ${alias}.name, ${alias}.task_type, ${alias}.workflow_instance_id, ${alias}.workflow_instance_name, ${alias}.project_code, ${alias}.task_code, ${alias}.task_definition_version, ${alias}.state, ${alias}.submit_time, - ${alias}.start_time, ${alias}.end_time, ${alias}.host, ${alias}.alert_flag, ${alias}.retry_times, ${alias}.pid, ${alias}.app_link, - ${alias}.flag, ${alias}.retry_interval, ${alias}.max_retry_times, ${alias}.task_instance_priority, ${alias}.worker_group, ${alias}.environment_code, ${alias}.executor_id, ${alias}.executor_name, + ${alias}.start_time, ${alias}.end_time, ${alias}.host, ${alias}.execute_path, ${alias}.alert_flag, ${alias}.retry_times, ${alias}.pid, ${alias}.app_link, + ${alias}.flag, ${alias}.retry_interval, ${alias}.max_retry_times, ${alias}.task_instance_priority, ${alias}.worker_group,${alias}.environment_code, ${alias}.executor_id, ${alias}.first_submit_time, ${alias}.delay_time, ${alias}.dry_run, ${alias}.task_group_id, ${alias}.cpu_quota, ${alias}.memory_max, ${alias}.task_execute_type select - + from t_ds_task_instance where workflow_instance_id = #{workflowInstanceId} and task_code = #{taskCode} and flag = 1 limit 1 - select from t_ds_task_instance @@ -163,7 +163,7 @@ - select from t_ds_task_instance @@ -206,7 +206,7 @@ order by submit_time desc - select from t_ds_task_instance @@ -248,7 +248,7 @@