From c29ab7bf2d1cac384228b53fbae11847414abf75 Mon Sep 17 00:00:00 2001 From: deardeng Date: Mon, 20 Jul 2026 16:38:27 +0800 Subject: [PATCH] [refactor](cloud) Rename cloud compute group metadata class ### What problem does this PR solve? Issue Number: None Related PR: None Problem Summary: The FE cloud catalog metadata entity shared the ComputeGroup class name with the resource routing abstraction, making imports and usages ambiguous. Rename the cloud metadata entity to CloudComputeGroupMeta and update all production and unit-test references without changing runtime behavior or persisted metadata. ### Release note None ### Check List (For Author) - Test: Unit Test - ./run-fe-ut.sh --run org.apache.doris.cloud.catalog.CloudComputeGroupMetaTest,org.apache.doris.cloud.catalog.CloudInstanceStatusCheckerTest,org.apache.doris.cloud.system.CloudSystemInfoServiceTest,org.apache.doris.cloud.WarmUpClusterOnTablesParseTest,org.apache.doris.mysql.privilege.CloudAuthTest,org.apache.doris.nereids.trees.plans.commands.AlterComputeGroupCommandTest - cd fe && mvn checkstyle:check -pl fe-core - Behavior changed: No - Does this need documentation: No --- .../doris/cloud/catalog/BalanceTypeEnum.java | 3 +- ...eGroup.java => CloudComputeGroupMeta.java} | 28 +++- .../catalog/CloudInstanceStatusChecker.java | 30 ++-- .../cloud/catalog/CloudTabletRebalancer.java | 10 +- .../cloud/system/CloudSystemInfoService.java | 43 +++--- .../commands/AlterComputeGroupCommand.java | 4 +- .../plans/commands/ShowClustersCommand.java | 8 +- .../plans/commands/WarmUpClusterCommand.java | 6 +- .../cloud/WarmUpClusterOnTablesParseTest.java | 20 +-- ...st.java => CloudComputeGroupMetaTest.java} | 135 +++++++++--------- .../CloudInstanceStatusCheckerTest.java | 10 +- .../system/CloudSystemInfoServiceTest.java | 104 +++++++------- .../doris/mysql/privilege/CloudAuthTest.java | 8 +- .../AlterComputeGroupCommandTest.java | 8 +- 14 files changed, 220 insertions(+), 197 deletions(-) rename fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/{ComputeGroup.java => CloudComputeGroupMeta.java} (89%) rename fe/fe-core/src/test/java/org/apache/doris/cloud/catalog/{ComputeGroupTest.java => CloudComputeGroupMetaTest.java} (61%) diff --git a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/BalanceTypeEnum.java b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/BalanceTypeEnum.java index d66e3126d5bb59..55dbfba1cc52ba 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/BalanceTypeEnum.java +++ b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/BalanceTypeEnum.java @@ -64,6 +64,7 @@ public static boolean isValid(String value) { */ public static BalanceTypeEnum getCloudWarmUpForRebalanceTypeEnum() { return fromString(Config.cloud_warm_up_for_rebalance_type) == null - ? ComputeGroup.DEFAULT_COMPUTE_GROUP_BALANCE_ENUM : fromString(Config.cloud_warm_up_for_rebalance_type); + ? CloudComputeGroupMeta.DEFAULT_COMPUTE_GROUP_BALANCE_ENUM + : fromString(Config.cloud_warm_up_for_rebalance_type); } } diff --git a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/ComputeGroup.java b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudComputeGroupMeta.java similarity index 89% rename from fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/ComputeGroup.java rename to fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudComputeGroupMeta.java index c895f13f7a03e8..4d258ed262e391 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/ComputeGroup.java +++ b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudComputeGroupMeta.java @@ -34,8 +34,26 @@ import java.util.List; import java.util.Map; -public class ComputeGroup { - private static final Logger LOG = LogManager.getLogger(ComputeGroup.class); +/** + * FE-side in-memory metadata for a cloud compute group. + * + *

This class models both physical and virtual cloud compute groups and keeps cloud-only + * control-plane state, including the compute group type, active-standby policy, sub compute + * groups, availability timestamps, and cache warm-up properties. Its instances are refreshed + * from the meta service by {@link CloudInstanceStatusChecker} and + * {@link org.apache.doris.cloud.system.CloudSystemInfoService}. + * + *

Do not confuse this class with {@link org.apache.doris.resource.computegroup.ComputeGroup}. + * The resource-layer class is a runtime routing abstraction shared by cloud and non-cloud + * deployments: it selects backends and resolves the workload group namespace for a request. + * In contrast, this class is the long-lived cloud control-plane metadata consulted during + * routing, failover, and cache warm-up. + * + *

The meta service is the source of truth. This class is only an FE in-memory mirror and is + * not persisted in the FE edit log or image. + */ +public class CloudComputeGroupMeta { + private static final Logger LOG = LogManager.getLogger(CloudComputeGroupMeta.class); public static final String BALANCE_TYPE = "balance_type"; @@ -139,7 +157,7 @@ public Cloud.ClusterPolicy toPb() { @Setter private Map properties = new LinkedHashMap<>(ALL_PROPERTIES_DEFAULT_VALUE_MAP); - public ComputeGroup(String id, String name, ComputeTypeEnum type) { + public CloudComputeGroupMeta(String id, String name, ComputeTypeEnum type) { this.id = id; this.name = name; this.type = type; @@ -305,10 +323,10 @@ public boolean equals(Object o) { if (this == o) { return true; } - if (!(o instanceof ComputeGroup)) { + if (!(o instanceof CloudComputeGroupMeta)) { return false; } - ComputeGroup that = (ComputeGroup) o; + CloudComputeGroupMeta that = (CloudComputeGroupMeta) o; return unavailableSince == that.unavailableSince && availableSince == that.availableSince && id.equals(that.id) 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..0c7c7774445252 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 @@ -113,7 +113,7 @@ private void processVirtualClusters(List clusters) { private void handleComputeClusters(List computeClusters) { for (Cloud.ClusterPB computeClusterInMs : computeClusters) { - ComputeGroup computeGroupInFe = cloudSystemInfoService + CloudComputeGroupMeta computeGroupInFe = cloudSystemInfoService .getComputeGroupById(computeClusterInMs.getClusterId()); if (computeGroupInFe == null) { // cluster checker will sync it @@ -131,7 +131,7 @@ private void handleComputeClusters(List computeClusters) { * Compare properties between compute cluster in MS and compute group in FE, * update only the changed key-value pairs to avoid unnecessary updates. */ - private void updatePropertiesIfChanged(ComputeGroup computeGroupInFe, Cloud.ClusterPB computeClusterInMs) { + private void updatePropertiesIfChanged(CloudComputeGroupMeta computeGroupInFe, Cloud.ClusterPB computeClusterInMs) { Map propertiesInMs = computeClusterInMs.getPropertiesMap(); Map propertiesInFe = computeGroupInFe.getProperties(); @@ -181,7 +181,7 @@ private void categorizeClusters(List clusters, private void handleVirtualClusters(List virtualGroups, List computeClusters) { for (Cloud.ClusterPB virtualGroupInMs : virtualGroups) { - ComputeGroup virtualGroupInFe = cloudSystemInfoService + CloudComputeGroupMeta virtualGroupInFe = cloudSystemInfoService .getComputeGroupById(virtualGroupInMs.getClusterId()); if (virtualGroupInFe != null) { handleExistingVirtualComputeGroup(virtualGroupInMs, virtualGroupInFe); @@ -203,7 +203,7 @@ private void handleVirtualClusters(List virtualGroups, List jobIds) { + private void cancelCacheJobs(CloudComputeGroupMeta vcgInFe, List jobIds) { CacheHotspotManager cacheHotspotManager = ((CloudEnv) Env.getCurrentEnv()).getCacheHotspotMgr(); if (!jobIds.isEmpty()) { LOG.info("warmup-vcg cancel-cache-jobs vcgName={} activeComputeGroup={} standbyComputeGroup={} " @@ -224,7 +224,7 @@ private void cancelCacheJobs(ComputeGroup vcgInFe, List jobIds) { } } - private void checkNeedRebuildFileCache(ComputeGroup virtualGroupInFe, List jobIdsInMs) { + private void checkNeedRebuildFileCache(CloudComputeGroupMeta virtualGroupInFe, List jobIdsInMs) { CacheHotspotManager cacheHotspotManager = ((CloudEnv) Env.getCurrentEnv()).getCacheHotspotMgr(); // check jobIds in Ms valid, if been cancelled, start new jobs for (String jobId : jobIdsInMs) { @@ -272,7 +272,8 @@ private void checkNeedRebuildFileCache(ComputeGroup virtualGroupInFe, List subComputeGroups = clusterInMs.getClusterNamesList(); if (subComputeGroups.isEmpty() || virtualGroupInFe.getSubComputeGroups() == null) { LOG.warn("virtual compute err, please check it, verbose {}", virtualGroupInFe); @@ -396,7 +398,7 @@ private boolean areSubComputeGroupsValid(Cloud.ClusterPB clusterInMs, ComputeGro return true; } - private void diffAndUpdateComputeGroup(Cloud.ClusterPB cluster, ComputeGroup computeGroup) { + private void diffAndUpdateComputeGroup(Cloud.ClusterPB cluster, CloudComputeGroupMeta computeGroup) { // vcg rename logic, here cluster_id same, but cluster_name changed, so vcg renamed String clusterNameInMs = cluster.getClusterName(); String computeGroupNameInFe = computeGroup.getName(); @@ -490,10 +492,10 @@ private void handleNewVirtualComputeGroup(Cloud.ClusterPB cluster, List(subComputeGroups)); - ComputeGroup.Policy policy = new ComputeGroup.Policy(); + CloudComputeGroupMeta.Policy policy = new CloudComputeGroupMeta.Policy(); policy.setActiveComputeGroup(cluster.getClusterPolicy().getActiveClusterName()); policy.setStandbyComputeGroup(cluster.getClusterPolicy().getStandbyClusterNames(0)); policy.setFailoverFailureThreshold(cluster.getClusterPolicy().getFailoverFailureThreshold()); @@ -541,7 +543,7 @@ private void handleFailedSync(Cloud.ClusterPB cluster, String subClusterName, private void removeObsoleteVirtualGroups(List virtualClusters) { List msVirtualClusters = virtualClusters.stream().map(Cloud.ClusterPB::getClusterId) .collect(Collectors.toList()); - for (ComputeGroup computeGroup : cloudSystemInfoService.getComputeGroups(true)) { + for (CloudComputeGroupMeta computeGroup : cloudSystemInfoService.getComputeGroups(true)) { // in fe mem, but not in meta server if (!msVirtualClusters.contains(computeGroup.getId())) { LOG.info("virtual compute group {} will be removed.", computeGroup.getName()); diff --git a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudTabletRebalancer.java b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudTabletRebalancer.java index ee9eb1eb568bef..6acaad24e0700a 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudTabletRebalancer.java +++ b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudTabletRebalancer.java @@ -202,7 +202,7 @@ public int hashCode() { * Get the current balance type for a compute group, falling back to global balance type if not found */ private BalanceTypeEnum getCurrentBalanceType(String clusterId) { - ComputeGroup cg = cloudSystemInfoService.getComputeGroupById(clusterId); + CloudComputeGroupMeta cg = cloudSystemInfoService.getComputeGroupById(clusterId); if (cg == null) { LOG.debug("compute group not found, use global balance type, id {}", clusterId); return globalBalanceTypeEnum; @@ -219,7 +219,7 @@ private BalanceTypeEnum getCurrentBalanceType(String clusterId) { * Get the current task timeout for a compute group, falling back to global timeout if not found */ private int getCurrentTaskTimeout(String clusterId) { - ComputeGroup cg = cloudSystemInfoService.getComputeGroupById(clusterId); + CloudComputeGroupMeta cg = cloudSystemInfoService.getComputeGroupById(clusterId); if (cg == null) { return Config.cloud_pre_heating_time_limit_sec; } @@ -233,15 +233,15 @@ private int getCurrentTaskTimeout(String clusterId) { } private boolean isComputeGroupBalanceChanged(String clusterId) { - ComputeGroup cg = cloudSystemInfoService.getComputeGroupById(clusterId); + CloudComputeGroupMeta cg = cloudSystemInfoService.getComputeGroupById(clusterId); if (cg == null) { return false; } BalanceTypeEnum computeGroupBalanceType = cg.getBalanceType(); int computeGroupTimeout = cg.getBalanceWarmUpTaskTimeout(); - return computeGroupBalanceType != ComputeGroup.DEFAULT_COMPUTE_GROUP_BALANCE_ENUM - || computeGroupTimeout != ComputeGroup.DEFAULT_BALANCE_WARM_UP_TASK_TIMEOUT; + return computeGroupBalanceType != CloudComputeGroupMeta.DEFAULT_COMPUTE_GROUP_BALANCE_ENUM + || computeGroupTimeout != CloudComputeGroupMeta.DEFAULT_BALANCE_WARM_UP_TASK_TIMEOUT; } public CloudTabletRebalancer(CloudSystemInfoService cloudSystemInfoService) { 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 13b347694c94c9..b512871d85e0a8 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 @@ -23,8 +23,8 @@ import org.apache.doris.catalog.Env; import org.apache.doris.catalog.ReplicaAllocation; import org.apache.doris.cloud.catalog.CloudColocatePlacement; +import org.apache.doris.cloud.catalog.CloudComputeGroupMeta; import org.apache.doris.cloud.catalog.CloudEnv; -import org.apache.doris.cloud.catalog.ComputeGroup; import org.apache.doris.cloud.proto.Cloud; import org.apache.doris.cloud.proto.Cloud.ClusterPB; import org.apache.doris.cloud.proto.Cloud.InstanceInfoPB; @@ -99,8 +99,8 @@ public class CloudSystemInfoService extends SystemInfoService { // clusterName -> clusterId protected Map clusterNameToId = new ConcurrentHashMap<>(); - // clusterId -> ComputeGroup - protected Map computeGroupIdToComputeGroup = new ConcurrentHashMap<>(); + // clusterId -> CloudComputeGroupMeta + protected Map computeGroupIdToComputeGroup = new ConcurrentHashMap<>(); private final Map colocatePlacementCache = new ConcurrentHashMap<>(); @@ -260,7 +260,7 @@ public void renameVirtualClusterInfoFromMapsNoLock(String clusterId, String oldC clusterNameToId.remove(oldClusterName); } - public ComputeGroup getComputeGroupByName(String computeGroupName) { + public CloudComputeGroupMeta getComputeGroupByName(String computeGroupName) { // rlock guards the compound name->id->group lookup: writers (add/remove/rename) // update both maps under wlock, and the read must observe a consistent snapshot // so callers like getPhysicalCluster don't transiently see a virtual group name @@ -355,7 +355,7 @@ public String resolveClusterIdByName(String cluster) throws ComputeGroupExceptio return getCloudClusterIdByName(cluster); } - public ComputeGroup getComputeGroupById(String computeGroupId) { + public CloudComputeGroupMeta getComputeGroupById(String computeGroupId) { try { rlock.lock(); return computeGroupIdToComputeGroup.get(computeGroupId); @@ -364,7 +364,7 @@ public ComputeGroup getComputeGroupById(String computeGroupId) { } } - public void addComputeGroup(String computeGroupId, ComputeGroup computeGroup) { + public void addComputeGroup(String computeGroupId, CloudComputeGroupMeta computeGroup) { LOG.debug("add id {} computeGroupIdToComputeGroup : {} ", computeGroupId, computeGroupIdToComputeGroup); try { wlock.lock(); @@ -376,8 +376,8 @@ public void addComputeGroup(String computeGroupId, ComputeGroup computeGroup) { } public boolean isStandByComputeGroup(String clusterName) { - List virtualGroups = getComputeGroups(true); - for (ComputeGroup vcg : virtualGroups) { + List virtualGroups = getComputeGroups(true); + for (CloudComputeGroupMeta vcg : virtualGroups) { if (vcg.getPolicy().getStandbyComputeGroup().equals(clusterName)) { return true; } @@ -385,7 +385,7 @@ public boolean isStandByComputeGroup(String clusterName) { return false; } - public List getComputeGroups(boolean virtual) { + public List getComputeGroups(boolean virtual) { LOG.debug("get virtual {} computeGroupIdToComputeGroup : {} ", virtual, computeGroupIdToComputeGroup); try { rlock.lock(); @@ -404,7 +404,7 @@ public List getComputeGroups(boolean virtual) { public String ownedByVirtualComputeGroup(String computeGroupName) { try { rlock.lock(); - for (ComputeGroup vcg : getComputeGroups(true)) { + for (CloudComputeGroupMeta vcg : getComputeGroups(true)) { if (computeGroupName.equals(vcg.getPolicy().getActiveComputeGroup())) { return vcg.getName(); } @@ -433,7 +433,7 @@ public void removeComputeGroup(String computeGroupId, String computeGroupName) { } public void renameVirtualComputeGroup(String computeGroupId, String oldComputeGroupName, - ComputeGroup newComputeGroup) { + CloudComputeGroupMeta newComputeGroup) { try { wlock.lock(); computeGroupIdToComputeGroup.put(computeGroupId, newComputeGroup); @@ -587,7 +587,8 @@ public void updateCloudClusterMapNoLock(List toAdd, List toDel clusterNameToId.put(clusterName, clusterId); // add to computeGroupIdToComputeGroup - ComputeGroup cg = new ComputeGroup(clusterId, clusterName, ComputeGroup.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta cg = new CloudComputeGroupMeta( + clusterId, clusterName, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); addComputeGroup(clusterId, cg); List be = clusterIdToBackend.get(clusterId); @@ -685,7 +686,7 @@ public synchronized void updateFrontends(List toAdd, List to } } - public static void updateFileCacheJobIds(ComputeGroup cg, List jobIds) { + public static void updateFileCacheJobIds(CloudComputeGroupMeta cg, List jobIds) { Cloud.ClusterPolicy policy = Cloud.ClusterPolicy.newBuilder() .setType(Cloud.ClusterPolicy.PolicyType.ActiveStandby) .addAllCacheWarmupJobids(jobIds).build(); @@ -727,7 +728,7 @@ enum PolicyType { } */ - private void switchActiveStandby(ComputeGroup cg, String active, String standby) { + private void switchActiveStandby(CloudComputeGroupMeta cg, String active, String standby) { Cloud.ClusterPolicy policy = cg.getPolicy().toPb().toBuilder() .clearStandbyClusterNames() .addStandbyClusterNames(active) @@ -1060,7 +1061,7 @@ public List getBackendsByClusterId(final String clusterId) { } public String getPhysicalCluster(String clusterName) { - ComputeGroup cg = getComputeGroupByName(clusterName); + CloudComputeGroupMeta cg = getComputeGroupByName(clusterName); if (cg == null) { return clusterName; } @@ -1069,11 +1070,11 @@ public String getPhysicalCluster(String clusterName) { return clusterName; } - ComputeGroup.Policy policy = cg.getPolicy(); + CloudComputeGroupMeta.Policy policy = cg.getPolicy(); // todo check policy String acgName = policy.getActiveComputeGroup(); if (acgName != null) { - ComputeGroup acg = getComputeGroupByName(acgName); + CloudComputeGroupMeta acg = getComputeGroupByName(acgName); if (acg != null) { if (isComputeGroupAvailable(acgName, policy.getUnhealthyNodeThresholdPercent())) { acg.setUnavailableSince(-1); @@ -1089,11 +1090,11 @@ public String getPhysicalCluster(String clusterName) { String scgName = policy.getStandbyComputeGroup(); if (scgName != null) { - ComputeGroup scg = getComputeGroupByName(scgName); + CloudComputeGroupMeta scg = getComputeGroupByName(scgName); if (scg != null) { if (isComputeGroupAvailable(scgName, policy.getUnhealthyNodeThresholdPercent())) { scg.setUnavailableSince(-1); - ComputeGroup acg = getComputeGroupByName(acgName); + CloudComputeGroupMeta acg = getComputeGroupByName(acgName); if (acg == null || System.currentTimeMillis() - acg.getUnavailableSince() > policy.getFailoverFailureThreshold() * Config.heartbeat_interval_second * 1000) { switchActiveStandby(cg, acgName, scgName); @@ -1378,7 +1379,7 @@ public Map> getCloudClusterIdToBackend(boolean needVirtual clusterId, computeGroupIdToComputeGroup); continue; } - ComputeGroup computeGroup = computeGroupIdToComputeGroup.get(clusterId); + CloudComputeGroupMeta computeGroup = computeGroupIdToComputeGroup.get(clusterId); if (!needVirtual && computeGroup.isVirtual()) { continue; } @@ -1421,7 +1422,7 @@ public Map getCloudClusterNameToId(boolean needVirtual) { try { for (Map.Entry nameAndId : clusterNameToId.entrySet()) { String clusterId = nameAndId.getValue(); - ComputeGroup computeGroup = computeGroupIdToComputeGroup.get(clusterId); + CloudComputeGroupMeta computeGroup = computeGroupIdToComputeGroup.get(clusterId); if (computeGroup == null) { LOG.warn("cant find clusterId {} in computeGroupIdToComputeGroup {}", clusterId, computeGroupIdToComputeGroup); diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/AlterComputeGroupCommand.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/AlterComputeGroupCommand.java index 8a295678f4826e..d5f520661b8f8e 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/AlterComputeGroupCommand.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/AlterComputeGroupCommand.java @@ -18,7 +18,7 @@ package org.apache.doris.nereids.trees.plans.commands; import org.apache.doris.catalog.Env; -import org.apache.doris.cloud.catalog.ComputeGroup; +import org.apache.doris.cloud.catalog.CloudComputeGroupMeta; import org.apache.doris.cloud.system.CloudSystemInfoService; import org.apache.doris.common.AnalysisException; import org.apache.doris.common.Config; @@ -68,7 +68,7 @@ public void validate(ConnectContext connectContext) throws UserException { CloudSystemInfoService cloudSys = ((CloudSystemInfoService) Env.getCurrentSystemInfo()); // check compute group exist - ComputeGroup cg = cloudSys.getComputeGroupByName(computeGroupName); + CloudComputeGroupMeta cg = cloudSys.getComputeGroupByName(computeGroupName); if (cg == null) { throw new AnalysisException("Compute Group " + computeGroupName + " does not exist"); } diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/ShowClustersCommand.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/ShowClustersCommand.java index a4c389941e009d..c774849ed48af4 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/ShowClustersCommand.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/ShowClustersCommand.java @@ -22,7 +22,7 @@ import org.apache.doris.catalog.Column; import org.apache.doris.catalog.Env; import org.apache.doris.catalog.ScalarType; -import org.apache.doris.cloud.catalog.ComputeGroup; +import org.apache.doris.cloud.catalog.CloudComputeGroupMeta; import org.apache.doris.cloud.qe.ComputeGroupException; import org.apache.doris.cloud.system.CloudSystemInfoService; import org.apache.doris.common.AnalysisException; @@ -96,9 +96,9 @@ public ShowResultSet doRun(ConnectContext ctx, StmtExecutor executor) throws Exc CloudSystemInfoService cloudSys = ((CloudSystemInfoService) Env.getCurrentSystemInfo()); clusterNames = cloudSys.getCloudClusterNames(); // virtual cluster info - List virtualComputeGroup = cloudSys.getComputeGroups(true); + List virtualComputeGroup = cloudSys.getComputeGroups(true); List virtualComputeGroupNames = virtualComputeGroup.stream() - .map(ComputeGroup::getName).collect(Collectors.toList()); + .map(CloudComputeGroupMeta::getName).collect(Collectors.toList()); clusterNames.addAll(virtualComputeGroupNames); @@ -112,7 +112,7 @@ public ShowResultSet doRun(ConnectContext ctx, StmtExecutor executor) throws Exc PrivPredicate.USAGE, ResourceTypeEnum.CLUSTER)) { continue; } - ComputeGroup cg = cloudSys.getComputeGroupByName(clusterName); + CloudComputeGroupMeta cg = cloudSys.getComputeGroupByName(clusterName); if (cg == null) { continue; } diff --git a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/WarmUpClusterCommand.java b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/WarmUpClusterCommand.java index 4c6d74c898f54f..fd77fb41779de4 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/WarmUpClusterCommand.java +++ b/fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/WarmUpClusterCommand.java @@ -24,8 +24,8 @@ import org.apache.doris.catalog.ScalarType; import org.apache.doris.catalog.info.TableNameInfo; import org.apache.doris.cloud.OnTablesFilter.TableFilterRule; +import org.apache.doris.cloud.catalog.CloudComputeGroupMeta; import org.apache.doris.cloud.catalog.CloudEnv; -import org.apache.doris.cloud.catalog.ComputeGroup; import org.apache.doris.cloud.system.CloudSystemInfoService; import org.apache.doris.common.AnalysisException; import org.apache.doris.common.Config; @@ -141,7 +141,7 @@ public void run(ConnectContext ctx, StmtExecutor executor) throws Exception { private void checkWarmupCgs(CloudSystemInfoService cloudSys) throws AnalysisException { if (!Strings.isNullOrEmpty(srcCluster)) { - ComputeGroup srcCg = cloudSys.getComputeGroupByName(srcCluster); + CloudComputeGroupMeta srcCg = cloudSys.getComputeGroupByName(srcCluster); if (srcCg != null && srcCg.isVirtual()) { throw new AnalysisException("The srcClusterName " + srcCluster + " is a virtual compute group, not support"); @@ -149,7 +149,7 @@ private void checkWarmupCgs(CloudSystemInfoService cloudSys) throws AnalysisExce } if (!Strings.isNullOrEmpty(dstCluster)) { - ComputeGroup dstCg = cloudSys.getComputeGroupByName(dstCluster); + CloudComputeGroupMeta dstCg = cloudSys.getComputeGroupByName(dstCluster); if (dstCg != null && dstCg.isVirtual()) { throw new AnalysisException("The dstClusterName " + dstCluster + " is a virtual compute group, not support"); diff --git a/fe/fe-core/src/test/java/org/apache/doris/cloud/WarmUpClusterOnTablesParseTest.java b/fe/fe-core/src/test/java/org/apache/doris/cloud/WarmUpClusterOnTablesParseTest.java index 8bec1d2d5eb394..ceee9533386d4a 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/cloud/WarmUpClusterOnTablesParseTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/cloud/WarmUpClusterOnTablesParseTest.java @@ -20,7 +20,7 @@ import org.apache.doris.catalog.Env; import org.apache.doris.cloud.OnTablesFilter.TableFilterRule; import org.apache.doris.cloud.OnTablesFilter.TableFilterRule.RuleType; -import org.apache.doris.cloud.catalog.ComputeGroup; +import org.apache.doris.cloud.catalog.CloudComputeGroupMeta; import org.apache.doris.cloud.system.CloudSystemInfoService; import org.apache.doris.common.AnalysisException; import org.apache.doris.common.Config; @@ -103,14 +103,14 @@ private CloudSystemInfoService buildCloudSystemInfoWithVirtualComputeGroup( private void addVirtualComputeGroup(CloudSystemInfoService cloudSys, String virtualComputeGroupName, String activeComputeGroupName, String standbyComputeGroupName) { - ComputeGroup activeComputeGroup = new ComputeGroup(activeComputeGroupName + "_id", - activeComputeGroupName, ComputeGroup.ComputeTypeEnum.COMPUTE); - ComputeGroup standbyComputeGroup = new ComputeGroup(standbyComputeGroupName + "_id", - standbyComputeGroupName, ComputeGroup.ComputeTypeEnum.COMPUTE); - ComputeGroup virtualComputeGroup = new ComputeGroup(virtualComputeGroupName + "_id", - virtualComputeGroupName, ComputeGroup.ComputeTypeEnum.VIRTUAL); + CloudComputeGroupMeta activeComputeGroup = new CloudComputeGroupMeta(activeComputeGroupName + "_id", + activeComputeGroupName, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta standbyComputeGroup = new CloudComputeGroupMeta(standbyComputeGroupName + "_id", + standbyComputeGroupName, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta virtualComputeGroup = new CloudComputeGroupMeta(virtualComputeGroupName + "_id", + virtualComputeGroupName, CloudComputeGroupMeta.ComputeTypeEnum.VIRTUAL); virtualComputeGroup.setSubComputeGroups(Arrays.asList(activeComputeGroupName, standbyComputeGroupName)); - ComputeGroup.Policy policy = new ComputeGroup.Policy(); + CloudComputeGroupMeta.Policy policy = new CloudComputeGroupMeta.Policy(); policy.setActiveComputeGroup(activeComputeGroupName); policy.setStandbyComputeGroup(standbyComputeGroupName); virtualComputeGroup.setPolicy(policy); @@ -370,7 +370,7 @@ public void testOnTablesLoadEventValidateAllowsDestinationComputeGroupOwnedByVir CloudSystemInfoService cloudSys = buildCloudSystemInfoWithVirtualComputeGroup( "vcg", "active_cg", "standby_cg"); cloudSys.addComputeGroup("outside_cg_id", - new ComputeGroup("outside_cg_id", "outside_cg", ComputeGroup.ComputeTypeEnum.COMPUTE)); + new CloudComputeGroupMeta("outside_cg_id", "outside_cg", CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE)); setField(env, Env.class, "systemInfo", cloudSys); WarmUpClusterCommand cmd = parse( "WARM UP CLUSTER standby_cg WITH CLUSTER outside_cg " @@ -398,7 +398,7 @@ public void testOnTablesLoadEventValidateAllowsSourceComputeGroupOwnedByVirtualC CloudSystemInfoService cloudSys = buildCloudSystemInfoWithVirtualComputeGroup( "vcg", "active_cg", "standby_cg"); cloudSys.addComputeGroup("outside_cg_id", - new ComputeGroup("outside_cg_id", "outside_cg", ComputeGroup.ComputeTypeEnum.COMPUTE)); + new CloudComputeGroupMeta("outside_cg_id", "outside_cg", CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE)); setField(env, Env.class, "systemInfo", cloudSys); WarmUpClusterCommand cmd = parse( "WARM UP CLUSTER outside_cg WITH CLUSTER active_cg " diff --git a/fe/fe-core/src/test/java/org/apache/doris/cloud/catalog/ComputeGroupTest.java b/fe/fe-core/src/test/java/org/apache/doris/cloud/catalog/CloudComputeGroupMetaTest.java similarity index 61% rename from fe/fe-core/src/test/java/org/apache/doris/cloud/catalog/ComputeGroupTest.java rename to fe/fe-core/src/test/java/org/apache/doris/cloud/catalog/CloudComputeGroupMetaTest.java index c61f2abefa42fc..2e4d84a4b7f0e2 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/cloud/catalog/ComputeGroupTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/cloud/catalog/CloudComputeGroupMetaTest.java @@ -26,12 +26,13 @@ import java.util.Map; -public class ComputeGroupTest { - private ComputeGroup computeGroup; +public class CloudComputeGroupMetaTest { + private CloudComputeGroupMeta computeGroup; @BeforeEach public void setUp() { - computeGroup = new ComputeGroup("test_id", "test_group", ComputeGroup.ComputeTypeEnum.COMPUTE); + computeGroup = new CloudComputeGroupMeta("test_id", "test_group", + CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); } @Test @@ -44,16 +45,16 @@ public void testCheckPropertiesWithNull() throws DdlException { public void testCheckPropertiesWithValidBalanceType() throws DdlException { // 测试有效的balance_type Map properties = Maps.newHashMap(); - properties.put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.WITHOUT_WARMUP.getValue()); + properties.put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.WITHOUT_WARMUP.getValue()); computeGroup.checkProperties(properties); - properties.put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.ASYNC_WARMUP.getValue()); + properties.put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.ASYNC_WARMUP.getValue()); computeGroup.checkProperties(properties); - properties.put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.SYNC_WARMUP.getValue()); + properties.put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.SYNC_WARMUP.getValue()); computeGroup.checkProperties(properties); - properties.put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.PEER_READ_ASYNC_WARMUP.getValue()); + properties.put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.PEER_READ_ASYNC_WARMUP.getValue()); computeGroup.checkProperties(properties); } @@ -61,7 +62,7 @@ public void testCheckPropertiesWithValidBalanceType() throws DdlException { public void testCheckPropertiesWithInvalidBalanceType() { // 测试无效的balance_type Map properties = Maps.newHashMap(); - properties.put(ComputeGroup.BALANCE_TYPE, "invalid_type"); + properties.put(CloudComputeGroupMeta.BALANCE_TYPE, "invalid_type"); Assertions.assertThrows(DdlException.class, () -> { computeGroup.checkProperties(properties); @@ -72,11 +73,11 @@ public void testCheckPropertiesWithInvalidBalanceType() { public void testCheckPropertiesWithValidTimeout() throws DdlException { // 测试有效的timeout Map properties = Maps.newHashMap(); - properties.put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.ASYNC_WARMUP.getValue()); - properties.put(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT, "300"); + properties.put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.ASYNC_WARMUP.getValue()); + properties.put(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT, "300"); computeGroup.checkProperties(properties); - properties.put(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT, "1"); + properties.put(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT, "1"); computeGroup.checkProperties(properties); } @@ -84,13 +85,13 @@ public void testCheckPropertiesWithValidTimeout() throws DdlException { public void testCheckPropertiesWithInvalidTimeout() { // 测试无效的timeout Map properties = Maps.newHashMap(); - properties.put(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT, "-1"); + properties.put(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT, "-1"); Assertions.assertThrows(DdlException.class, () -> { computeGroup.checkProperties(properties); }); - properties.put(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT, "invalid"); + properties.put(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT, "invalid"); Assertions.assertThrows(DdlException.class, () -> { computeGroup.checkProperties(properties); }); @@ -111,62 +112,62 @@ public void testCheckPropertiesWithUnsupportedProperty() { public void testModifyPropertiesWithDirectSwitch() throws DdlException { // 测试without_warmup类型,应该删除timeout Map inputProperties = Maps.newHashMap(); - inputProperties.put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.WITHOUT_WARMUP.getValue()); - inputProperties.put(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT, - String.valueOf(ComputeGroup.DEFAULT_BALANCE_WARM_UP_TASK_TIMEOUT)); + inputProperties.put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.WITHOUT_WARMUP.getValue()); + inputProperties.put(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT, + String.valueOf(CloudComputeGroupMeta.DEFAULT_BALANCE_WARM_UP_TASK_TIMEOUT)); // 先设置timeout到properties中 - computeGroup.getProperties().put(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT, - String.valueOf(ComputeGroup.DEFAULT_BALANCE_WARM_UP_TASK_TIMEOUT)); - Assertions.assertTrue(computeGroup.getProperties().containsKey(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT)); + computeGroup.getProperties().put(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT, + String.valueOf(CloudComputeGroupMeta.DEFAULT_BALANCE_WARM_UP_TASK_TIMEOUT)); + Assertions.assertTrue(computeGroup.getProperties().containsKey(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT)); computeGroup.modifyProperties(inputProperties); // 验证timeout被删除 - Assertions.assertFalse(computeGroup.getProperties().containsKey(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT)); + Assertions.assertFalse(computeGroup.getProperties().containsKey(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT)); } @Test public void testModifyPropertiesWithSyncCache() throws DdlException { // 测试sync_cache类型,应该删除timeout Map inputProperties = Maps.newHashMap(); - inputProperties.put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.SYNC_WARMUP.getValue()); - inputProperties.put(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT, - String.valueOf(ComputeGroup.DEFAULT_BALANCE_WARM_UP_TASK_TIMEOUT)); + inputProperties.put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.SYNC_WARMUP.getValue()); + inputProperties.put(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT, + String.valueOf(CloudComputeGroupMeta.DEFAULT_BALANCE_WARM_UP_TASK_TIMEOUT)); // 先设置timeout到properties中 - computeGroup.getProperties().put(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT, - String.valueOf(ComputeGroup.DEFAULT_BALANCE_WARM_UP_TASK_TIMEOUT)); - Assertions.assertTrue(computeGroup.getProperties().containsKey(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT)); + computeGroup.getProperties().put(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT, + String.valueOf(CloudComputeGroupMeta.DEFAULT_BALANCE_WARM_UP_TASK_TIMEOUT)); + Assertions.assertTrue(computeGroup.getProperties().containsKey(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT)); computeGroup.modifyProperties(inputProperties); // 验证timeout被删除 - Assertions.assertFalse(computeGroup.getProperties().containsKey(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT)); + Assertions.assertFalse(computeGroup.getProperties().containsKey(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT)); } @Test public void testCheckPropertiesWithBalanceTypeTransition() throws DdlException { - computeGroup.getProperties().put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.WITHOUT_WARMUP.getValue()); + computeGroup.getProperties().put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.WITHOUT_WARMUP.getValue()); Map inputProperties = Maps.newHashMap(); - inputProperties.put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.ASYNC_WARMUP.getValue()); + inputProperties.put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.ASYNC_WARMUP.getValue()); computeGroup.checkProperties(inputProperties); } @Test public void testCheckPropertiesWithWarmupCacheToWarmupCache() throws DdlException { - computeGroup.getProperties().put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.ASYNC_WARMUP.getValue()); + computeGroup.getProperties().put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.ASYNC_WARMUP.getValue()); Map inputProperties = Maps.newHashMap(); - inputProperties.put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.ASYNC_WARMUP.getValue()); + inputProperties.put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.ASYNC_WARMUP.getValue()); computeGroup.checkProperties(inputProperties); } @Test public void testCheckPropertiesWithDirectSwitchToDirectSwitch() throws DdlException { // 测试从direct_switch转换到direct_switch,不需要设置timeout - computeGroup.getProperties().put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.WITHOUT_WARMUP.getValue()); + computeGroup.getProperties().put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.WITHOUT_WARMUP.getValue()); Map inputProperties = Maps.newHashMap(); - inputProperties.put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.WITHOUT_WARMUP.getValue()); + inputProperties.put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.WITHOUT_WARMUP.getValue()); computeGroup.checkProperties(inputProperties); } @@ -174,34 +175,34 @@ public void testCheckPropertiesWithDirectSwitchToDirectSwitch() throws DdlExcept public void testModifyPropertiesWithWarmupCacheAndExistingTimeout() throws DdlException { // 测试async_warmup类型,已存在timeout,不应该添加默认值 Map inputProperties = Maps.newHashMap(); - inputProperties.put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.ASYNC_WARMUP.getValue()); - inputProperties.put(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT, "600"); + inputProperties.put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.ASYNC_WARMUP.getValue()); + inputProperties.put(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT, "600"); // 先设置timeout到properties中 - computeGroup.getProperties().put(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT, "500"); - String originalTimeout = computeGroup.getProperties().get(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT); + computeGroup.getProperties().put(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT, "500"); + String originalTimeout = computeGroup.getProperties().get(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT); computeGroup.modifyProperties(inputProperties); // 验证timeout没有被修改, 这里的意思是用户已经设置过timeout了,就不应该被覆盖 - Assertions.assertEquals(originalTimeout, computeGroup.getProperties().get(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT)); + Assertions.assertEquals(originalTimeout, computeGroup.getProperties().get(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT)); } @Test public void testModifyPropertiesWithWarmupCacheAndNoTimeout() throws DdlException { // 测试async_warmup类型,不存在timeout,应该添加默认值 Map inputProperties = Maps.newHashMap(); - inputProperties.put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.ASYNC_WARMUP.getValue()); + inputProperties.put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.ASYNC_WARMUP.getValue()); // 确保properties中没有timeout - computeGroup.getProperties().remove(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT); - Assertions.assertFalse(computeGroup.getProperties().containsKey(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT)); + computeGroup.getProperties().remove(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT); + Assertions.assertFalse(computeGroup.getProperties().containsKey(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT)); computeGroup.modifyProperties(inputProperties); // 验证默认值被添加 - Assertions.assertTrue(computeGroup.getProperties().containsKey(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT)); - Assertions.assertEquals(String.valueOf(ComputeGroup.DEFAULT_BALANCE_WARM_UP_TASK_TIMEOUT), - computeGroup.getProperties().get(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT)); + Assertions.assertTrue(computeGroup.getProperties().containsKey(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT)); + Assertions.assertEquals(String.valueOf(CloudComputeGroupMeta.DEFAULT_BALANCE_WARM_UP_TASK_TIMEOUT), + computeGroup.getProperties().get(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT)); } @Test @@ -233,56 +234,56 @@ public void testModifyPropertiesWithEmptyInput() throws DdlException { @Test public void testCheckPropertiesWithSyncCacheToWarmupCache() throws DdlException { - computeGroup.getProperties().put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.SYNC_WARMUP.getValue()); + computeGroup.getProperties().put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.SYNC_WARMUP.getValue()); Map inputProperties = Maps.newHashMap(); - inputProperties.put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.ASYNC_WARMUP.getValue()); + inputProperties.put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.ASYNC_WARMUP.getValue()); computeGroup.checkProperties(inputProperties); } @Test public void testValidateTimeoutRestrictionWithNoCurrentBalanceType() throws DdlException { // 测试当前没有设置balance_type的情况 - computeGroup.getProperties().remove(ComputeGroup.BALANCE_TYPE); + computeGroup.getProperties().remove(CloudComputeGroupMeta.BALANCE_TYPE); Map inputProperties = Maps.newHashMap(); - inputProperties.put(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT, "500"); + inputProperties.put(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT, "500"); computeGroup.checkProperties(inputProperties); } @Test public void testValidateTimeoutRestrictionWithCurrentWarmupCache() throws DdlException { // 测试当前balance_type是warmup_cache的情况 - computeGroup.getProperties().put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.ASYNC_WARMUP.getValue()); + computeGroup.getProperties().put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.ASYNC_WARMUP.getValue()); Map inputProperties = Maps.newHashMap(); - inputProperties.put(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT, "500"); + inputProperties.put(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT, "500"); computeGroup.checkProperties(inputProperties); } @Test public void testValidateTimeoutRestrictionWithDirectSwitchToWarmupCache() throws DdlException { // 测试从direct_switch转换到warmup_cache并设置timeout - computeGroup.getProperties().put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.WITHOUT_WARMUP.getValue()); + computeGroup.getProperties().put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.WITHOUT_WARMUP.getValue()); Map inputProperties = Maps.newHashMap(); - inputProperties.put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.ASYNC_WARMUP.getValue()); - inputProperties.put(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT, "500"); + inputProperties.put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.ASYNC_WARMUP.getValue()); + inputProperties.put(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT, "500"); computeGroup.checkProperties(inputProperties); } @Test public void testValidateTimeoutRestrictionWithSyncCacheToWarmupCacheWithTimeout() throws DdlException { // 测试从sync_cache转换到warmup_cache并设置timeout - computeGroup.getProperties().put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.SYNC_WARMUP.getValue()); + computeGroup.getProperties().put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.SYNC_WARMUP.getValue()); Map inputProperties = Maps.newHashMap(); - inputProperties.put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.ASYNC_WARMUP.getValue()); - inputProperties.put(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT, "500"); + inputProperties.put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.ASYNC_WARMUP.getValue()); + inputProperties.put(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT, "500"); computeGroup.checkProperties(inputProperties); } @Test public void testValidateTimeoutRestrictionWithDirectSwitchAndOnlyTimeout() { // 测试当前是direct_switch,仅设置timeout应该失败 - computeGroup.getProperties().put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.WITHOUT_WARMUP.getValue()); + computeGroup.getProperties().put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.WITHOUT_WARMUP.getValue()); Map inputProperties = Maps.newHashMap(); - inputProperties.put(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT, "500"); + inputProperties.put(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT, "500"); Assertions.assertThrows(DdlException.class, () -> { computeGroup.checkProperties(inputProperties); @@ -292,9 +293,9 @@ public void testValidateTimeoutRestrictionWithDirectSwitchAndOnlyTimeout() { @Test public void testValidateTimeoutRestrictionWithSyncCacheAndOnlyTimeout() { // 测试当前是sync_cache,仅设置timeout应该失败 - computeGroup.getProperties().put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.SYNC_WARMUP.getValue()); + computeGroup.getProperties().put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.SYNC_WARMUP.getValue()); Map inputProperties = Maps.newHashMap(); - inputProperties.put(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT, "500"); + inputProperties.put(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT, "500"); Assertions.assertThrows(DdlException.class, () -> { computeGroup.checkProperties(inputProperties); @@ -304,10 +305,10 @@ public void testValidateTimeoutRestrictionWithSyncCacheAndOnlyTimeout() { @Test public void testValidateTimeoutRestrictionWithDirectSwitchToSyncCacheAndTimeout() { // 测试从direct_switch转换到sync_cache并设置timeout应该失败 - computeGroup.getProperties().put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.WITHOUT_WARMUP.getValue()); + computeGroup.getProperties().put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.WITHOUT_WARMUP.getValue()); Map inputProperties = Maps.newHashMap(); - inputProperties.put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.SYNC_WARMUP.getValue()); - inputProperties.put(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT, "500"); + inputProperties.put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.SYNC_WARMUP.getValue()); + inputProperties.put(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT, "500"); Assertions.assertThrows(DdlException.class, () -> { computeGroup.checkProperties(inputProperties); @@ -317,10 +318,10 @@ public void testValidateTimeoutRestrictionWithDirectSwitchToSyncCacheAndTimeout( @Test public void testValidateTimeoutRestrictionWithDirectSwitchToSameAndTimeout() { // 测试从direct_switch转换到direct_switch并设置timeout应该失败 - computeGroup.getProperties().put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.WITHOUT_WARMUP.getValue()); + computeGroup.getProperties().put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.WITHOUT_WARMUP.getValue()); Map inputProperties = Maps.newHashMap(); - inputProperties.put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.WITHOUT_WARMUP.getValue()); - inputProperties.put(ComputeGroup.BALANCE_WARM_UP_TASK_TIMEOUT, "500"); + inputProperties.put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.WITHOUT_WARMUP.getValue()); + inputProperties.put(CloudComputeGroupMeta.BALANCE_WARM_UP_TASK_TIMEOUT, "500"); Assertions.assertThrows(DdlException.class, () -> { computeGroup.checkProperties(inputProperties); @@ -330,7 +331,7 @@ public void testValidateTimeoutRestrictionWithDirectSwitchToSameAndTimeout() { @Test public void testValidateTimeoutRestrictionWithNoInputTimeout() throws DdlException { // 测试输入中没有timeout的情况 - computeGroup.getProperties().put(ComputeGroup.BALANCE_TYPE, BalanceTypeEnum.WITHOUT_WARMUP.getValue()); + computeGroup.getProperties().put(CloudComputeGroupMeta.BALANCE_TYPE, BalanceTypeEnum.WITHOUT_WARMUP.getValue()); Map inputProperties = Maps.newHashMap(); inputProperties.put("other_property", "value"); Assertions.assertThrows(DdlException.class, () -> { 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..fe4060d2ab44e0 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 @@ -102,7 +102,7 @@ public void testSyncInstanceCreatesVirtualComputeGroup() { new CloudInstanceStatusChecker(cloudSystemInfoService).runAfterCatalogReady(); - ComputeGroup virtualComputeGroup = cloudSystemInfoService.getComputeGroupById("vcg_id"); + CloudComputeGroupMeta virtualComputeGroup = cloudSystemInfoService.getComputeGroupById("vcg_id"); Assertions.assertNotNull(virtualComputeGroup); Assertions.assertTrue(virtualComputeGroup.isVirtual()); Assertions.assertEquals("vcg", virtualComputeGroup.getName()); @@ -129,17 +129,17 @@ 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(CloudComputeGroupMeta.class), Mockito.anyList())).thenAnswer(invocation -> null); new CloudInstanceStatusChecker(cloudSystemInfoService).runAfterCatalogReady(); mockedCloudSystemInfoService.verify(() -> CloudSystemInfoService.updateFileCacheJobIds( - Mockito.any(ComputeGroup.class), Mockito.anyList())); + Mockito.any(CloudComputeGroupMeta.class), Mockito.anyList())); } finally { logger.removeAppender(appender); appender.stop(); } - ComputeGroup virtualComputeGroup = cloudSystemInfoService.getComputeGroupById("vcg_id"); + CloudComputeGroupMeta virtualComputeGroup = cloudSystemInfoService.getComputeGroupById("vcg_id"); Assertions.assertNotNull(virtualComputeGroup); Assertions.assertTrue(virtualComputeGroup.isVirtual()); Assertions.assertFalse(virtualComputeGroup.isNeedRebuildFileCache()); @@ -174,7 +174,7 @@ public void testSyncInstanceCreatesVirtualComputeGroupAndCancelsTableLevelLoadEv private void addComputeGroup(String computeGroupId, String computeGroupName) { cloudSystemInfoService.addComputeGroup(computeGroupId, - new ComputeGroup(computeGroupId, computeGroupName, ComputeGroup.ComputeTypeEnum.COMPUTE)); + new CloudComputeGroupMeta(computeGroupId, computeGroupName, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE)); } private Cloud.GetInstanceResponse instanceResponseWithVirtualComputeGroup() { diff --git a/fe/fe-core/src/test/java/org/apache/doris/cloud/system/CloudSystemInfoServiceTest.java b/fe/fe-core/src/test/java/org/apache/doris/cloud/system/CloudSystemInfoServiceTest.java index 83f62b351eec8d..28a4decd6be456 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/cloud/system/CloudSystemInfoServiceTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/cloud/system/CloudSystemInfoServiceTest.java @@ -19,8 +19,8 @@ import org.apache.doris.analysis.UserIdentity; import org.apache.doris.catalog.Env; +import org.apache.doris.cloud.catalog.CloudComputeGroupMeta; import org.apache.doris.cloud.catalog.CloudEnv; -import org.apache.doris.cloud.catalog.ComputeGroup; import org.apache.doris.cloud.proto.Cloud; import org.apache.doris.cloud.rpc.MetaServiceProxy; import org.apache.doris.common.Config; @@ -71,7 +71,7 @@ public void testGetPhysicalClusterPhysicalCluster() { //public void testGetPhysicalClusterEmptyVirtualCluster() { // infoService = new CloudSystemInfoService(); // String vcgName = "v_cluster_1"; - // ComputeGroup vcg = new ComputeGroup("id1", vcgName, ComputeGroup.ComputeTypeEnum.VIRTUAL); + // CloudComputeGroupMeta vcg = new CloudComputeGroupMeta("id1", vcgName, CloudComputeGroupMeta.ComputeTypeEnum.VIRTUAL); // infoService.addComputeGroup(vcgName, vcg); // String res = infoService.getPhysicalCluster(vcgName); @@ -86,14 +86,14 @@ public void testGetPhysicalClusterEmptyCluster() { String pcgName1 = "p_cluster_1"; String pcgName2 = "p_cluster_2"; - ComputeGroup vcg = new ComputeGroup("id1", vcgName, ComputeGroup.ComputeTypeEnum.VIRTUAL); - ComputeGroup.Policy policy = new ComputeGroup.Policy(); + CloudComputeGroupMeta vcg = new CloudComputeGroupMeta("id1", vcgName, CloudComputeGroupMeta.ComputeTypeEnum.VIRTUAL); + CloudComputeGroupMeta.Policy policy = new CloudComputeGroupMeta.Policy(); policy.setActiveComputeGroup(pcgName1); policy.setStandbyComputeGroup(pcgName2); vcg.setPolicy(policy); - ComputeGroup pcg1 = new ComputeGroup("id2", pcgName1, ComputeGroup.ComputeTypeEnum.COMPUTE); - ComputeGroup pcg2 = new ComputeGroup("id3", pcgName2, ComputeGroup.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta pcg1 = new CloudComputeGroupMeta("id2", pcgName1, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta pcg2 = new CloudComputeGroupMeta("id3", pcgName2, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); infoService.addComputeGroup(vcgName, vcg); infoService.addComputeGroup(pcgName1, pcg1); infoService.addComputeGroup(pcgName2, pcg2); @@ -111,14 +111,14 @@ public void testGetPhysicalClusterStandbyAvailable() { String pcgName1 = "p_cluster_1"; String pcgName2 = "p_cluster_2"; - ComputeGroup vcg = new ComputeGroup("id1", vcgName, ComputeGroup.ComputeTypeEnum.VIRTUAL); - ComputeGroup.Policy policy = new ComputeGroup.Policy(); + CloudComputeGroupMeta vcg = new CloudComputeGroupMeta("id1", vcgName, CloudComputeGroupMeta.ComputeTypeEnum.VIRTUAL); + CloudComputeGroupMeta.Policy policy = new CloudComputeGroupMeta.Policy(); policy.setActiveComputeGroup(pcgName1); policy.setStandbyComputeGroup(pcgName2); vcg.setPolicy(policy); - ComputeGroup pcg1 = new ComputeGroup("id2", pcgName1, ComputeGroup.ComputeTypeEnum.COMPUTE); - ComputeGroup pcg2 = new ComputeGroup("id3", pcgName2, ComputeGroup.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta pcg1 = new CloudComputeGroupMeta("id2", pcgName1, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta pcg2 = new CloudComputeGroupMeta("id3", pcgName2, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); infoService.addComputeGroup(vcgName, vcg); infoService.addComputeGroup(pcgName1, pcg1); infoService.addComputeGroup(pcgName2, pcg2); @@ -148,14 +148,14 @@ public void testGetPhysicalClusterActiveAvailable() { String pcgName1 = "p_cluster_1"; String pcgName2 = "p_cluster_2"; - ComputeGroup vcg = new ComputeGroup("id1", vcgName, ComputeGroup.ComputeTypeEnum.VIRTUAL); - ComputeGroup.Policy policy = new ComputeGroup.Policy(); + CloudComputeGroupMeta vcg = new CloudComputeGroupMeta("id1", vcgName, CloudComputeGroupMeta.ComputeTypeEnum.VIRTUAL); + CloudComputeGroupMeta.Policy policy = new CloudComputeGroupMeta.Policy(); policy.setActiveComputeGroup(pcgName1); policy.setStandbyComputeGroup(pcgName2); vcg.setPolicy(policy); - ComputeGroup pcg1 = new ComputeGroup("id2", pcgName1, ComputeGroup.ComputeTypeEnum.COMPUTE); - ComputeGroup pcg2 = new ComputeGroup("id3", pcgName2, ComputeGroup.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta pcg1 = new CloudComputeGroupMeta("id2", pcgName1, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta pcg2 = new CloudComputeGroupMeta("id3", pcgName2, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); infoService.addComputeGroup(vcgName, vcg); infoService.addComputeGroup(pcgName1, pcg1); infoService.addComputeGroup(pcgName2, pcg2); @@ -185,14 +185,14 @@ public void testGetPhysicalClusterActive3AliveBe() { String pcgName1 = "p_cluster_1"; String pcgName2 = "p_cluster_2"; - ComputeGroup vcg = new ComputeGroup("id1", vcgName, ComputeGroup.ComputeTypeEnum.VIRTUAL); - ComputeGroup.Policy policy = new ComputeGroup.Policy(); + CloudComputeGroupMeta vcg = new CloudComputeGroupMeta("id1", vcgName, CloudComputeGroupMeta.ComputeTypeEnum.VIRTUAL); + CloudComputeGroupMeta.Policy policy = new CloudComputeGroupMeta.Policy(); policy.setActiveComputeGroup(pcgName1); policy.setStandbyComputeGroup(pcgName2); vcg.setPolicy(policy); - ComputeGroup pcg1 = new ComputeGroup("id2", pcgName1, ComputeGroup.ComputeTypeEnum.COMPUTE); - ComputeGroup pcg2 = new ComputeGroup("id3", pcgName2, ComputeGroup.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta pcg1 = new CloudComputeGroupMeta("id2", pcgName1, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta pcg2 = new CloudComputeGroupMeta("id3", pcgName2, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); infoService.addComputeGroup(vcgName, vcg); infoService.addComputeGroup(pcgName1, pcg1); infoService.addComputeGroup(pcgName2, pcg2); @@ -234,14 +234,14 @@ public void testGetPhysicalClusterStandby3AliveBe() { String pcgName1 = "p_cluster_1"; String pcgName2 = "p_cluster_2"; - ComputeGroup vcg = new ComputeGroup("id1", vcgName, ComputeGroup.ComputeTypeEnum.VIRTUAL); - ComputeGroup.Policy policy = new ComputeGroup.Policy(); + CloudComputeGroupMeta vcg = new CloudComputeGroupMeta("id1", vcgName, CloudComputeGroupMeta.ComputeTypeEnum.VIRTUAL); + CloudComputeGroupMeta.Policy policy = new CloudComputeGroupMeta.Policy(); policy.setActiveComputeGroup(pcgName1); policy.setStandbyComputeGroup(pcgName2); vcg.setPolicy(policy); - ComputeGroup pcg1 = new ComputeGroup("id2", pcgName1, ComputeGroup.ComputeTypeEnum.COMPUTE); - ComputeGroup pcg2 = new ComputeGroup("id3", pcgName2, ComputeGroup.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta pcg1 = new CloudComputeGroupMeta("id2", pcgName1, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta pcg2 = new CloudComputeGroupMeta("id3", pcgName2, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); infoService.addComputeGroup(vcgName, vcg); infoService.addComputeGroup(pcgName1, pcg1); infoService.addComputeGroup(pcgName2, pcg2); @@ -283,14 +283,14 @@ public void testGetPhysicalClusterSwitchActiveStandbyMetric() throws Exception { String pcgName1 = "p_cluster_1"; String pcgName2 = "p_cluster_2"; - ComputeGroup vcg = new ComputeGroup(vcgId, vcgName, ComputeGroup.ComputeTypeEnum.VIRTUAL); - ComputeGroup.Policy policy = new ComputeGroup.Policy(); + CloudComputeGroupMeta vcg = new CloudComputeGroupMeta(vcgId, vcgName, CloudComputeGroupMeta.ComputeTypeEnum.VIRTUAL); + CloudComputeGroupMeta.Policy policy = new CloudComputeGroupMeta.Policy(); policy.setActiveComputeGroup(pcgName1); policy.setStandbyComputeGroup(pcgName2); policy.setUnhealthyNodeThresholdPercent(100); vcg.setPolicy(policy); - ComputeGroup pcg2 = new ComputeGroup("id3", pcgName2, ComputeGroup.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta pcg2 = new CloudComputeGroupMeta("id3", pcgName2, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); infoService.addComputeGroup(vcgId, vcg); infoService.clusterNameToId.put(pcgName1, "id2"); infoService.addComputeGroup("id3", pcg2); @@ -342,14 +342,14 @@ public void testGetPhysicalClusterActive1AliveBe2DeadBe() { String pcgName1 = "p_cluster_1"; String pcgName2 = "p_cluster_2"; - ComputeGroup vcg = new ComputeGroup("id1", vcgName, ComputeGroup.ComputeTypeEnum.VIRTUAL); - ComputeGroup.Policy policy = new ComputeGroup.Policy(); + CloudComputeGroupMeta vcg = new CloudComputeGroupMeta("id1", vcgName, CloudComputeGroupMeta.ComputeTypeEnum.VIRTUAL); + CloudComputeGroupMeta.Policy policy = new CloudComputeGroupMeta.Policy(); policy.setActiveComputeGroup(pcgName1); policy.setStandbyComputeGroup(pcgName2); vcg.setPolicy(policy); - ComputeGroup pcg1 = new ComputeGroup("id2", pcgName1, ComputeGroup.ComputeTypeEnum.COMPUTE); - ComputeGroup pcg2 = new ComputeGroup("id3", pcgName2, ComputeGroup.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta pcg1 = new CloudComputeGroupMeta("id2", pcgName1, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta pcg2 = new CloudComputeGroupMeta("id3", pcgName2, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); infoService.addComputeGroup(vcgName, vcg); infoService.addComputeGroup(pcgName1, pcg1); infoService.addComputeGroup(pcgName2, pcg2); @@ -395,15 +395,15 @@ public void testIsStandByComputeGroup() { String pcgName2 = "p_cluster_2"; String pcgName3 = "p_cluster_3"; - ComputeGroup vcg = new ComputeGroup("id1", vcgName, ComputeGroup.ComputeTypeEnum.VIRTUAL); - ComputeGroup.Policy policy = new ComputeGroup.Policy(); + CloudComputeGroupMeta vcg = new CloudComputeGroupMeta("id1", vcgName, CloudComputeGroupMeta.ComputeTypeEnum.VIRTUAL); + CloudComputeGroupMeta.Policy policy = new CloudComputeGroupMeta.Policy(); policy.setActiveComputeGroup(pcgName1); policy.setStandbyComputeGroup(pcgName2); vcg.setPolicy(policy); - ComputeGroup pcg1 = new ComputeGroup("id2", pcgName1, ComputeGroup.ComputeTypeEnum.COMPUTE); - ComputeGroup pcg2 = new ComputeGroup("id3", pcgName2, ComputeGroup.ComputeTypeEnum.COMPUTE); - ComputeGroup pcg3 = new ComputeGroup("id4", pcgName2, ComputeGroup.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta pcg1 = new CloudComputeGroupMeta("id2", pcgName1, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta pcg2 = new CloudComputeGroupMeta("id3", pcgName2, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta pcg3 = new CloudComputeGroupMeta("id4", pcgName2, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); infoService.addComputeGroup(vcgName, vcg); infoService.addComputeGroup(pcgName1, pcg1); infoService.addComputeGroup(pcgName2, pcg2); @@ -427,7 +427,7 @@ public void testGetMinPipelineExecutorSizeWithEmptyCluster() { String clusterId = "test_cluster_id"; // Mock an empty cluster (no backends) - ComputeGroup cg = new ComputeGroup(clusterId, clusterName, ComputeGroup.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta cg = new CloudComputeGroupMeta(clusterId, clusterName, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); infoService.addComputeGroup(clusterId, cg); // Set ConnectContext to select the cluster @@ -449,7 +449,7 @@ public void testGetMinPipelineExecutorSizeWithSingleBackend() { String clusterId = "test_cluster_id"; // Setup cluster - ComputeGroup cg = new ComputeGroup(clusterId, clusterName, ComputeGroup.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta cg = new CloudComputeGroupMeta(clusterId, clusterName, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); infoService.addComputeGroup(clusterId, cg); // Add a backend with pipeline executor size = 8 @@ -483,7 +483,7 @@ public void testGetMinPipelineExecutorSizeWithMultipleBackends() { String clusterId = "test_cluster_id"; // Setup cluster - ComputeGroup cg = new ComputeGroup(clusterId, clusterName, ComputeGroup.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta cg = new CloudComputeGroupMeta(clusterId, clusterName, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); infoService.addComputeGroup(clusterId, cg); // Add multiple backends with different pipeline executor sizes @@ -534,7 +534,7 @@ public void testGetMinPipelineExecutorSizeWithZeroSizeBackends() { String clusterId = "test_cluster_id"; // Setup cluster - ComputeGroup cg = new ComputeGroup(clusterId, clusterName, ComputeGroup.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta cg = new CloudComputeGroupMeta(clusterId, clusterName, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); infoService.addComputeGroup(clusterId, cg); // Add backends with zero and positive pipeline executor sizes @@ -585,7 +585,7 @@ public void testGetMinPipelineExecutorSizeWithAllZeroSizeBackends() { String clusterId = "test_cluster_id"; // Setup cluster - ComputeGroup cg = new ComputeGroup(clusterId, clusterName, ComputeGroup.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta cg = new CloudComputeGroupMeta(clusterId, clusterName, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); infoService.addComputeGroup(clusterId, cg); // Add backends with only zero or negative pipeline executor sizes @@ -645,7 +645,7 @@ public void testGetMinPipelineExecutorSizeWithMixedValidInvalidBackends() { String clusterId = "mixed_cluster_id"; // Setup cluster - ComputeGroup cg = new ComputeGroup(clusterId, clusterName, ComputeGroup.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta cg = new CloudComputeGroupMeta(clusterId, clusterName, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); infoService.addComputeGroup(clusterId, cg); // Add backends with mixed valid and invalid pipeline executor sizes @@ -708,7 +708,7 @@ public void testGetMinPipelineExecutorSizeWithLargeValues() { String clusterId = "large_cluster_id"; // Setup cluster - ComputeGroup cg = new ComputeGroup(clusterId, clusterName, ComputeGroup.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta cg = new CloudComputeGroupMeta(clusterId, clusterName, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); infoService.addComputeGroup(clusterId, cg); // Add backends with large pipeline executor sizes @@ -759,7 +759,7 @@ public void testGetMinPipelineExecutorSizeConsistency() { String clusterId = "consistency_cluster_id"; // Setup cluster - ComputeGroup cg = new ComputeGroup(clusterId, clusterName, ComputeGroup.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta cg = new CloudComputeGroupMeta(clusterId, clusterName, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); infoService.addComputeGroup(clusterId, cg); // Add backends with same pipeline executor sizes @@ -800,11 +800,11 @@ public void testGetMinPipelineExecutorSizeWithMultipleComputeGroups() { String cluster2Id = "cluster2_id"; // Setup cluster1 - ComputeGroup cg1 = new ComputeGroup(cluster1Id, cluster1Name, ComputeGroup.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta cg1 = new CloudComputeGroupMeta(cluster1Id, cluster1Name, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); infoService.addComputeGroup(cluster1Id, cg1); // Setup cluster2 - ComputeGroup cg2 = new ComputeGroup(cluster2Id, cluster2Name, ComputeGroup.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta cg2 = new CloudComputeGroupMeta(cluster2Id, cluster2Name, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); infoService.addComputeGroup(cluster2Id, cg2); // Add backends to cluster1 with smaller pipeline executor sizes @@ -872,20 +872,20 @@ public void testGetMinPipelineExecutorSizeWithVirtualComputeGroup() { String otherClusterId = "other_cluster_id"; // Setup virtual cluster - ComputeGroup virtualCg = new ComputeGroup(virtualClusterId, virtualClusterName, - ComputeGroup.ComputeTypeEnum.VIRTUAL); - ComputeGroup.Policy policy = new ComputeGroup.Policy(); + CloudComputeGroupMeta virtualCg = new CloudComputeGroupMeta(virtualClusterId, virtualClusterName, + CloudComputeGroupMeta.ComputeTypeEnum.VIRTUAL); + CloudComputeGroupMeta.Policy policy = new CloudComputeGroupMeta.Policy(); policy.setActiveComputeGroup(physicalClusterName); virtualCg.setPolicy(policy); infoService.addComputeGroup(virtualClusterId, virtualCg); // Setup physical cluster - ComputeGroup physicalCg = new ComputeGroup(physicalClusterId, physicalClusterName, - ComputeGroup.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta physicalCg = new CloudComputeGroupMeta(physicalClusterId, physicalClusterName, + CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); infoService.addComputeGroup(physicalClusterId, physicalCg); // Setup other cluster - ComputeGroup otherCg = new ComputeGroup(otherClusterId, otherClusterName, ComputeGroup.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta otherCg = new CloudComputeGroupMeta(otherClusterId, otherClusterName, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); infoService.addComputeGroup(otherClusterId, otherCg); // Add backends to physical cluster @@ -972,11 +972,11 @@ public void testGetMinPipelineExecutorSizeWithConnectContext() { String cluster2Id = "ctx_cluster2_id"; // Setup cluster1 - ComputeGroup cg1 = new ComputeGroup(cluster1Id, cluster1Name, ComputeGroup.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta cg1 = new CloudComputeGroupMeta(cluster1Id, cluster1Name, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); infoService.addComputeGroup(cluster1Id, cg1); // Setup cluster2 - ComputeGroup cg2 = new ComputeGroup(cluster2Id, cluster2Name, ComputeGroup.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta cg2 = new CloudComputeGroupMeta(cluster2Id, cluster2Name, CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); infoService.addComputeGroup(cluster2Id, cg2); // Add backends to cluster1 with smaller pipeline executor sizes diff --git a/fe/fe-core/src/test/java/org/apache/doris/mysql/privilege/CloudAuthTest.java b/fe/fe-core/src/test/java/org/apache/doris/mysql/privilege/CloudAuthTest.java index 31880402db5918..1c8d739aeac2d1 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/mysql/privilege/CloudAuthTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/mysql/privilege/CloudAuthTest.java @@ -24,8 +24,8 @@ import org.apache.doris.catalog.AccessPrivilege; import org.apache.doris.catalog.AccessPrivilegeWithCols; import org.apache.doris.catalog.Env; +import org.apache.doris.cloud.catalog.CloudComputeGroupMeta; import org.apache.doris.cloud.catalog.CloudEnv; -import org.apache.doris.cloud.catalog.ComputeGroup; import org.apache.doris.cloud.system.CloudSystemInfoService; import org.apache.doris.common.AnalysisException; import org.apache.doris.nereids.parser.NereidsParser; @@ -336,12 +336,12 @@ public void testVirtualComputeGroup() throws Exception { Assert.assertTrue(accessManager.checkCloudPriv(new UserIdentity("testUser", "%"), "vcg", PrivPredicate.USAGE, ResourceTypeEnum.CLUSTER)); // create vcg, sub cg(cg1, cg2), add to systemInfoService - ComputeGroup vcg = new ComputeGroup("vcg_id", "vcg", ComputeGroup.ComputeTypeEnum.VIRTUAL); + CloudComputeGroupMeta vcg = new CloudComputeGroupMeta("vcg_id", "vcg", CloudComputeGroupMeta.ComputeTypeEnum.VIRTUAL); vcg.setSubComputeGroups(Lists.newArrayList("cg2", "cg1")); systemInfoService.addComputeGroup("vcg_id", vcg); - ComputeGroup cg = new ComputeGroup("vcg_id", "vcg", ComputeGroup.ComputeTypeEnum.COMPUTE); + CloudComputeGroupMeta cg = new CloudComputeGroupMeta("vcg_id", "vcg", CloudComputeGroupMeta.ComputeTypeEnum.COMPUTE); systemInfoService.addComputeGroup("cg", cg); - ComputeGroup.Policy policy = new ComputeGroup.Policy(); + CloudComputeGroupMeta.Policy policy = new CloudComputeGroupMeta.Policy(); policy.setActiveComputeGroup("cg1"); policy.setStandbyComputeGroup("cg2"); vcg.setPolicy(policy); diff --git a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/AlterComputeGroupCommandTest.java b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/AlterComputeGroupCommandTest.java index ce11d4ce7f7723..dc5c595b77e88d 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/AlterComputeGroupCommandTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/nereids/trees/plans/commands/AlterComputeGroupCommandTest.java @@ -18,7 +18,7 @@ package org.apache.doris.nereids.trees.plans.commands; import org.apache.doris.catalog.Env; -import org.apache.doris.cloud.catalog.ComputeGroup; +import org.apache.doris.cloud.catalog.CloudComputeGroupMeta; import org.apache.doris.cloud.system.CloudSystemInfoService; import org.apache.doris.common.Config; import org.apache.doris.common.DdlException; @@ -113,7 +113,7 @@ public void testValidateVirtualComputeGroup() throws Exception { Config.deploy_mode = "cloud"; CloudSystemInfoService mockCloudSIS = Mockito.mock(CloudSystemInfoService.class); - ComputeGroup mockComputeGroup = Mockito.mock(ComputeGroup.class); + CloudComputeGroupMeta mockComputeGroup = Mockito.mock(CloudComputeGroupMeta.class); Mockito.doReturn(mockComputeGroup).when(mockCloudSIS).getComputeGroupByName("virtual_group"); Mockito.doReturn(true).when(mockComputeGroup).isVirtual(); Deencapsulation.setField(env, "systemInfo", mockCloudSIS); @@ -130,7 +130,7 @@ public void testValidateInvalidProperties() throws Exception { Config.deploy_mode = "cloud"; CloudSystemInfoService mockCloudSIS = Mockito.mock(CloudSystemInfoService.class); - ComputeGroup mockComputeGroup = Mockito.mock(ComputeGroup.class); + CloudComputeGroupMeta mockComputeGroup = Mockito.mock(CloudComputeGroupMeta.class); Mockito.doReturn(mockComputeGroup).when(mockCloudSIS).getComputeGroupByName("test_group"); Mockito.doReturn(false).when(mockComputeGroup).isVirtual(); Mockito.doThrow(new DdlException("Invalid property")).when(mockComputeGroup) @@ -149,7 +149,7 @@ public void testValidateSuccess() throws Exception { Config.deploy_mode = "cloud"; CloudSystemInfoService mockCloudSIS = Mockito.mock(CloudSystemInfoService.class); - ComputeGroup mockComputeGroup = Mockito.mock(ComputeGroup.class); + CloudComputeGroupMeta mockComputeGroup = Mockito.mock(CloudComputeGroupMeta.class); Mockito.doReturn(mockComputeGroup).when(mockCloudSIS).getComputeGroupByName("test_group"); Mockito.doReturn(false).when(mockComputeGroup).isVirtual(); Mockito.doNothing().when(mockComputeGroup).checkProperties(Mockito.anyMap());