This is an automated email from the ASF dual-hosted git repository.
morningman pushed a commit to branch branch-4.0
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.0 by this push:
new c6d9f2ec912 [fix](fe) Cancel rebuilt VCG warm up jobs on drop (#65426)
(#66166)
c6d9f2ec912 is described below
commit c6d9f2ec91211ca95aabf840fcece6c325d09ee0
Author: deardeng <[email protected]>
AuthorDate: Wed Jul 29 23:14:15 2026 +0800
[fix](fe) Cancel rebuilt VCG warm up jobs on drop (#65426) (#66166)
pick from https://github.com/apache/doris/pull/65426
---
.../cloud/catalog/CloudInstanceStatusChecker.java | 16 ++-
.../doris/cloud/system/CloudSystemInfoService.java | 5 +-
.../catalog/CloudInstanceStatusCheckerTest.java | 119 ++++++++++++++++++++-
3 files changed, 132 insertions(+), 8 deletions(-)
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudInstanceStatusChecker.java
b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudInstanceStatusChecker.java
index 4b3b7cf79b7..05757cbca9f 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudInstanceStatusChecker.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudInstanceStatusChecker.java
@@ -316,9 +316,19 @@ public class CloudInstanceStatusChecker extends
MasterDaemon {
long jobIdEvent =
cacheHotspotManager.createJob(eventStmtPeriodic);
// send jobIds to ms
List<String> newJobIds =
Arrays.asList(Long.toString(jobIdPeriodic), Long.toString(jobIdEvent));
- CloudSystemInfoService.updateFileCacheJobIds(virtualGroupInFe,
newJobIds);
- LOG.info("virtual compute group {}, generate new jobIds
periodic={}, event={}, and old jobIds {}",
- virtualGroupInFe, jobIdPeriodic, jobIdEvent,
jobIdsInMs);
+ boolean updated =
CloudSystemInfoService.updateFileCacheJobIds(virtualGroupInFe, newJobIds);
+ if (!updated) {
+ LOG.warn("warmup-vcg rebuild-failed vcgName={}
srcCluster={} dstCluster={} "
+ + "createdPeriodicJobId={}
createdEventJobId={} oldJobIds={} "
+ + "failureReason=failed to update new job
ids to ms",
+ virtualGroupInFe.getName(), srcCg, dstCg,
jobIdPeriodic, jobIdEvent, jobIdsInMs);
+ cancelCacheJobs(virtualGroupInFe, newJobIds);
+ return;
+ }
+ virtualGroupInFe.getPolicy().setCacheWarmupJobIds(newJobIds);
+ LOG.info("warmup-vcg rebuild-finish vcgName={} srcCluster={}
dstCluster={} "
+ + "createdPeriodicJobId={}
createdEventJobId={} oldJobIds={}",
+ virtualGroupInFe.getName(), srcCg, dstCg,
jobIdPeriodic, jobIdEvent, jobIdsInMs);
} catch (AnalysisException e) {
LOG.warn("virtual compute err, name: {}, failed to generate
file cache warm up jobs: {}",
virtualGroupInFe.getName(), e.getMessage(), e);
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/cloud/system/CloudSystemInfoService.java
b/fe/fe-core/src/main/java/org/apache/doris/cloud/system/CloudSystemInfoService.java
index c334db2b824..3b6b0c8d407 100644
---
a/fe/fe-core/src/main/java/org/apache/doris/cloud/system/CloudSystemInfoService.java
+++
b/fe/fe-core/src/main/java/org/apache/doris/cloud/system/CloudSystemInfoService.java
@@ -687,7 +687,7 @@ public class CloudSystemInfoService extends
SystemInfoService {
}
}
- public static void updateFileCacheJobIds(ComputeGroup cg, List<String>
jobIds) {
+ public static boolean updateFileCacheJobIds(ComputeGroup cg, List<String>
jobIds) {
Cloud.ClusterPolicy policy = Cloud.ClusterPolicy.newBuilder()
.setType(Cloud.ClusterPolicy.PolicyType.ActiveStandby)
.addAllCacheWarmupJobids(jobIds).build();
@@ -709,9 +709,12 @@ public class CloudSystemInfoService extends
SystemInfoService {
LOG.info("update file cache jobIds, request: {}, response: {}",
request, response);
if (response.getStatus().getCode() != Cloud.MetaServiceCode.OK) {
LOG.warn("update file cache jobIds, response: {}", response);
+ return false;
}
+ return true;
} catch (RpcException e) {
LOG.warn("failed to update file cache jobIds {}", cg, e);
+ return false;
}
}
diff --git
a/fe/fe-core/src/test/java/org/apache/doris/cloud/catalog/CloudInstanceStatusCheckerTest.java
b/fe/fe-core/src/test/java/org/apache/doris/cloud/catalog/CloudInstanceStatusCheckerTest.java
index ff19f67dcb6..24747b1882f 100644
---
a/fe/fe-core/src/test/java/org/apache/doris/cloud/catalog/CloudInstanceStatusCheckerTest.java
+++
b/fe/fe-core/src/test/java/org/apache/doris/cloud/catalog/CloudInstanceStatusCheckerTest.java
@@ -129,7 +129,7 @@ public class CloudInstanceStatusCheckerTest {
try (MockedStatic<CloudSystemInfoService> mockedCloudSystemInfoService
=
Mockito.mockStatic(CloudSystemInfoService.class,
Mockito.CALLS_REAL_METHODS)) {
mockedCloudSystemInfoService.when(() ->
CloudSystemInfoService.updateFileCacheJobIds(
- Mockito.any(ComputeGroup.class),
Mockito.anyList())).thenAnswer(invocation -> null);
+ Mockito.any(ComputeGroup.class),
Mockito.anyList())).thenReturn(true);
new
CloudInstanceStatusChecker(cloudSystemInfoService).runAfterCatalogReady();
mockedCloudSystemInfoService.verify(() ->
CloudSystemInfoService.updateFileCacheJobIds(
@@ -166,7 +166,91 @@ public class CloudInstanceStatusCheckerTest {
String logs = appender.messagesAsString();
Assertions.assertFalse(logs.contains("failed to create virtual compute
group vcg"), logs);
- Assertions.assertTrue(logs.contains("generate new jobIds"), logs);
+ Assertions.assertTrue(logs.contains("warmup-vcg rebuild-finish"),
logs);
+ Assertions.assertTrue(logs.contains("srcCluster=active_cg"), logs);
+ Assertions.assertTrue(logs.contains("dstCluster=standby_cg"), logs);
+ }
+
+ @Test
+ public void testDropVirtualComputeGroupCancelsRebuiltWarmUpJobs() throws
Exception {
+ addComputeGroup("active_cg_id", "active_cg");
+ addComputeGroup("standby_cg_id", "standby_cg");
+ long oldPeriodicJobId = cacheHotspotManager.createJob(
+ buildPeriodicStmt("active_cg", "standby_cg"));
+ long oldEventJobId = cacheHotspotManager.createJob(
+ buildEventDrivenStmt("active_cg", "standby_cg"));
+
+ ComputeGroup virtualComputeGroup = new ComputeGroup("vcg_id", "vcg",
ComputeGroup.ComputeTypeEnum.VIRTUAL);
+ virtualComputeGroup.setSubComputeGroups(Arrays.asList("active_cg",
"standby_cg"));
+ ComputeGroup.Policy policy = new ComputeGroup.Policy();
+ policy.setActiveComputeGroup("active_cg");
+ policy.setStandbyComputeGroup("standby_cg");
+ policy.setCacheWarmupJobIds(Arrays.asList(
+ Long.toString(oldPeriodicJobId),
Long.toString(oldEventJobId)));
+ virtualComputeGroup.setPolicy(policy);
+ cloudSystemInfoService.addComputeGroup("vcg_id", virtualComputeGroup);
+
+ Mockito.when(cloudEnv.isMaster()).thenReturn(true);
+ Mockito.doReturn(instanceResponseWithVirtualComputeGroup("standby_cg",
"active_cg"))
+ .doReturn(instanceResponseWithoutVirtualComputeGroup())
+ .when(cloudSystemInfoService).getCloudInstance();
+
+ try (MockedStatic<CloudSystemInfoService> mockedCloudSystemInfoService
=
+ Mockito.mockStatic(CloudSystemInfoService.class,
Mockito.CALLS_REAL_METHODS)) {
+ mockedCloudSystemInfoService.when(() ->
CloudSystemInfoService.updateFileCacheJobIds(
+ Mockito.any(ComputeGroup.class),
Mockito.anyList())).thenReturn(true);
+
+ CloudInstanceStatusChecker checker = new
CloudInstanceStatusChecker(cloudSystemInfoService);
+ checker.runAfterCatalogReady();
+ List<Long> rebuiltJobIds =
cacheHotspotManager.getCloudWarmUpJobs().values().stream()
+ .filter(job -> job.getJobType() ==
CloudWarmUpJob.JobType.CLUSTER)
+ .filter(job ->
"standby_cg".equals(job.getSrcClusterName()))
+ .filter(job -> "active_cg".equals(job.getDstClusterName()))
+ .map(CloudWarmUpJob::getJobId)
+ .collect(java.util.stream.Collectors.toList());
+ Assertions.assertEquals(2, rebuiltJobIds.size());
+
+ checker.runAfterCatalogReady();
+
+ for (long rebuiltJobId : rebuiltJobIds) {
+ CloudWarmUpJob job =
cacheHotspotManager.getCloudWarmUpJob(rebuiltJobId);
+ Assertions.assertEquals(CloudWarmUpJob.JobState.CANCELLED,
job.getJobState(),
+ "rebuilt warm up job should be cancelled when VCG is
dropped: " + rebuiltJobId);
+ }
+ }
+ }
+
+ @Test
+ public void testRebuildCancelsNewWarmUpJobsWhenUpdateJobIdsFails() {
+ addComputeGroup("active_cg_id", "active_cg");
+ addComputeGroup("standby_cg_id", "standby_cg");
+ Mockito.when(cloudEnv.isMaster()).thenReturn(true);
+
Mockito.doReturn(instanceResponseWithVirtualComputeGroup()).when(cloudSystemInfoService).getCloudInstance();
+
+ try (MockedStatic<CloudSystemInfoService> mockedCloudSystemInfoService
=
+ Mockito.mockStatic(CloudSystemInfoService.class,
Mockito.CALLS_REAL_METHODS)) {
+ mockedCloudSystemInfoService.when(() ->
CloudSystemInfoService.updateFileCacheJobIds(
+ Mockito.any(ComputeGroup.class),
Mockito.anyList())).thenReturn(false);
+
+ new
CloudInstanceStatusChecker(cloudSystemInfoService).runAfterCatalogReady();
+ mockedCloudSystemInfoService.verify(() ->
CloudSystemInfoService.updateFileCacheJobIds(
+ Mockito.any(ComputeGroup.class), Mockito.anyList()));
+ }
+
+ ComputeGroup virtualComputeGroup =
cloudSystemInfoService.getComputeGroupById("vcg_id");
+ Assertions.assertNotNull(virtualComputeGroup);
+ Assertions.assertTrue(virtualComputeGroup.isNeedRebuildFileCache());
+
Assertions.assertTrue(virtualComputeGroup.getPolicy().getCacheWarmupJobIds().isEmpty());
+
+ List<CloudWarmUpJob> newWarmUpJobs =
cacheHotspotManager.getCloudWarmUpJobs().values().stream()
+ .filter(job -> job.getJobType() ==
CloudWarmUpJob.JobType.CLUSTER)
+ .filter(job -> "active_cg".equals(job.getSrcClusterName()))
+ .filter(job -> "standby_cg".equals(job.getDstClusterName()))
+ .collect(java.util.stream.Collectors.toList());
+ Assertions.assertEquals(2, newWarmUpJobs.size());
+ for (CloudWarmUpJob job : newWarmUpJobs) {
+ Assertions.assertEquals(CloudWarmUpJob.JobState.CANCELLED,
job.getJobState());
+ }
}
private void addComputeGroup(String computeGroupId, String
computeGroupName) {
@@ -175,6 +259,10 @@ public class CloudInstanceStatusCheckerTest {
}
private Cloud.GetInstanceResponse
instanceResponseWithVirtualComputeGroup() {
+ return instanceResponseWithVirtualComputeGroup("active_cg",
"standby_cg");
+ }
+
+ private Cloud.GetInstanceResponse
instanceResponseWithVirtualComputeGroup(String active, String standby) {
Cloud.ClusterPB activeComputeGroup = computeGroup("active_cg_id",
"active_cg");
Cloud.ClusterPB standbyComputeGroup = computeGroup("standby_cg_id",
"standby_cg");
Cloud.ClusterPB virtualComputeGroup = Cloud.ClusterPB.newBuilder()
@@ -185,8 +273,8 @@ public class CloudInstanceStatusCheckerTest {
.addClusterNames("standby_cg")
.setClusterPolicy(Cloud.ClusterPolicy.newBuilder()
.setType(Cloud.ClusterPolicy.PolicyType.ActiveStandby)
- .setActiveClusterName("active_cg")
- .addStandbyClusterNames("standby_cg")
+ .setActiveClusterName(active)
+ .addStandbyClusterNames(standby)
.build())
.build();
return Cloud.GetInstanceResponse.newBuilder()
@@ -203,6 +291,22 @@ public class CloudInstanceStatusCheckerTest {
.build();
}
+ private Cloud.GetInstanceResponse
instanceResponseWithoutVirtualComputeGroup() {
+ Cloud.ClusterPB activeComputeGroup = computeGroup("active_cg_id",
"active_cg");
+ Cloud.ClusterPB standbyComputeGroup = computeGroup("standby_cg_id",
"standby_cg");
+ return Cloud.GetInstanceResponse.newBuilder()
+ .setStatus(Cloud.MetaServiceResponseStatus.newBuilder()
+ .setCode(Cloud.MetaServiceCode.OK)
+ .setMsg("OK")
+ .build())
+ .setInstance(Cloud.InstanceInfoPB.newBuilder()
+ .setStatus(Cloud.InstanceInfoPB.Status.NORMAL)
+ .addClusters(activeComputeGroup)
+ .addClusters(standbyComputeGroup)
+ .build())
+ .build();
+ }
+
private Cloud.ClusterPB computeGroup(String computeGroupId, String
computeGroupName) {
return Cloud.ClusterPB.newBuilder()
.setClusterId(computeGroupId)
@@ -241,6 +345,13 @@ public class CloudInstanceStatusCheckerTest {
properties, Arrays.asList(rules));
}
+ private WarmUpClusterCommand buildPeriodicStmt(String src, String dst) {
+ Map<String, String> properties = new HashMap<>();
+ properties.put("sync_mode", "periodic");
+ properties.put("sync_interval_sec", "600");
+ return new WarmUpClusterCommand(new ArrayList<>(), src, dst, false,
false, properties);
+ }
+
private static class RecordingAppender extends AbstractAppender {
private final List<String> messages = new ArrayList<>();
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]