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());