This is an automated email from the ASF dual-hosted git repository.

gavinchou pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/master by this push:
     new 94ad62a4aea [fix](cloud) prevent warm-up job scheduling starvation 
(#67527)
94ad62a4aea is described below

commit 94ad62a4aeaefbb3f41336f50fa0ecf8b4cdf203
Author: bobhan1 <[email protected]>
AuthorDate: Tue Sep 8 15:20:10 2026 +0800

    [fix](cloud) prevent warm-up job scheduling starvation (#67527)
    
    Problem Summary:
    
    When cloud warm-up jobs outnumber the configured active slots, the
    scheduler can over-submit tasks because a job is recorded as active only
    after its worker starts. The direct-handoff thread pool previously
    discarded rejected tasks, and unordered map traversal could repeatedly
    favor the same jobs, leaving other runnable jobs pending indefinitely.
    
    This PR keeps the direct-handoff pool and makes scheduling bounded and
    retryable:
    
    - Scan runnable jobs once and keep only the jobs needed by the current
    slots in a bounded priority heap.
    - Reserve the active slot before submission and roll it back when
    submission is rejected.
    - Schedule the least-recently-run jobs first, with never-scheduled
    one-time jobs taking the first turn.
    - Surface pool rejection so the scheduler can retry the job in the next
    cycle.
    - Apply mutable scheduler concurrency and interval configuration without
    requiring an FE restart.
    
    ### Release note
    
    Fix cloud warm-up jobs remaining pending when the scheduler thread pool
    is saturated.
---
 .../main/java/org/apache/doris/common/Config.java  |  14 +-
 .../java/org/apache/doris/common/ConfigTest.java   |  21 +++
 .../apache/doris/cloud/CacheHotspotManager.java    | 120 ++++++++++++++---
 .../cloud/CacheHotspotManagerSchedulerTest.java    | 149 +++++++++++++++++++++
 4 files changed, 282 insertions(+), 22 deletions(-)

diff --git a/fe/fe-common/src/main/java/org/apache/doris/common/Config.java 
b/fe/fe-common/src/main/java/org/apache/doris/common/Config.java
index 4f8918a54c9..de5d152a009 100644
--- a/fe/fe-common/src/main/java/org/apache/doris/common/Config.java
+++ b/fe/fe-common/src/main/java/org/apache/doris/common/Config.java
@@ -3345,9 +3345,21 @@ public class Config extends ConfigBase {
     @ConfField(mutable = true, masterOnly = true)
     public static int cloud_warm_up_timeout_second = 86400 * 30; // 30 days
 
-    @ConfField(mutable = true, masterOnly = true)
+    @ConfField(mutable = true, masterOnly = true,
+            callback = PositiveCloudWarmUpSchedulerIntervalConfHandler.class)
     public static int cloud_warm_up_job_scheduler_interval_millisecond = 1000; 
// 1 seconds
 
+    public static class PositiveCloudWarmUpSchedulerIntervalConfHandler 
implements ConfHandler {
+        @Override
+        public void handle(Field field, String value) throws Exception {
+            int parsedValue = Integer.parseInt(value.trim());
+            if (parsedValue <= 0) {
+                throw new ConfigException(field.getName() + " must be greater 
than 0");
+            }
+            field.setInt(null, parsedValue);
+        }
+    }
+
     @ConfField(mutable = true, masterOnly = true)
     public static long cloud_warm_up_job_max_bytes_per_batch = 21474836480L; 
// 20GB
 
diff --git a/fe/fe-common/src/test/java/org/apache/doris/common/ConfigTest.java 
b/fe/fe-common/src/test/java/org/apache/doris/common/ConfigTest.java
index 64778f41091..809f0512b7c 100644
--- a/fe/fe-common/src/test/java/org/apache/doris/common/ConfigTest.java
+++ b/fe/fe-common/src/test/java/org/apache/doris/common/ConfigTest.java
@@ -165,6 +165,27 @@ public class ConfigTest {
         }
     }
 
+    @Test
+    public void testCloudWarmUpSchedulerIntervalMustBePositive() throws 
ConfigException {
+        int original = Config.cloud_warm_up_job_scheduler_interval_millisecond;
+        try {
+            
ConfigBase.setMutableConfig("cloud_warm_up_job_scheduler_interval_millisecond", 
"2000");
+            Assert.assertEquals(2000, 
Config.cloud_warm_up_job_scheduler_interval_millisecond);
+
+            ConfigException zeroException = 
Assert.assertThrows(ConfigException.class,
+                    () -> 
ConfigBase.setMutableConfig("cloud_warm_up_job_scheduler_interval_millisecond", 
"0"));
+            Assert.assertTrue(zeroException.getMessage().contains("must be 
greater than 0"));
+            Assert.assertEquals(2000, 
Config.cloud_warm_up_job_scheduler_interval_millisecond);
+
+            ConfigException negativeException = 
Assert.assertThrows(ConfigException.class,
+                    () -> 
ConfigBase.setMutableConfig("cloud_warm_up_job_scheduler_interval_millisecond", 
"-1"));
+            Assert.assertTrue(negativeException.getMessage().contains("must be 
greater than 0"));
+            Assert.assertEquals(2000, 
Config.cloud_warm_up_job_scheduler_interval_millisecond);
+        } finally {
+            Config.cloud_warm_up_job_scheduler_interval_millisecond = original;
+        }
+    }
+
     @Test
     public void testValidateWebSqlStartupConfig() throws ConfigException {
         int originalIdleTimeout = Config.web_sql_session_idle_timeout_seconds;
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/cloud/CacheHotspotManager.java 
b/fe/fe-core/src/main/java/org/apache/doris/cloud/CacheHotspotManager.java
index bd08f3847d6..528837151e4 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/cloud/CacheHotspotManager.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/cloud/CacheHotspotManager.java
@@ -86,6 +86,7 @@ import java.util.Map;
 import java.util.Map.Entry;
 import java.util.Objects;
 import java.util.Optional;
+import java.util.PriorityQueue;
 import java.util.Set;
 import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.ConcurrentHashMap;
@@ -95,9 +96,11 @@ import java.util.concurrent.ExecutorCompletionService;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
 import java.util.concurrent.Future;
+import java.util.concurrent.RejectedExecutionException;
 import java.util.concurrent.ThreadPoolExecutor;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.TimeoutException;
+import java.util.concurrent.atomic.AtomicLong;
 import java.util.concurrent.locks.ReentrantLock;
 import java.util.stream.Collectors;
 
@@ -139,11 +142,16 @@ public class CacheHotspotManager extends MasterDaemon {
 
     private ConcurrentMap<Long, CloudWarmUpJob> runnableCloudWarmUpJobs = 
Maps.newConcurrentMap();
 
+    // Keep scheduling history in memory only. Jobs that have never been 
scheduled have sequence 0
+    // and are selected before jobs that ran in previous cycles.
+    private final ConcurrentMap<Long, Long> cloudWarmUpJobLastScheduleSeq = 
Maps.newConcurrentMap();
+
+    private final AtomicLong cloudWarmUpJobScheduleSeq = new AtomicLong(0);
+
     private final ConcurrentMap<OncePendingJobKey, 
RefCountedPendingCreateLock> oncePendingCreateLocks
             = Maps.newConcurrentMap();
 
-    private final ThreadPoolExecutor cloudWarmUpThreadPool = 
ThreadPoolManager.newDaemonCacheThreadPool(
-            Config.max_active_cloud_warm_up_job, "cloud-warm-up-pool", true);
+    private final ThreadPoolExecutor cloudWarmUpThreadPool;
 
     private static class JobKey {
         private final String srcName;
@@ -593,8 +601,15 @@ public class CacheHotspotManager extends MasterDaemon {
     }
 
     public CacheHotspotManager(CloudSystemInfoService nodeMgr) {
+        this(nodeMgr, ThreadPoolManager.newDaemonCacheThreadPoolThrowException(
+                Config.max_active_cloud_warm_up_job, "cloud-warm-up-pool", 
true));
+    }
+
+    @VisibleForTesting
+    CacheHotspotManager(CloudSystemInfoService nodeMgr, ThreadPoolExecutor 
cloudWarmUpThreadPool) {
         super("CacheHotspotManager", 
Config.fetch_cluster_cache_hotspot_interval_ms);
         this.nodeMgr = nodeMgr;
+        this.cloudWarmUpThreadPool = cloudWarmUpThreadPool;
     }
 
     @Override
@@ -1003,6 +1018,10 @@ public class CacheHotspotManager extends MasterDaemon {
 
         @Override
         public void runAfterCatalogReady() {
+            if (getInterval() != 
Config.cloud_warm_up_job_scheduler_interval_millisecond) {
+                
setInterval(Config.cloud_warm_up_job_scheduler_interval_millisecond);
+                LOG.info("update cloud warm up job daemon interval to {}ms", 
getInterval());
+            }
             if (cycleCount >= CYCLE_COUNT_TO_CHECK_EXPIRE_CLOUD_WARM_UP_JOB) {
                 clearFinishedOrCancelCloudWarmUpJob();
                 cycleCount = 0;
@@ -1202,6 +1221,7 @@ public class CacheHotspotManager extends MasterDaemon {
             CloudWarmUpJob cloudWarmUpJob = iterator.next().getValue();
             if (cloudWarmUpJob.isDone()) {
                 iterator.remove();
+                
cloudWarmUpJobLastScheduleSeq.remove(cloudWarmUpJob.getJobId());
             }
         }
         Iterator<Map.Entry<Long, CloudWarmUpJob>> iterator2 = 
cloudWarmUpJobs.entrySet().iterator();
@@ -1504,28 +1524,86 @@ public class CacheHotspotManager extends MasterDaemon {
         }
     }
 
-    private void runCloudWarmUpJob() {
-        runnableCloudWarmUpJobs.values().forEach(cloudWarmUpJob -> {
-            if (cloudWarmUpJob.shouldWait()) {
-                return;
+    @VisibleForTesting
+    void runCloudWarmUpJob() {
+        int maxActiveJobs = Config.max_active_cloud_warm_up_job;
+        if (maxActiveJobs <= 0) {
+            return;
+        }
+
+        if (cloudWarmUpThreadPool.getMaximumPoolSize() != maxActiveJobs) {
+            cloudWarmUpThreadPool.setMaximumPoolSize(maxActiveJobs);
+            LOG.info("resize cloud warm up thread pool to {}", maxActiveJobs);
+        }
+
+        int availableSlots = maxActiveJobs - activeCloudWarmUpJobs.size();
+        if (availableSlots <= 0) {
+            return;
+        }
+
+        // A smaller last-scheduled sequence has higher priority, and 
never-scheduled jobs use sequence 0.
+        // For the same sequence, prefer ONCE jobs, then earlier creation 
time, then smaller job ID.
+        Comparator<CloudWarmUpJob> schedulePriority = Comparator
+                .comparingLong((CloudWarmUpJob job) ->
+                        
cloudWarmUpJobLastScheduleSeq.getOrDefault(job.getJobId(), 0L))
+                .thenComparingInt(job -> job.isOnce() ? 0 : 1)
+                .thenComparingLong(CloudWarmUpJob::getCreateTimeMs)
+                .thenComparingLong(CloudWarmUpJob::getJobId);
+
+        // Keep only the highest-priority jobs needed by this cycle. The 
reversed comparator keeps
+        // the lowest-priority selected job at the heap top so it can be 
replaced during the scan.
+        PriorityQueue<CloudWarmUpJob> candidates = new 
PriorityQueue<>(schedulePriority.reversed());
+        for (CloudWarmUpJob job : runnableCloudWarmUpJobs.values()) {
+            if (job.shouldWait() || job.isDone() || 
activeCloudWarmUpJobs.containsKey(job.getJobId())) {
+                continue;
             }
-            if (!cloudWarmUpJob.isDone() && 
!activeCloudWarmUpJobs.containsKey(cloudWarmUpJob.getJobId())
-                    && activeCloudWarmUpJobs.size() < 
Config.max_active_cloud_warm_up_job) {
-                if (FeConstants.runningUnitTest) {
-                    cloudWarmUpJob.run();
-                } else {
-                    cloudWarmUpThreadPool.submit(() -> {
-                        if 
(activeCloudWarmUpJobs.putIfAbsent(cloudWarmUpJob.getJobId(), cloudWarmUpJob) 
== null) {
-                            try {
-                                cloudWarmUpJob.run();
-                            } finally {
-                                
activeCloudWarmUpJobs.remove(cloudWarmUpJob.getJobId());
-                            }
-                        }
-                    });
+            if (candidates.size() < availableSlots) {
+                candidates.offer(job);
+            } else if (schedulePriority.compare(job, candidates.peek()) < 0) {
+                candidates.poll();
+                candidates.offer(job);
+            }
+        }
+
+        List<CloudWarmUpJob> jobsToSchedule = new ArrayList<>(candidates);
+        jobsToSchedule.sort(schedulePriority);
+
+        for (CloudWarmUpJob job : jobsToSchedule) {
+            if (availableSlots <= 0) {
+                break;
+            }
+            long jobId = job.getJobId();
+            if (activeCloudWarmUpJobs.putIfAbsent(jobId, job) != null) {
+                continue;
+            }
+
+            if (FeConstants.runningUnitTest) {
+                cloudWarmUpJobLastScheduleSeq.put(jobId, 
cloudWarmUpJobScheduleSeq.incrementAndGet());
+                try {
+                    job.run();
+                } finally {
+                    activeCloudWarmUpJobs.remove(jobId, job);
                 }
+                --availableSlots;
+                continue;
             }
-        });
+
+            try {
+                cloudWarmUpThreadPool.execute(() -> {
+                    try {
+                        job.run();
+                    } finally {
+                        activeCloudWarmUpJobs.remove(jobId, job);
+                    }
+                });
+                cloudWarmUpJobLastScheduleSeq.put(jobId, 
cloudWarmUpJobScheduleSeq.incrementAndGet());
+                --availableSlots;
+            } catch (RejectedExecutionException e) {
+                activeCloudWarmUpJobs.remove(jobId, job);
+                LOG.warn("failed to schedule cloud warm up job {}, retry in 
next cycle", jobId);
+                break;
+            }
+        }
     }
 
     public void replayCloudWarmUpJob(CloudWarmUpJob cloudWarmUpJob) throws 
Exception {
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/cloud/CacheHotspotManagerSchedulerTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/cloud/CacheHotspotManagerSchedulerTest.java
new file mode 100644
index 00000000000..b974d6a9d85
--- /dev/null
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/cloud/CacheHotspotManagerSchedulerTest.java
@@ -0,0 +1,149 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package org.apache.doris.cloud;
+
+import org.apache.doris.cloud.system.CloudSystemInfoService;
+import org.apache.doris.common.Config;
+import org.apache.doris.common.FeConstants;
+
+import org.junit.After;
+import org.junit.Assert;
+import org.junit.Before;
+import org.junit.Test;
+import org.mockito.Mockito;
+
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.List;
+import java.util.concurrent.RejectedExecutionException;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.atomic.AtomicInteger;
+
+public class CacheHotspotManagerSchedulerTest {
+    private boolean originalRunningUnitTest;
+    private int originalMaxActiveCloudWarmUpJob;
+    private ThreadPoolExecutor executor;
+    private CacheHotspotManager manager;
+
+    @Before
+    public void setUp() {
+        originalRunningUnitTest = FeConstants.runningUnitTest;
+        originalMaxActiveCloudWarmUpJob = Config.max_active_cloud_warm_up_job;
+        FeConstants.runningUnitTest = false;
+        Config.max_active_cloud_warm_up_job = 2;
+
+        executor = Mockito.mock(ThreadPoolExecutor.class);
+        Mockito.when(executor.getMaximumPoolSize()).thenReturn(2);
+        manager = new 
CacheHotspotManager(Mockito.mock(CloudSystemInfoService.class), executor);
+    }
+
+    @After
+    public void tearDown() {
+        FeConstants.runningUnitTest = originalRunningUnitTest;
+        Config.max_active_cloud_warm_up_job = originalMaxActiveCloudWarmUpJob;
+    }
+
+    @Test
+    public void testNewOnceJobGetsFirstTurnAndJobsRotate() throws Exception {
+        List<Long> runOrder = new ArrayList<>();
+        Mockito.doAnswer(invocation -> {
+            ((Runnable) invocation.getArgument(0)).run();
+            return null;
+        }).when(executor).execute(Mockito.any(Runnable.class));
+
+        manager.addCloudWarmUpJob(mockJob(1L, false, 1L, runOrder));
+        manager.addCloudWarmUpJob(mockJob(2L, false, 2L, runOrder));
+        manager.addCloudWarmUpJob(mockJob(3L, false, 3L, runOrder));
+        manager.addCloudWarmUpJob(mockJob(4L, true, 4L, runOrder));
+
+        manager.runCloudWarmUpJob();
+        manager.runCloudWarmUpJob();
+        manager.runCloudWarmUpJob();
+        manager.runCloudWarmUpJob();
+
+        Assert.assertEquals(Arrays.asList(4L, 1L, 2L, 3L, 4L, 1L, 2L, 3L), 
runOrder);
+    }
+
+    @Test
+    public void testActiveJobIsNotSubmittedAgain() throws Exception {
+        List<Runnable> submittedTasks = new ArrayList<>();
+        Mockito.doAnswer(invocation -> {
+            submittedTasks.add(invocation.getArgument(0));
+            return null;
+        }).when(executor).execute(Mockito.any(Runnable.class));
+
+        CloudWarmUpJob job = mockJob(1L, true, 1L, new ArrayList<>());
+        manager.addCloudWarmUpJob(job);
+
+        manager.runCloudWarmUpJob();
+        manager.runCloudWarmUpJob();
+        Assert.assertEquals(1, submittedTasks.size());
+
+        submittedTasks.get(0).run();
+        manager.runCloudWarmUpJob();
+        Assert.assertEquals(2, submittedTasks.size());
+        submittedTasks.get(1).run();
+        Mockito.verify(job, Mockito.times(2)).run();
+    }
+
+    @Test
+    public void testRejectedJobIsRetried() throws Exception {
+        AtomicInteger submitCount = new AtomicInteger();
+        Mockito.doAnswer(invocation -> {
+            if (submitCount.incrementAndGet() == 1) {
+                throw new RejectedExecutionException("injected rejection");
+            }
+            ((Runnable) invocation.getArgument(0)).run();
+            return null;
+        }).when(executor).execute(Mockito.any(Runnable.class));
+
+        CloudWarmUpJob job = mockJob(1L, true, 1L, new ArrayList<>());
+        manager.addCloudWarmUpJob(job);
+
+        manager.runCloudWarmUpJob();
+        Mockito.verify(job, Mockito.never()).run();
+        manager.runCloudWarmUpJob();
+
+        Assert.assertEquals(2, submitCount.get());
+        Mockito.verify(job, Mockito.times(1)).run();
+    }
+
+    @Test
+    public void testThreadPoolSizeFollowsMutableConfig() {
+        Config.max_active_cloud_warm_up_job = 3;
+        Mockito.when(executor.getMaximumPoolSize()).thenReturn(1);
+
+        manager.runCloudWarmUpJob();
+
+        Mockito.verify(executor).setMaximumPoolSize(3);
+    }
+
+    private CloudWarmUpJob mockJob(long jobId, boolean once, long 
createTimeMs, List<Long> runOrder) {
+        CloudWarmUpJob job = Mockito.mock(CloudWarmUpJob.class);
+        Mockito.when(job.getJobId()).thenReturn(jobId);
+        Mockito.when(job.isOnce()).thenReturn(once);
+        Mockito.when(job.getCreateTimeMs()).thenReturn(createTimeMs);
+        Mockito.when(job.shouldWait()).thenReturn(false);
+        Mockito.when(job.isDone()).thenReturn(false);
+        Mockito.doAnswer(invocation -> {
+            runOrder.add(jobId);
+            return null;
+        }).when(job).run();
+        return job;
+    }
+}


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to