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);