From 5592e7bb7b770c8d1503630c044cbb3f0b10b340 Mon Sep 17 00:00:00 2001 From: wind Date: Fri, 7 Jan 2022 10:42:01 +0800 Subject: [PATCH] [Bug-7788][MasterServer] fix submit duplicate tasks sometimes when retry (#7808) (#7866) * [Bug-7788] fix submit duplicate tasks sometimes when retry * add exist check when add task to standby list * update * put queue contain judge first Co-authored-by: caishunfeng <534328519@qq.com> Co-authored-by: caishunfeng <534328519@qq.com> --- .../master/runner/WorkflowExecuteThread.java | 31 ++++++++++++++++--- 1 file changed, 27 insertions(+), 4 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 085c1d569d..92862ae619 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 @@ -76,6 +76,7 @@ import java.util.HashMap; import java.util.Iterator; import java.util.List; import java.util.Map; +import java.util.Map.Entry; import java.util.Objects; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; @@ -1171,13 +1172,35 @@ public class WorkflowExecuteThread implements Runnable { * @param taskInstance task instance */ private void addTaskToStandByList(TaskInstance taskInstance) { - logger.info("add task to stand by list, task name: {} , task id:{}", taskInstance.getName(), taskInstance.getId()); try { - if (!readyToSubmitTaskQueue.contains(taskInstance)) { - readyToSubmitTaskQueue.put(taskInstance); + if (readyToSubmitTaskQueue.contains(taskInstance)) { + logger.warn("task was found in ready submit queue, task code:{}", taskInstance.getTaskCode()); + return; } + // need to check if the tasks with same task code is active + boolean active = false; + Map taskInstanceMap = taskInstanceHashMap.column(taskInstance.getTaskCode()); + if (taskInstanceMap != null && taskInstanceMap.size() > 0) { + for (Entry entry : taskInstanceMap.entrySet()) { + Integer taskInstanceId = entry.getKey(); + if (activeTaskProcessorMaps.containsKey(taskInstanceId)) { + TaskInstance latestTaskInstance = processService.findTaskInstanceById(taskInstanceId); + if (latestTaskInstance != null && !latestTaskInstance.getState().typeIsFailure()) { + active = true; + break; + } + } + } + } + if (active) { + logger.warn("task was found in active task list, task code:{}", taskInstance.getTaskCode()); + return; + } + logger.info("add task to stand by list, task name:{}, task id:{}, task code:{}", + taskInstance.getName(), taskInstance.getId(), taskInstance.getTaskCode()); + readyToSubmitTaskQueue.put(taskInstance); } catch (Exception e) { - logger.error("add task instance to readyToSubmitTaskQueue, taskName: {}, task id: {}", taskInstance.getName(), taskInstance.getId(), e); + logger.error("add task instance to readyToSubmitTaskQueue, taskName:{}, task id:{}", taskInstance.getName(), taskInstance.getId(), e); } }