From feeda4aa3273e4c2badc16b4bd54d24c1747a624 Mon Sep 17 00:00:00 2001 From: JinYong Li <42576980+JinyLeeChina@users.noreply.github.com> Date: Fri, 13 Jan 2023 21:06:09 +0800 Subject: [PATCH] [Fix-13244] [Master] Stopping a workflow does not update task status correctly (#13375) Co-authored-by: JinyLeeChina --- .../common/utils/HadoopUtils.java | 11 ++- .../dao/mapper/TaskInstanceMapper.java | 2 + .../dao/mapper/TaskInstanceMapper.xml | 12 +++ .../processor/TaskKillResponseProcessor.java | 6 +- .../processor/queue/TaskResponseEvent.java | 13 ++++ .../queue/TaskResponsePersistThread.java | 8 +- .../processor/queue/TaskResponseService.java | 6 +- .../master/runner/EventExecuteService.java | 2 +- .../master/runner/WorkflowExecuteThread.java | 18 +++-- .../processor/TaskExecuteProcessor.java | 2 - .../processor/TaskKillAckProcessor.java | 11 ++- .../worker/processor/TaskKillProcessor.java | 16 ++-- .../worker/runner/TaskExecuteThread.java | 76 +++++++------------ .../service/process/ProcessService.java | 15 +++- .../TaskExecutionContextCacheManager.java | 8 ++ .../spi/task/request/TaskRequest.java | 14 ++++ 16 files changed, 135 insertions(+), 85 deletions(-) diff --git a/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/utils/HadoopUtils.java b/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/utils/HadoopUtils.java index 3b6ac17019..508ca31985 100644 --- a/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/utils/HadoopUtils.java +++ b/dolphinscheduler-common/src/main/java/org/apache/dolphinscheduler/common/utils/HadoopUtils.java @@ -92,8 +92,9 @@ public class HadoopUtils implements Closeable { private FileSystem fs; private HadoopUtils() { - init(); - initHdfsPath(); + if(init()) { + initHdfsPath(); + } } public static HadoopUtils getInstance() { @@ -120,7 +121,7 @@ public class HadoopUtils implements Closeable { /** * init hadoop configuration */ - private void init() { + private boolean init() { try { configuration = new HdfsConfiguration(); @@ -168,11 +169,13 @@ public class HadoopUtils implements Closeable { configuration.set(Constants.FS_S3A_ACCESS_KEY, PropertyUtils.getString(Constants.FS_S3A_ACCESS_KEY)); configuration.set(Constants.FS_S3A_SECRET_KEY, PropertyUtils.getString(Constants.FS_S3A_SECRET_KEY)); fs = FileSystem.get(configuration); + } else { + return false; } - } catch (Exception e) { logger.error(e.getMessage(), e); } + return true; } /** 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 5291d6c020..b2a766ce49 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 @@ -84,4 +84,6 @@ public interface TaskInstanceMapper extends BaseMapper { TaskInstance queryLastTaskInstance(@Param("taskCode") long taskCode, @Param("startTime") Date startTime, @Param("endTime") Date endTime); List queryLastTaskInstanceList(@Param("taskCodes") Set taskCodes, @Param("startTime") Date startTime, @Param("endTime") Date endTime); + + List queryTaskInstanceListByIds(@Param("ids") Set ids); } 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 46dfde8543..5ddec3f78d 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 @@ -196,4 +196,16 @@ and start_time = ]]> #{startTime} and start_time #{endTime} + diff --git a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/processor/TaskKillResponseProcessor.java b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/processor/TaskKillResponseProcessor.java index 36dde2982c..24101108f0 100644 --- a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/processor/TaskKillResponseProcessor.java +++ b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/processor/TaskKillResponseProcessor.java @@ -71,10 +71,8 @@ public class TaskKillResponseProcessor implements NettyRequestProcessor { TaskKillResponseCommand responseCommand = JSONUtils.parseObject(command.getBody(), TaskKillResponseCommand.class); logger.info("received task kill response command : {}", responseCommand); // TaskResponseEvent - TaskResponseEvent taskResponseEvent = TaskResponseEvent.newActionStop(ExecutionStatus.of(responseCommand.getStatus()), - responseCommand.getTaskInstanceId(), - responseCommand.getProcessInstanceId() - ); + TaskResponseEvent taskResponseEvent = TaskResponseEvent.newKillResponse(ExecutionStatus.of(responseCommand.getStatus()), + responseCommand.getTaskInstanceId(), channel, responseCommand.getProcessInstanceId()); taskResponseService.addResponse(taskResponseEvent); } diff --git a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/processor/queue/TaskResponseEvent.java b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/processor/queue/TaskResponseEvent.java index ebf4a017f1..da6300fcbb 100644 --- a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/processor/queue/TaskResponseEvent.java +++ b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/processor/queue/TaskResponseEvent.java @@ -100,6 +100,19 @@ public class TaskResponseEvent { */ private long opaque; + public static TaskResponseEvent newKillResponse(ExecutionStatus state, + int taskInstanceId, + Channel channel, + int processInstanceId) { + TaskResponseEvent event = new TaskResponseEvent(); + event.setState(state); + event.setTaskInstanceId(taskInstanceId); + event.setEvent(Event.ACTION_STOP); + event.setChannel(channel); + event.setProcessInstanceId(processInstanceId); + return event; + } + public static TaskResponseEvent newActionStop(ExecutionStatus state, int taskInstanceId, int processInstanceId) { diff --git a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/processor/queue/TaskResponsePersistThread.java b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/processor/queue/TaskResponsePersistThread.java index f337757110..715f64a0c1 100644 --- a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/processor/queue/TaskResponsePersistThread.java +++ b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/processor/queue/TaskResponsePersistThread.java @@ -26,7 +26,6 @@ import org.apache.dolphinscheduler.remote.command.DBTaskAckCommand; import org.apache.dolphinscheduler.remote.command.DBTaskResponseCommand; import org.apache.dolphinscheduler.remote.command.TaskKillAckCommand; import org.apache.dolphinscheduler.remote.command.TaskRecallAckCommand; -import org.apache.dolphinscheduler.remote.processor.NettyRemoteChannel; import org.apache.dolphinscheduler.server.master.runner.WorkflowExecuteThread; import org.apache.dolphinscheduler.server.master.runner.task.ITaskProcessor; import org.apache.dolphinscheduler.server.master.runner.task.TaskAction; @@ -159,6 +158,10 @@ public class TaskResponsePersistThread implements Runnable { taskProcessor.persist(TaskAction.STOP); logger.debug("ACTION_STOP: task instance id:{}, process instance id:{}", taskResponseEvent.getTaskInstanceId(), taskResponseEvent.getProcessInstanceId()); } + workflowExecuteThread.getActiveTaskProcessorMaps().remove(taskResponseEvent.getTaskInstanceId()); + if (workflowExecuteThread.activeTaskFinish()) { + this.processInstanceMapper.remove(taskResponseEvent.getProcessInstanceId()); + } } if (channel != null) { @@ -197,7 +200,8 @@ public class TaskResponsePersistThread implements Runnable { } WorkflowExecuteThread workflowExecuteThread = this.processInstanceMapper.get(taskResponseEvent.getProcessInstanceId()); - if (workflowExecuteThread != null && taskResponseEvent.getState().typeIsFinished()) { + if (workflowExecuteThread != null && taskResponseEvent.getState().typeIsFinished() + && event != Event.ACTION_STOP && !workflowExecuteThread.getProcessInstance().getState().typeIsStop()) { StateEvent stateEvent = new StateEvent(); stateEvent.setProcessInstanceId(taskResponseEvent.getProcessInstanceId()); stateEvent.setTaskInstanceId(taskResponseEvent.getTaskInstanceId()); diff --git a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/processor/queue/TaskResponseService.java b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/processor/queue/TaskResponseService.java index b5e70eedc8..74c88ac588 100644 --- a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/processor/queue/TaskResponseService.java +++ b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/processor/queue/TaskResponseService.java @@ -177,7 +177,7 @@ public class TaskResponseService { Thread.currentThread().interrupt(); break; } catch (Exception e) { - logger.error("persist task error", e); + logger.error("handle task error", e); } } logger.info("StateEventResponseWorker stopped"); @@ -227,7 +227,7 @@ public class TaskResponseService { FutureCallback futureCallback = new FutureCallback() { @Override public void onSuccess(Object o) { - logger.info("persist events {} succeeded.", taskResponsePersistThread.getProcessInstanceId()); + logger.info("handle events {} succeeded.", taskResponsePersistThread.getProcessInstanceId()); if (!processInstanceMap.containsKey(taskResponsePersistThread.getProcessInstanceId())) { processTaskResponseMap.remove(taskResponsePersistThread.getProcessInstanceId()); logger.info("remove process instance: {}", taskResponsePersistThread.getProcessInstanceId()); @@ -237,7 +237,7 @@ public class TaskResponseService { @Override public void onFailure(Throwable throwable) { - logger.error("persist events failed: {}", throwable); + logger.error("handle events failed: {}", throwable.getMessage()); if (!processInstanceMap.containsKey(taskResponsePersistThread.getProcessInstanceId())) { processTaskResponseMap.remove(taskResponsePersistThread.getProcessInstanceId()); logger.info("remove process instance: {}", taskResponsePersistThread.getProcessInstanceId()); diff --git a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/runner/EventExecuteService.java b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/runner/EventExecuteService.java index ae17babf35..2279cad7b1 100644 --- a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/runner/EventExecuteService.java +++ b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/runner/EventExecuteService.java @@ -133,7 +133,7 @@ public class EventExecuteService extends Thread { FutureCallback futureCallback = new FutureCallback() { @Override public void onSuccess(Object o) { - if (workflowExecuteThread.workFlowFinish()) { + if (workflowExecuteThread.workFlowFinish() && workflowExecuteThread.activeTaskFinish()) { processInstanceExecMaps.remove(processInstanceId); notifyProcessChanged(); logger.info("process instance {} finished.", processInstanceId); diff --git a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/runner/WorkflowExecuteThread.java b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/runner/WorkflowExecuteThread.java index 8f96c20c7d..cc6078a1e6 100644 --- a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/runner/WorkflowExecuteThread.java +++ b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/runner/WorkflowExecuteThread.java @@ -942,11 +942,6 @@ public class WorkflowExecuteThread implements Runnable { taskInstance.setVarPool(JSONUtils.toJsonString(allProperty.values())); } } -// else { -// if (StringUtils.isNotEmpty(processInstance.getVarPool())) { -// taskInstance.setVarPool(processInstance.getVarPool()); -// } -// } } private void setVarPoolValue(Map allProperty, Map allTaskInstance, TaskInstance preTaskInstance, Property thisProperty) { @@ -1457,6 +1452,19 @@ public class WorkflowExecuteThread implements Runnable { return this.processInstance.getState().typeIsFinished(); } + public boolean activeTaskFinish() { + if (activeTaskProcessorMaps.isEmpty()) { + return true; + } + List taskInstanceList = processService.findTaskInstanceListByIds(activeTaskProcessorMaps.keySet()); + for (TaskInstance taskInstance : taskInstanceList) { + if (!taskInstance.getState().typeIsFinished()) { + return false; + } + } + return true; + } + /** * handling the list of tasks to be submitted */ diff --git a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/worker/processor/TaskExecuteProcessor.java b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/worker/processor/TaskExecuteProcessor.java index 0ac829293b..b587235883 100644 --- a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/worker/processor/TaskExecuteProcessor.java +++ b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/worker/processor/TaskExecuteProcessor.java @@ -88,8 +88,6 @@ public class TaskExecuteProcessor implements NettyRequestProcessor { * @param taskExecutionContext task */ private void setTaskCache(TaskExecutionContext taskExecutionContext) { - TaskExecutionContext preTaskCache = new TaskExecutionContext(); - preTaskCache.setTaskInstanceId(taskExecutionContext.getTaskInstanceId()); TaskRequest taskRequest = JSONUtils.parseObject(JSONUtils.toJsonString(taskExecutionContext), TaskRequest.class); TaskExecutionContextCacheManager.cacheTaskExecutionContext(taskRequest); } diff --git a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/worker/processor/TaskKillAckProcessor.java b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/worker/processor/TaskKillAckProcessor.java index c381af3c06..dff97191a2 100644 --- a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/worker/processor/TaskKillAckProcessor.java +++ b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/worker/processor/TaskKillAckProcessor.java @@ -41,20 +41,19 @@ public class TaskKillAckProcessor implements NettyRequestProcessor { public void process(Channel channel, Command command) { Preconditions.checkArgument(CommandType.TASK_KILL_RESPONSE_ACK == command.getType(), String.format("invalid command type : %s", command.getType())); - - TaskKillAckCommand taskKillAckCommand = JSONUtils.parseObject( - command.getBody(), TaskKillAckCommand.class); - + TaskKillAckCommand taskKillAckCommand = JSONUtils.parseObject(command.getBody(), TaskKillAckCommand.class); if (taskKillAckCommand == null) { + logger.warn("Cannot parse command, command type: {}", command.getType()); return; } + logger.info("received kill ack command : {}", taskKillAckCommand); if (taskKillAckCommand.getStatus() == ExecutionStatus.SUCCESS.getCode()) { ResponceCache.get().removeKillResponseCache(taskKillAckCommand.getTaskInstanceId()); TaskExecutionContextCacheManager.removeByTaskInstanceId(taskKillAckCommand.getTaskInstanceId()); - logger.debug("removeKillResponseCache: task instance id:{}", taskKillAckCommand.getTaskInstanceId()); + logger.info("removeKillResponseCache: task instance id:{}", taskKillAckCommand.getTaskInstanceId()); TaskCallbackService.remove(taskKillAckCommand.getTaskInstanceId()); - logger.debug("remove REMOTE_CHANNELS, task instance id:{}", taskKillAckCommand.getTaskInstanceId()); + logger.info("remove REMOTE_CHANNELS, task instance id:{}", taskKillAckCommand.getTaskInstanceId()); } } } diff --git a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/worker/processor/TaskKillProcessor.java b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/worker/processor/TaskKillProcessor.java index 8a7045d2ed..1caa23dc36 100644 --- a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/worker/processor/TaskKillProcessor.java +++ b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/worker/processor/TaskKillProcessor.java @@ -95,16 +95,18 @@ public class TaskKillProcessor implements NettyRequestProcessor { TaskKillRequestCommand killCommand = JSONUtils.parseObject(command.getBody(), TaskKillRequestCommand.class); logger.info("received kill command : {}", killCommand); - taskCallbackService.addRemoteChannel(killCommand.getTaskInstanceId(), - new NettyRemoteChannel(channel, command.getOpaque())); + TaskRequest taskRequest = TaskExecutionContextCacheManager.getByTaskInstanceId(killCommand.getTaskInstanceId()); + if (taskRequest == null) { + logger.warn("Cannot find taskInstanceId {} in taskContextCacheManager", killCommand.getTaskInstanceId()); + return; + } + taskRequest.setCurrentExecutionStatus(org.apache.dolphinscheduler.spi.task.ExecutionStatus.STOP); + TaskExecutionContextCacheManager.updateTaskExecutionContext(taskRequest); + + taskCallbackService.addRemoteChannel(killCommand.getTaskInstanceId(), new NettyRemoteChannel(channel, command.getOpaque())); Pair> result = doKill(killCommand); - TaskRequest taskRequest = TaskExecutionContextCacheManager.getByTaskInstanceId(killCommand.getTaskInstanceId()); - - if (taskRequest == null) { - return; - } TaskKillResponseCommand taskKillResponseCommand = buildKillTaskResponseCommand(taskRequest, result); ResponceCache.get().cache(taskKillResponseCommand.getTaskInstanceId(), taskKillResponseCommand.convert2Command(), Event.ACTION_STOP); taskCallbackService.sendResult(taskKillResponseCommand.getTaskInstanceId(), taskKillResponseCommand.convert2Command()); diff --git a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/worker/runner/TaskExecuteThread.java b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/worker/runner/TaskExecuteThread.java index 631131601a..90c5a41923 100644 --- a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/worker/runner/TaskExecuteThread.java +++ b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/worker/runner/TaskExecuteThread.java @@ -20,10 +20,14 @@ package org.apache.dolphinscheduler.server.worker.runner; import org.apache.dolphinscheduler.common.Constants; import org.apache.dolphinscheduler.common.enums.Event; import org.apache.dolphinscheduler.common.enums.ExecutionStatus; -import org.apache.dolphinscheduler.common.enums.TaskType; import org.apache.dolphinscheduler.common.process.Property; -import org.apache.dolphinscheduler.common.utils.*; -import org.apache.dolphinscheduler.remote.command.Command; +import org.apache.dolphinscheduler.common.utils.CommonUtils; +import org.apache.dolphinscheduler.common.utils.DateUtils; +import org.apache.dolphinscheduler.common.utils.FileUtils; +import org.apache.dolphinscheduler.common.utils.HadoopUtils; +import org.apache.dolphinscheduler.common.utils.JSONUtils; +import org.apache.dolphinscheduler.common.utils.LoggerUtils; +import org.apache.dolphinscheduler.common.utils.OSUtils; import org.apache.dolphinscheduler.remote.command.TaskExecuteAckCommand; import org.apache.dolphinscheduler.remote.command.TaskExecuteResponseCommand; import org.apache.dolphinscheduler.server.utils.LogUtils; @@ -51,15 +55,12 @@ import java.util.List; import java.util.Map; import java.util.Set; import java.util.concurrent.Delayed; -import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import com.github.rholder.retry.RetryException; - /** * task scheduler thread */ @@ -141,11 +142,11 @@ public class TaskExecuteThread implements Runnable, Delayed { taskExecutionContext.setStartTime(new Date()); } if (taskExecutionContext.getCurrentExecutionStatus() != ExecutionStatus.RUNNING_EXECUTION) { - changeTaskExecutionStatusToRunning(); + //changeTaskExecutionStatusToRunning(); + logger.info("the task begins to execute. task instance id: {}", taskExecutionContext.getTaskInstanceId()); + taskExecutionContext.setCurrentExecutionStatus(ExecutionStatus.RUNNING_EXECUTION); + sendTaskExecuteRunningCommand(taskExecutionContext); } - logger.info("the task begins to execute. task instance id: {}", taskExecutionContext.getTaskInstanceId()); - taskExecutionContext.setCurrentExecutionStatus(ExecutionStatus.RUNNING_EXECUTION); - sendTaskExecuteRunningCommand(taskExecutionContext); int dryRun = taskExecutionContext.getDryRun(); // copy hdfs/minio file to local if (dryRun == Constants.DRY_RUN_FLAG_NO) { @@ -168,13 +169,19 @@ public class TaskExecuteThread implements Runnable, Delayed { throw new RuntimeException(String.format("%s Task Plugin Not Found,Please Check Config File.", taskExecutionContext.getTaskType())); } TaskRequest taskRequest = JSONUtils.parseObject(JSONUtils.toJsonString(taskExecutionContext), TaskRequest.class); + if (null == taskRequest) { + throw new RuntimeException("The taskExecutionContext parse error"); + } String taskLogName = LoggerUtils.buildTaskId(LoggerUtils.TASK_LOGGER_INFO_PREFIX, taskExecutionContext.getProcessDefineCode(), taskExecutionContext.getProcessDefineVersion(), taskExecutionContext.getProcessInstanceId(), taskExecutionContext.getTaskInstanceId()); taskRequest.setTaskLogName(taskLogName); - + if (!TaskExecutionContextCacheManager.updateTaskExecutionContext(taskRequest)) { + TaskExecutionContextCacheManager.cacheTaskExecutionContext(taskRequest); + logger.info("taskRequest reCache successfully, taskInstanceId: {}", taskExecutionContext.getTaskInstanceId()); + } // set the name of the current thread Thread.currentThread().setName(String.format(TaskConstants.TASK_LOGGER_THREAD_NAME_FORMAT,taskLogName)); @@ -212,9 +219,14 @@ public class TaskExecuteThread implements Runnable, Delayed { responseCommand.setProcessId(task.getProcessId()); responseCommand.setAppIds(task.getAppIds()); } finally { - TaskExecutionContextCacheManager.removeByTaskInstanceId(taskExecutionContext.getTaskInstanceId()); - ResponceCache.get().cache(taskExecutionContext.getTaskInstanceId(), responseCommand.convert2Command(), Event.RESULT); - taskCallbackService.sendResult(taskExecutionContext.getTaskInstanceId(), responseCommand.convert2Command()); + if (TaskExecutionContextCacheManager.statusIsStop(taskExecutionContext.getTaskInstanceId())) { + logger.info("task has exited, taskInstanceId:{}, exitStatusCode:{}, task executionStatus:{}", + taskExecutionContext.getTaskInstanceId(), this.task.getExitStatusCode(), ExecutionStatus.STOP); + } else { + TaskExecutionContextCacheManager.removeByTaskInstanceId(taskExecutionContext.getTaskInstanceId()); + ResponceCache.get().cache(taskExecutionContext.getTaskInstanceId(), responseCommand.convert2Command(), Event.RESULT); + taskCallbackService.sendResult(taskExecutionContext.getTaskInstanceId(), responseCommand.convert2Command()); + } clearTaskExecPath(); } } @@ -352,42 +364,6 @@ public class TaskExecuteThread implements Runnable, Delayed { } } - /** - * send an ack to change the status of the task. - */ - private void changeTaskExecutionStatusToRunning() { - taskExecutionContext.setCurrentExecutionStatus(ExecutionStatus.RUNNING_EXECUTION); - Command ackCommand = buildAckCommand().convert2Command(); - try { - RetryerUtils.retryCall(() -> { - taskCallbackService.sendAck(taskExecutionContext.getTaskInstanceId(), ackCommand); - return Boolean.TRUE; - }); - } catch (ExecutionException | RetryException e) { - logger.error(e.getMessage(), e); - } - } - - /** - * build ack command. - * - * @return TaskExecuteAckCommand - */ - private TaskExecuteAckCommand buildAckCommand() { - TaskExecuteAckCommand ackCommand = new TaskExecuteAckCommand(); - ackCommand.setTaskInstanceId(taskExecutionContext.getTaskInstanceId()); - ackCommand.setStatus(taskExecutionContext.getCurrentExecutionStatus().getCode()); - ackCommand.setStartTime(taskExecutionContext.getStartTime()); - ackCommand.setLogPath(taskExecutionContext.getLogPath()); - ackCommand.setHost(taskExecutionContext.getHost()); - if (TaskType.SQL.getDesc().equalsIgnoreCase(taskExecutionContext.getTaskType()) || TaskType.PROCEDURE.getDesc().equalsIgnoreCase(taskExecutionContext.getTaskType())) { - ackCommand.setExecutePath(null); - } else { - ackCommand.setExecutePath(taskExecutionContext.getExecutePath()); - } - return ackCommand; - } - /** * get current TaskExecutionContext * diff --git a/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/process/ProcessService.java b/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/process/ProcessService.java index 1d017eb4db..3200fd185d 100644 --- a/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/process/ProcessService.java +++ b/dolphinscheduler-service/src/main/java/org/apache/dolphinscheduler/service/process/ProcessService.java @@ -1489,12 +1489,25 @@ public class ProcessService { * find task instance by id * * @param taskId task id - * @return task intance + * @return task instance */ public TaskInstance findTaskInstanceById(Integer taskId) { return taskInstanceMapper.selectById(taskId); } + /** + * find task instance list by ids + * + * @param taskIds task id list + * @return task instance list + */ + public List findTaskInstanceListByIds(Set taskIds) { + if (CollectionUtils.isEmpty(taskIds)) { + return new ArrayList<>(); + } + return taskInstanceMapper.queryTaskInstanceListByIds(taskIds); + } + /** * package task instance,associate processInstance and processDefine * diff --git a/dolphinscheduler-spi/src/main/java/org/apache/dolphinscheduler/spi/task/TaskExecutionContextCacheManager.java b/dolphinscheduler-spi/src/main/java/org/apache/dolphinscheduler/spi/task/TaskExecutionContextCacheManager.java index e2ab195a4b..aa7c16926f 100644 --- a/dolphinscheduler-spi/src/main/java/org/apache/dolphinscheduler/spi/task/TaskExecutionContextCacheManager.java +++ b/dolphinscheduler-spi/src/main/java/org/apache/dolphinscheduler/spi/task/TaskExecutionContextCacheManager.java @@ -71,4 +71,12 @@ public class TaskExecutionContextCacheManager { public static Collection getAllTaskRequestList() { return taskRequestContextCache.values(); } + + public static boolean statusIsStop(Integer taskInstanceId) { + TaskRequest taskRequest = taskRequestContextCache.get(taskInstanceId); + if (taskRequest == null) { + return true; + } + return taskRequest.getCurrentExecutionStatus().typeIsStop(); + } } diff --git a/dolphinscheduler-spi/src/main/java/org/apache/dolphinscheduler/spi/task/request/TaskRequest.java b/dolphinscheduler-spi/src/main/java/org/apache/dolphinscheduler/spi/task/request/TaskRequest.java index 3fa9442174..76cbcf8b08 100644 --- a/dolphinscheduler-spi/src/main/java/org/apache/dolphinscheduler/spi/task/request/TaskRequest.java +++ b/dolphinscheduler-spi/src/main/java/org/apache/dolphinscheduler/spi/task/request/TaskRequest.java @@ -18,6 +18,7 @@ package org.apache.dolphinscheduler.spi.task.request; import org.apache.dolphinscheduler.spi.enums.TaskTimeoutStrategy; +import org.apache.dolphinscheduler.spi.task.ExecutionStatus; import org.apache.dolphinscheduler.spi.task.Property; import java.util.Date; @@ -183,6 +184,11 @@ public class TaskRequest { */ private int delayTime; + /** + * current execution status + */ + private ExecutionStatus currentExecutionStatus; + /** * Task Logger name should be like: Task-{processDefinitionId}-{processInstanceId}-{taskInstanceId} */ @@ -473,6 +479,14 @@ public class TaskRequest { this.delayTime = delayTime; } + public ExecutionStatus getCurrentExecutionStatus() { + return currentExecutionStatus; + } + + public void setCurrentExecutionStatus(ExecutionStatus currentExecutionStatus) { + this.currentExecutionStatus = currentExecutionStatus; + } + public SQLTaskExecutionContext getSqlTaskExecutionContext() { return sqlTaskExecutionContext; }