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 c0f2a1b76be716..d39d333900a76a 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 @@ -3331,9 +3331,21 @@ public static int metaServiceRpcRetryTimes() { @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 395ff41f620a0b..e64e778c4280b1 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 @@ -167,6 +167,27 @@ public void testSetWebSqlMaxResultBytes() throws ConfigException { } } + @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 bd08f3847d639f..528837151e467f 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.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.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 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 cloudWarmUpJobLastScheduleSeq = Maps.newConcurrentMap(); + + private final AtomicLong cloudWarmUpJobScheduleSeq = new AtomicLong(0); + private final ConcurrentMap 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 void notifyJobStop(CloudWarmUpJob job) { } 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 @@ private class JobDaemon 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 @@ private void clearFinishedOrCancelCloudWarmUpJob() { CloudWarmUpJob cloudWarmUpJob = iterator.next().getValue(); if (cloudWarmUpJob.isDone()) { iterator.remove(); + cloudWarmUpJobLastScheduleSeq.remove(cloudWarmUpJob.getJobId()); } } Iterator> iterator2 = cloudWarmUpJobs.entrySet().iterator(); @@ -1504,28 +1524,86 @@ public void cancelTableFilterJobsForClusterChange(String clusterName, String rea } } - 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 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 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 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 00000000000000..b974d6a9d85fee --- /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 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 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 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; + } +}