This is an automated email from the ASF dual-hosted git repository.
leonbao pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/dolphinscheduler.git
The following commit(s) were added to refs/heads/dev by this push:
new b6824b4 [DS-6737][MasterServer] fix event handle twice (#6738)
b6824b4 is described below
commit b6824b47419a5abc7ab4279d648417249a74f122
Author: wind <[email protected]>
AuthorDate: Mon Nov 8 16:54:58 2021 +0800
[DS-6737][MasterServer] fix event handle twice (#6738)
Co-authored-by: caishunfeng <[email protected]>
---
.../server/master/runner/WorkflowExecuteThread.java | 15 ++++++---------
1 file changed, 6 insertions(+), 9 deletions(-)
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 fa60030..6e287ad 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
@@ -105,10 +105,7 @@ public class WorkflowExecuteThread implements Runnable {
* runing TaskNode
*/
private final Map<Integer, ITaskProcessor> activeTaskProcessorMaps = new
ConcurrentHashMap<>();
- /**
- * task exec service
- */
- private final ExecutorService taskExecService;
+
/**
* process instance
*/
@@ -217,9 +214,6 @@ public class WorkflowExecuteThread implements Runnable {
this.processInstance = processInstance;
this.masterConfig = masterConfig;
- int masterTaskExecNum = masterConfig.getMasterExecTaskNum();
- this.taskExecService =
ThreadUtils.newDaemonFixedThreadExecutor("Master-Task-Exec-Thread",
- masterTaskExecNum);
this.nettyExecutorManager = nettyExecutorManager;
this.processAlertManager = processAlertManager;
this.taskTimeoutCheckList = taskTimeoutCheckList;
@@ -228,8 +222,11 @@ public class WorkflowExecuteThread implements Runnable {
@Override
public void run() {
try {
- startProcess();
- handleEvents();
+ if (!this.isStart()) {
+ startProcess();
+ } else {
+ handleEvents();
+ }
} catch (Exception e) {
logger.error("handler error:", e);
}