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 73805a950ceae8..2e30928c8381ac 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 @@ -325,7 +325,16 @@ private void syncFileCacheTasksForVirtualGroup(Cloud.ClusterPB virtualGroupInMs, long jobIdEvent = cacheHotspotManager.createJob(eventStmtPeriodic); // send jobIds to ms List 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 5660d9697e71ea..0eccb05c4c7bba 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 @@ -544,7 +544,7 @@ public synchronized void updateFrontends(List toAdd, List to } } - public static void updateFileCacheJobIds(ComputeGroup cg, List jobIds) { + public static boolean updateFileCacheJobIds(ComputeGroup cg, List jobIds) { Cloud.ClusterPolicy policy = Cloud.ClusterPolicy.newBuilder() .setType(Cloud.ClusterPolicy.PolicyType.ActiveStandby) .addAllCacheWarmupJobids(jobIds).build(); @@ -566,9 +566,12 @@ public static void updateFileCacheJobIds(ComputeGroup cg, List jobIds) { 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 95eb739de64867..a38b7ad277d528 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 void testSyncInstanceCreatesVirtualComputeGroupAndCancelsTableLevelLoadEv try (MockedStatic 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( @@ -172,12 +172,98 @@ public void testSyncInstanceCreatesVirtualComputeGroupAndCancelsTableLevelLoadEv 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 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 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 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 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 ComputeGroup(computeGroupId, computeGroupName, ComputeGroup.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 +274,8 @@ private Cloud.GetInstanceResponse instanceResponseWithVirtualComputeGroup() { .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 +292,22 @@ private Cloud.GetInstanceResponse instanceResponseWithVirtualComputeGroup() { .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 +346,13 @@ private WarmUpClusterCommand buildEventDrivenStmt(String src, String dst, TableF properties, Arrays.asList(rules)); } + private WarmUpClusterCommand buildPeriodicStmt(String src, String dst) { + Map 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 messages = new ArrayList<>();