From c1b434e19796e39c788cc3d42fc3fba3a70642d5 Mon Sep 17 00:00:00 2001 From: pegasas Date: Mon, 20 May 2024 13:47:30 +0800 Subject: [PATCH] [Improvement][Master] Add master explicit overload control --- .../server/master/runner/WorkflowExecuteThreadPool.java | 9 +++++++++ 1 file changed, 9 insertions(+) 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(); + } }