Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 13 additions & 1 deletion fe/fe-common/src/main/java/org/apache/doris/common/Config.java
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
21 changes: 21 additions & 0 deletions fe/fe-common/src/test/java/org/apache/doris/common/ConfigTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;

Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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);
Comment thread
bobhan1 marked this conversation as resolved.
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;
Expand Down Expand Up @@ -1202,6 +1221,7 @@ private void clearFinishedOrCancelCloudWarmUpJob() {
CloudWarmUpJob cloudWarmUpJob = iterator.next().getValue();
if (cloudWarmUpJob.isDone()) {
iterator.remove();
cloudWarmUpJobLastScheduleSeq.remove(cloudWarmUpJob.getJobId());
}
}
Iterator<Map.Entry<Long, CloudWarmUpJob>> iterator2 = cloudWarmUpJobs.entrySet().iterator();
Expand Down Expand Up @@ -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<CloudWarmUpJob> schedulePriority = Comparator
.comparingLong((CloudWarmUpJob job) ->
cloudWarmUpJobLastScheduleSeq.getOrDefault(job.getJobId(), 0L))
Comment thread
bobhan1 marked this conversation as resolved.
.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(() -> {
Comment thread
bobhan1 marked this conversation as resolved.
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 {
Expand Down
Original file line number Diff line number Diff line change
@@ -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);
Comment thread
bobhan1 marked this conversation as resolved.
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;
}
}
Loading