This is an automated email from the ASF dual-hosted git repository.
deardeng 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 60e56ec867a [fix](fe) Cancel rebuilt VCG warm up jobs on drop (#65426)
60e56ec867a is described below
commit 60e56ec867a3b8423b1e650ec199b81acdc1a870
Author: deardeng <[email protected]>
AuthorDate: Fri Jul 24 11:02:13 2026 +0800
[fix](fe) Cancel rebuilt VCG warm up jobs on drop (#65426)
Problem Summary: When active/standby switches for a virtual compute
group, FE rebuilds periodic and event-driven file cache warm up jobs and
writes new job ids to MS. FE memory kept stale cacheWarmupJobIds until a
later MS sync, so a quick DROP VCG cancelled old ids and left rebuilt
jobs running. This change makes updateFileCacheJobIds report MS update
success, synchronizes the FE-local policy with rebuilt job ids after a
successful update, and cancels newly created jobs if the MS update
fails.
---
.../cloud/catalog/CloudInstanceStatusChecker.java | 11 +-
.../doris/cloud/system/CloudSystemInfoService.java | 5 +-
.../catalog/CloudInstanceStatusCheckerTest.java | 116 ++++++++++++++++++++-
3 files changed, 127 insertions(+), 5 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 0c7c7774445..cc9d256d566 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
@@ -326,7 +326,16 @@ 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);
+ 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);
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 b512871d85e..a2f5c5c5100 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
@@ -686,7 +686,7 @@ public class CloudSystemInfoService extends
SystemInfoService {
}
}
- public static void updateFileCacheJobIds(CloudComputeGroupMeta cg,
List<String> jobIds) {
+ public static boolean updateFileCacheJobIds(CloudComputeGroupMeta cg,
List<String> jobIds) {
Cloud.ClusterPolicy policy = Cloud.ClusterPolicy.newBuilder()
.setType(Cloud.ClusterPolicy.PolicyType.ActiveStandby)
.addAllCacheWarmupJobids(jobIds).build();
@@ -708,9 +708,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 fe4060d2ab4..dd36c0f8b3b 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(CloudComputeGroupMeta.class),
Mockito.anyList())).thenAnswer(invocation -> null);
+ Mockito.any(CloudComputeGroupMeta.class),
Mockito.anyList())).thenReturn(true);
new
CloudInstanceStatusChecker(cloudSystemInfoService).runAfterCatalogReady();
mockedCloudSystemInfoService.verify(() ->
CloudSystemInfoService.updateFileCacheJobIds(
@@ -172,12 +172,99 @@ public class CloudInstanceStatusCheckerTest {
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"));
+
+ CloudComputeGroupMeta virtualComputeGroup = new
CloudComputeGroupMeta("vcg_id", "vcg",
+ CloudComputeGroupMeta.ComputeTypeEnum.VIRTUAL);
+ virtualComputeGroup.setSubComputeGroups(Arrays.asList("active_cg",
"standby_cg"));
+ CloudComputeGroupMeta.Policy policy = new
CloudComputeGroupMeta.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(CloudComputeGroupMeta.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(CloudComputeGroupMeta.class),
Mockito.anyList())).thenReturn(false);
+
+ new
CloudInstanceStatusChecker(cloudSystemInfoService).runAfterCatalogReady();
+ mockedCloudSystemInfoService.verify(() ->
CloudSystemInfoService.updateFileCacheJobIds(
+ Mockito.any(CloudComputeGroupMeta.class),
Mockito.anyList()));
+ }
+
+ CloudComputeGroupMeta 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) {
cloudSystemInfoService.addComputeGroup(computeGroupId,
new CloudComputeGroupMeta(computeGroupId, computeGroupName,
CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE));
}
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()
@@ -188,8 +275,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()
@@ -206,6 +293,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)
@@ -244,6 +347,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]