From 660e0b7eee1906a116f189880043be60845a65e3 Mon Sep 17 00:00:00 2001 From: bobhan1 Date: Thu, 17 Sep 2026 11:22:26 +0800 Subject: [PATCH] [fix](cloud) Release warm-up destination on initialization failure (#67924) A cloud warm-up job registers its destination compute group before initializing tablet batches. If initialization throws, the outer `run()` handler only logs the exception: the job stays `PENDING` and keeps the destination registration. Other ONCE/PERIODIC jobs targeting that group cannot start. A later successful retry can recover the original job, but repeated initialization failures can block the group indefinitely because the warm-up timeout only applies to `RUNNING` jobs. Catch initialization failures before transitioning to `RUNNING` and reuse `cancel(..., false)` to persist the error and release the destination registration. ONCE jobs become `CANCELLED`; PERIODIC jobs remain `PENDING` and retry at their existing interval. Initialization has not submitted work to BEs, so this path does not send cleanup RPCs. Successful initialization retains the destination registration as before. Release the destination compute group when cloud warm-up initialization fails, allowing subsequent warm-up jobs to proceed. Report the initialization error and preserve periodic retry scheduling. (cherry picked from commit f0affa953dff3d0e761bd6bb2f62b231b31a9e26) --- .../apache/doris/cloud/CloudWarmUpJob.java | 42 ++++--- .../doris/cloud/CloudWarmUpJobTest.java | 110 ++++++++++++++++++ 2 files changed, 135 insertions(+), 17 deletions(-) diff --git a/fe/fe-core/src/main/java/org/apache/doris/cloud/CloudWarmUpJob.java b/fe/fe-core/src/main/java/org/apache/doris/cloud/CloudWarmUpJob.java index e2f58cc49b04ab..6fbe6188d5e378 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/cloud/CloudWarmUpJob.java +++ b/fe/fe-core/src/main/java/org/apache/doris/cloud/CloudWarmUpJob.java @@ -932,23 +932,31 @@ private void runPendingJob() throws DdlException { return; } - // Todo: nothing to prepare yet - this.setJobDone = false; - this.lastBatchId = -1; - this.startTimeMs = System.currentTimeMillis(); - // reset clients to ensure we have the latest BE info - this.beToThriftAddress = null; - this.beToClient = null; - this.beToAddr = null; - MetricRepo.updateClusterWarmUpJobLatestStartTime(String.valueOf(jobId), srcClusterName, - dstClusterName, startTimeMs); - this.fetchBeToTabletIdBatches(); - long totalTablets = beToTabletIdBatches.values().stream() - .flatMap(List::stream) - .mapToLong(List::size) - .sum(); - MetricRepo.increaseClusterWarmUpJobRequestedTablets(dstClusterName, totalTablets); - MetricRepo.increaseClusterWarmUpJobExecCount(dstClusterName); + long totalTablets; + try { + this.setJobDone = false; + this.lastBatchId = -1; + this.startTimeMs = System.currentTimeMillis(); + // reset clients to ensure we have the latest BE info + this.beToThriftAddress = null; + this.beToClient = null; + this.beToAddr = null; + MetricRepo.updateClusterWarmUpJobLatestStartTime(String.valueOf(jobId), srcClusterName, + dstClusterName, startTimeMs); + this.fetchBeToTabletIdBatches(); + totalTablets = beToTabletIdBatches.values().stream() + .flatMap(List::stream) + .mapToLong(List::size) + .sum(); + MetricRepo.increaseClusterWarmUpJobRequestedTablets(dstClusterName, totalTablets); + MetricRepo.increaseClusterWarmUpJobExecCount(dstClusterName); + } catch (Exception e) { + LOG.warn("failed to initialize cloud warm up job {}", jobId, e); + // No BE job has started. Reuse cancellation to release the destination registration + // and preserve periodic jobs for their next scheduled attempt. + cancel("Failed to initialize warm up job: " + e.getMessage(), false); + return; + } this.jobState = JobState.RUNNING; Env.getCurrentEnv().getEditLog().logModifyCloudWarmUpJob(this); LOG.info("warmup-lock state-transition jobId={} srcCluster={} dstCluster={} syncMode={} jobType={} " diff --git a/fe/fe-core/src/test/java/org/apache/doris/cloud/CloudWarmUpJobTest.java b/fe/fe-core/src/test/java/org/apache/doris/cloud/CloudWarmUpJobTest.java index 800dd89c6b214c..9ace836e14771b 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/cloud/CloudWarmUpJobTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/cloud/CloudWarmUpJobTest.java @@ -56,6 +56,7 @@ import java.util.LinkedHashMap; import java.util.List; import java.util.Map; +import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.atomic.AtomicReference; public class CloudWarmUpJobTest { @@ -186,6 +187,104 @@ public void testPendingRetryKeepsErrMsgWhenJobStarts() throws Exception { Mockito.verify(editLog).logModifyCloudWarmUpJob(job); } + @Test + public void testOncePendingInitializationFailureReleasesDestinationLock() throws Exception { + checkPendingInitializationFailureReleasesDestinationLock(SyncMode.ONCE); + } + + @Test + public void testPeriodicPendingInitializationFailureReleasesDestinationLock() throws Exception { + checkPendingInitializationFailureReleasesDestinationLock(SyncMode.PERIODIC); + } + + private void checkPendingInitializationFailureReleasesDestinationLock(SyncMode syncMode) throws Exception { + CloudWarmUpJob job = Mockito.spy(createPendingJob(204L, syncMode)); + CloudWarmUpJob nextJob = createPendingJob(205L, SyncMode.ONCE); + CloudEnv cloudEnv = Mockito.mock(CloudEnv.class); + CacheHotspotManager manager = new CacheHotspotManager(Mockito.mock(CloudSystemInfoService.class), + Mockito.mock(ThreadPoolExecutor.class)); + EditLog editLog = Mockito.mock(EditLog.class); + Mockito.when(cloudEnv.getCacheHotspotMgr()).thenReturn(manager); + Mockito.when(cloudEnv.getEditLog()).thenReturn(editLog); + Mockito.doAnswer(invocation -> { + Assert.assertFalse(manager.tryRegisterRunningJob(nextJob)); + throw new IllegalStateException("initialization failed"); + }).when(job).fetchBeToTabletIdBatches(); + + try (MockedStatic mockedEnv = Mockito.mockStatic(Env.class)) { + mockedEnv.when(Env::getCurrentEnv).thenReturn(cloudEnv); + job.run(); + + Assert.assertTrue(manager.tryRegisterRunningJob(nextJob)); + Assert.assertEquals(syncMode == SyncMode.ONCE ? JobState.CANCELLED : JobState.PENDING, + job.getJobState()); + Assert.assertEquals("Failed to initialize warm up job: initialization failed", job.getErrMsg()); + Assert.assertTrue(job.getStartTimeMs() > 0); + Assert.assertTrue(job.getFinishedTimeMs() >= job.getStartTimeMs()); + Assert.assertEquals(syncMode == SyncMode.PERIODIC, job.shouldWait()); + Mockito.verify(editLog).logModifyCloudWarmUpJob(job); + + CloudWarmUpJob persistedJob = copyBySerialization(job); + Assert.assertEquals(job.getJobState(), persistedJob.getJobState()); + Assert.assertEquals(job.getErrMsg(), persistedJob.getErrMsg()); + Assert.assertEquals(job.getStartTimeMs(), persistedJob.getStartTimeMs()); + Assert.assertEquals(job.getFinishedTimeMs(), persistedJob.getFinishedTimeMs()); + + nextJob.run(); + Assert.assertEquals(JobState.RUNNING, nextJob.getJobState()); + Mockito.verifyNoInteractions(mockBackendPool); + + if (syncMode == SyncMode.PERIODIC) { + manager.notifyJobStop(nextJob); + Mockito.doCallRealMethod().when(job).fetchBeToTabletIdBatches(); + setStartTimeMs(job, System.currentTimeMillis() - 61_000L); + Assert.assertFalse(job.shouldWait()); + job.run(); + Assert.assertEquals(JobState.RUNNING, job.getJobState()); + Assert.assertFalse(manager.tryRegisterRunningJob(nextJob)); + } else { + job.run(); + Mockito.verify(job).fetchBeToTabletIdBatches(); + Assert.assertEquals(JobState.CANCELLED, job.getJobState()); + } + } + } + + @Test + public void testOncePendingInitializationKeepsDestinationLockOnSuccess() { + checkPendingInitializationKeepsDestinationLockOnSuccess(SyncMode.ONCE); + } + + @Test + public void testPeriodicPendingInitializationKeepsDestinationLockOnSuccess() { + checkPendingInitializationKeepsDestinationLockOnSuccess(SyncMode.PERIODIC); + } + + private void checkPendingInitializationKeepsDestinationLockOnSuccess(SyncMode syncMode) { + CloudWarmUpJob job = createPendingJob(206L, syncMode); + CloudWarmUpJob nextJob = Mockito.spy(createPendingJob(207L, SyncMode.ONCE)); + CloudEnv cloudEnv = Mockito.mock(CloudEnv.class); + CacheHotspotManager manager = new CacheHotspotManager(Mockito.mock(CloudSystemInfoService.class), + Mockito.mock(ThreadPoolExecutor.class)); + EditLog editLog = Mockito.mock(EditLog.class); + Mockito.when(cloudEnv.getCacheHotspotMgr()).thenReturn(manager); + Mockito.when(cloudEnv.getEditLog()).thenReturn(editLog); + + try (MockedStatic mockedEnv = Mockito.mockStatic(Env.class)) { + mockedEnv.when(Env::getCurrentEnv).thenReturn(cloudEnv); + job.run(); + Assert.assertEquals(JobState.RUNNING, job.getJobState()); + + nextJob.run(); + Assert.assertEquals(JobState.PENDING, nextJob.getJobState()); + Assert.assertEquals(-1L, nextJob.getStartTimeMs()); + Assert.assertFalse(manager.tryRegisterRunningJob(nextJob)); + Mockito.verify(nextJob, Mockito.never()).fetchBeToTabletIdBatches(); + Mockito.verify(editLog, Mockito.never()).logModifyCloudWarmUpJob(nextJob); + Mockito.verify(editLog).logModifyCloudWarmUpJob(job); + } + } + @Test public void testEventDrivenSuccessfulRetryClearsErrMsg() throws Exception { CloudSystemInfoService cloudSystemInfoService = Mockito.mock(CloudSystemInfoService.class); @@ -325,6 +424,17 @@ public void testRunningRetryClearsErrMsgWhenJobFinishes() throws Exception { Mockito.verify(mockBackendPool).returnObject(address, client); } + private CloudWarmUpJob createPendingJob(long jobId, SyncMode syncMode) { + return new CloudWarmUpJob.Builder() + .setJobId(jobId) + .setSrcClusterName("source_cluster") + .setDstClusterName("target_cluster") + .setJobType(JobType.CLUSTER) + .setSyncMode(syncMode) + .setSyncInterval(60L) + .build(); + } + private CloudWarmUpJob createRunningJob(long jobId, TNetworkAddress firstAddress, TNetworkAddress secondAddress) { CloudWarmUpJob job = new CloudWarmUpJob.Builder()