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 4b3b7cf79b735d..8c355afcb7fb80 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(); for (String jobId : jobIds) { try { @@ -219,7 +219,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) { @@ -267,7 +267,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); @@ -385,7 +387,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(); @@ -479,10 +481,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()); @@ -530,7 +532,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 69df53fcda94e1..c4d7617acab6d9 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 @@ -199,7 +199,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; @@ -216,7 +216,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; } @@ -230,15 +230,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 c334db2b824c0a..4ef5098a6d50c3 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 @@ -25,8 +25,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; @@ -101,8 +101,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<>(); @@ -262,7 +262,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 @@ -357,7 +357,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); @@ -366,7 +366,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(); @@ -378,8 +378,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; } @@ -387,7 +387,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(); @@ -406,7 +406,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(); } @@ -435,7 +435,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); @@ -589,7 +589,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); @@ -687,7 +688,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(); @@ -729,7 +730,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) @@ -1072,7 +1073,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; } @@ -1081,11 +1082,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); @@ -1101,11 +1102,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); @@ -1387,7 +1388,7 @@ public Map> getCloudClusterIdToBackend(boolean needVirtual clusterId, computeGroupIdToComputeGroup); continue; } - ComputeGroup computeGroup = computeGroupIdToComputeGroup.get(clusterId); + CloudComputeGroupMeta computeGroup = computeGroupIdToComputeGroup.get(clusterId); if (!needVirtual && computeGroup.isVirtual()) { continue; } @@ -1430,7 +1431,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 7ea77b085a856d..366fc59c168be8 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.cluster.ClusterNamespace; @@ -97,9 +97,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); @@ -113,7 +113,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 d915207fb3370f..de22ae9899cc15 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 @@ -23,8 +23,8 @@ import org.apache.doris.catalog.OlapTable; import org.apache.doris.catalog.ScalarType; 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 ff19f67dcb64dc..5c0a1f84b2d7db 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()); @@ -171,7 +171,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 11b288dfc2a218..36f2b776a77493 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,7 +19,7 @@ import org.apache.doris.analysis.UserIdentity; import org.apache.doris.catalog.Env; -import org.apache.doris.cloud.catalog.ComputeGroup; +import org.apache.doris.cloud.catalog.CloudComputeGroupMeta; import org.apache.doris.common.Config; import org.apache.doris.qe.ConnectContext; import org.apache.doris.resource.Tag; @@ -64,7 +64,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); @@ -79,14 +79,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); @@ -104,14 +104,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); @@ -141,14 +141,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); @@ -178,14 +178,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); @@ -227,14 +227,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); @@ -276,14 +276,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); @@ -329,15 +329,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); @@ -361,7 +361,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 @@ -383,7 +383,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 @@ -417,7 +417,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 @@ -468,7 +468,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 @@ -519,7 +519,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 @@ -579,7 +579,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 @@ -642,7 +642,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 @@ -693,7 +693,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 @@ -734,11 +734,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 @@ -806,20 +806,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 @@ -906,11 +906,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 5047357b4a65ff..b086138681dd19 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; @@ -364,12 +364,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 8ff8b60ffe05b6..be55e5ac1ba88a 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 @@ -20,7 +20,7 @@ import org.apache.doris.backup.CatalogMocker; import org.apache.doris.catalog.Database; 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; @@ -52,7 +52,7 @@ public class AlterComputeGroupCommandTest { @Mocked private CloudSystemInfoService cloudSystemInfoService; @Mocked - private ComputeGroup computeGroup; + private CloudComputeGroupMeta computeGroup; private Database db; private void runBefore() throws Exception {