This is an automated email from the ASF dual-hosted git repository.

zhongjiajie pushed a commit to branch 2.0.8-prepare
in repository https://gitbox.apache.org/repos/asf/dolphinscheduler.git


The following commit(s) were added to refs/heads/2.0.8-prepare by this push:
     new feeda4aa32 [Fix-13244] [Master] Stopping a workflow does not update 
task status correctly (#13375)
feeda4aa32 is described below

commit feeda4aa3273e4c2badc16b4bd54d24c1747a624
Author: JinYong Li <[email protected]>
AuthorDate: Fri Jan 13 21:06:09 2023 +0800

    [Fix-13244] [Master] Stopping a workflow does not update task status 
correctly (#13375)
    
    Co-authored-by: JinyLeeChina <[email protected]>
---
 .../dolphinscheduler/common/utils/HadoopUtils.java | 11 ++--
 .../dao/mapper/TaskInstanceMapper.java             |  2 +
 .../dao/mapper/TaskInstanceMapper.xml              | 12 ++++
 .../processor/TaskKillResponseProcessor.java       |  6 +-
 .../master/processor/queue/TaskResponseEvent.java  | 13 ++++
 .../processor/queue/TaskResponsePersistThread.java |  8 ++-
 .../processor/queue/TaskResponseService.java       |  6 +-
 .../server/master/runner/EventExecuteService.java  |  2 +-
 .../master/runner/WorkflowExecuteThread.java       | 18 +++--
 .../worker/processor/TaskExecuteProcessor.java     |  2 -
 .../worker/processor/TaskKillAckProcessor.java     | 11 ++--
 .../server/worker/processor/TaskKillProcessor.java | 14 ++--
 .../server/worker/runner/TaskExecuteThread.java    | 76 ++++++++--------------
 .../service/process/ProcessService.java            | 15 ++++-
 .../spi/task/TaskExecutionContextCacheManager.java |  8 +++
 .../spi/task/request/TaskRequest.java              | 14 ++++
 16 files changed, 134 insertions(+), 84 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> {
     TaskInstance queryLastTaskInstance(@Param("taskCode") long taskCode, 
@Param("startTime") Date startTime, @Param("endTime") Date endTime);
 
     List<TaskInstance> queryLastTaskInstanceList(@Param("taskCodes") Set<Long> 
taskCodes, @Param("startTime") Date startTime, @Param("endTime") Date endTime);
+
+    List<TaskInstance> queryTaskInstanceListByIds(@Param("ids") Set<Integer> 
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 <![CDATA[ >= ]]> #{startTime} and start_time 
<![CDATA[ <= ]]> #{endTime}
         </if>
     </select>
+    <select id="queryTaskInstanceListByIds" 
resultType="org.apache.dolphinscheduler.dao.entity.TaskInstance">
+        select
+        <include refid="baseSql"/>
+        from t_ds_task_instance
+        where 1=1
+        <if test="ids != null and ids.size() != 0">
+            and id in
+            <foreach collection="ids" index="index" item="i" open="(" 
separator="," close=")">
+                #{i}
+            </foreach>
+        </if>
+    </select>
 </mapper>
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<String, Property> allProperty, 
Map<String, TaskInstance> 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<TaskInstance> 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()));
-
-        Pair<Boolean, List<String>> result = doKill(killCommand);
-
         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<Boolean, List<String>> result = doKill(killCommand);
+
         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<TaskInstance> findTaskInstanceListByIds(Set<Integer> 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<TaskRequest> 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;
     }

Reply via email to