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

lidongdai pushed a commit to branch dev-1.3.0
in repository https://gitbox.apache.org/repos/asf/incubator-dolphinscheduler.git


The following commit(s) were added to refs/heads/dev-1.3.0 by this push:
     new 009fd01  fix [BUG] TaskExecutionContextCacheManagerImpl Do not execute 
removeByTaskInstanceId #2745 (#2754)
009fd01 is described below

commit 009fd01047283ca11d60714b128c3b52d61fc705
Author: yaoyao <[email protected]>
AuthorDate: Sat May 30 23:23:01 2020 +0800

    fix [BUG] TaskExecutionContextCacheManagerImpl Do not execute 
removeByTaskInstanceId #2745 (#2754)
---
 .../server/worker/processor/TaskKillProcessor.java             |  1 +
 .../server/worker/runner/TaskExecuteThread.java                | 10 ++++++++++
 2 files changed, 11 insertions(+)

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 b6f5827..cf0c051 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
@@ -93,6 +93,7 @@ public class TaskKillProcessor implements 
NettyRequestProcessor {
 
         TaskKillResponseCommand taskKillResponseCommand = 
buildKillTaskResponseCommand(killCommand,result);
         
taskCallbackService.sendResult(taskKillResponseCommand.getTaskInstanceId(), 
taskKillResponseCommand.convert2Command());
+        
taskExecutionContextCacheManager.removeByTaskInstanceId(taskKillResponseCommand.getTaskInstanceId());
     }
 
     /**
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 d314c55..d2d783a 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
@@ -27,9 +27,12 @@ import org.apache.dolphinscheduler.common.thread.ThreadUtils;
 import org.apache.dolphinscheduler.common.utils.*;
 import org.apache.dolphinscheduler.remote.command.TaskExecuteResponseCommand;
 import org.apache.dolphinscheduler.server.entity.TaskExecutionContext;
+import 
org.apache.dolphinscheduler.server.worker.cache.TaskExecutionContextCacheManager;
+import 
org.apache.dolphinscheduler.server.worker.cache.impl.TaskExecutionContextCacheManagerImpl;
 import org.apache.dolphinscheduler.server.worker.processor.TaskCallbackService;
 import org.apache.dolphinscheduler.server.worker.task.AbstractTask;
 import org.apache.dolphinscheduler.server.worker.task.TaskManager;
+import org.apache.dolphinscheduler.service.bean.SpringApplicationContext;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -64,6 +67,11 @@ public class TaskExecuteThread implements Runnable {
     private TaskCallbackService taskCallbackService;
 
     /**
+     * taskExecutionContextCacheManager
+     */
+    private TaskExecutionContextCacheManager taskExecutionContextCacheManager;
+
+    /**
      *  constructor
      * @param taskExecutionContext taskExecutionContext
      * @param taskCallbackService taskCallbackService
@@ -71,6 +79,7 @@ public class TaskExecuteThread implements Runnable {
     public TaskExecuteThread(TaskExecutionContext taskExecutionContext, 
TaskCallbackService taskCallbackService){
         this.taskExecutionContext = taskExecutionContext;
         this.taskCallbackService = taskCallbackService;
+        this.taskExecutionContextCacheManager = 
SpringApplicationContext.getBean(TaskExecutionContextCacheManagerImpl.class);
     }
 
     @Override
@@ -134,6 +143,7 @@ public class TaskExecuteThread implements Runnable {
             responseCommand.setAppIds(task.getAppIds());
         } finally {
             try {
+                
taskExecutionContextCacheManager.removeByTaskInstanceId(taskExecutionContext.getTaskInstanceId());
                 
taskCallbackService.sendResult(taskExecutionContext.getTaskInstanceId(), 
responseCommand.convert2Command());
             }catch (Exception e){
                 ThreadUtils.sleep(Constants.SLEEP_TIME_MILLIS);

Reply via email to