diff --git a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/runner/WorkflowExecuteThreadPool.java b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/runner/WorkflowExecuteThreadPool.java index 1aa4232865..1ab455ff49 100644 --- a/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/runner/WorkflowExecuteThreadPool.java +++ b/dolphinscheduler-master/src/main/java/org/apache/dolphinscheduler/server/master/runner/WorkflowExecuteThreadPool.java @@ -85,6 +85,12 @@ public class WorkflowExecuteThreadPool extends ThreadPoolTaskExecutor { if (!workflowExecuteThread.isStart() || workflowExecuteThread.eventSize() == 0) { return; } + + if (isOverload()) { + log.warn("WorkflowExecuteThreadPool is overload, cannot submit new workflowExecuteThread"); + return; + } + IWorkflowExecuteContext workflowExecuteRunnableContext = workflowExecuteThread.getWorkflowExecuteContext(); Integer workflowInstanceId = workflowExecuteRunnableContext.getWorkflowInstance().getId(); @@ -131,4 +137,7 @@ public class WorkflowExecuteThreadPool extends ThreadPoolTaskExecutor { }); } + public boolean isOverload() { + return this.getThreadPoolExecutor().getQueue().size() >= masterConfig.getExecThreads(); + } }