This is an automated email from the ASF dual-hosted git repository.
wchevreuil pushed a commit to branch branch-3
in repository https://gitbox.apache.org/repos/asf/hbase.git
The following commit(s) were added to refs/heads/branch-3 by this push:
new 123cedef712 HBASE-30300 CacheAwareLoadBalancer disproportionally
favouring cache ratio over skewness (#8502)
123cedef712 is described below
commit 123cedef7128b7977c2a1df04c743350cafea510
Author: Wellington Ramos Chevreuil <[email protected]>
AuthorDate: Wed Jul 29 20:55:56 2026 +0100
HBASE-30300 CacheAwareLoadBalancer disproportionally favouring cache ratio
over skewness (#8502)
Co-authored-by: Claude Code Opus 4.6 <[email protected]>
Signed-off-by: Tak Lon (Stephen) Wu <[email protected]>
Reviewed-by: Kevin Geiszler <[email protected]>
Change-Id: I263b26675495effa06a5231144b89a6f4d1634e7
---
.../master/balancer/BalancerClusterState.java | 16 +-
.../master/balancer/CacheAwareLoadBalancer.java | 166 ++++---
.../balancer/TestCacheAwareLoadBalancer.java | 542 ++++++++++++++++++++-
...lancerWithCacheAwareLoadBalancerAsInternal.java | 84 +++-
4 files changed, 698 insertions(+), 110 deletions(-)
diff --git
a/hbase-balancer/src/main/java/org/apache/hadoop/hbase/master/balancer/BalancerClusterState.java
b/hbase-balancer/src/main/java/org/apache/hadoop/hbase/master/balancer/BalancerClusterState.java
index aa73b52a404..b1d40475493 100644
---
a/hbase-balancer/src/main/java/org/apache/hadoop/hbase/master/balancer/BalancerClusterState.java
+++
b/hbase-balancer/src/main/java/org/apache/hadoop/hbase/master/balancer/BalancerClusterState.java
@@ -400,8 +400,8 @@ class BalancerClusterState {
serversPerRack);
}
- this.serverBlockCacheFreeSize = new long[numServers];
if (serverBlockCacheFreeByServer != null) {
+ this.serverBlockCacheFreeSize = new long[numServers];
for (int i = 0; i < numServers; i++) {
ServerName sn = servers[i];
this.serverBlockCacheFreeSize[i] =
@@ -909,6 +909,20 @@ class BalancerClusterState {
}
numRegionsPerServerPerTable[tableIndex][newServer]++;
+ // Update server block cache free size to reflect the region move. The
destination server
+ // will need to cache this region (reducing its free space), and the
source server frees
+ // space when blocks are evicted on region close.
+ if (serverBlockCacheFreeSize != null) {
+ long regionSizeBytes = (long) getRegionSizeMinusColdDataMB(region) *
1024L * 1024L;
+ if (regionSizeBytes > 0) {
+ serverBlockCacheFreeSize[newServer] =
+ Math.max(0, serverBlockCacheFreeSize[newServer] - regionSizeBytes);
+ if (oldServer >= 0) {
+ serverBlockCacheFreeSize[oldServer] += regionSizeBytes;
+ }
+ }
+ }
+
// update for servers
int primary = regionIndexToPrimaryIndex[region];
if (oldServer >= 0) {
diff --git
a/hbase-balancer/src/main/java/org/apache/hadoop/hbase/master/balancer/CacheAwareLoadBalancer.java
b/hbase-balancer/src/main/java/org/apache/hadoop/hbase/master/balancer/CacheAwareLoadBalancer.java
index 4064c2575c2..9c8080a9185 100644
---
a/hbase-balancer/src/main/java/org/apache/hadoop/hbase/master/balancer/CacheAwareLoadBalancer.java
+++
b/hbase-balancer/src/main/java/org/apache/hadoop/hbase/master/balancer/CacheAwareLoadBalancer.java
@@ -27,6 +27,7 @@ package org.apache.hadoop.hbase.master.balancer;
*/
import static
org.apache.hadoop.hbase.HConstants.BUCKET_CACHE_PERSISTENT_PATH_KEY;
+import static org.apache.hadoop.hbase.HConstants.BUCKET_CACHE_SIZE_KEY;
import java.math.BigDecimal;
import java.text.DecimalFormat;
@@ -94,6 +95,8 @@ public class CacheAwareLoadBalancer extends
StochasticLoadBalancer {
private float lowCacheRatioThreshold;
private float potentialCacheRatioAfterMove;
private float minFreeCacheSpaceFactor;
+ private long cachePrefetchOverheadBytes;
+ private boolean cacheSpaceTrackingEnabled;
private BigDecimal simulatedRatio = BigDecimal.ZERO;
@@ -111,6 +114,16 @@ public class CacheAwareLoadBalancer extends
StochasticLoadBalancer {
POTENTIAL_CACHE_RATIO_AFTER_MOVE_DEFAULT);
minFreeCacheSpaceFactor =
configuration.getFloat(MIN_FREE_CACHE_SPACE_FACTOR_KEY,
MIN_FREE_CACHE_SPACE_FACTOR_DEFAULT);
+ float bucketCacheSizeMB = configuration.getFloat(BUCKET_CACHE_SIZE_KEY,
0F);
+ float acceptableFactor =
configuration.getFloat("hbase.bucketcache.acceptfactor", 0.95f);
+ cachePrefetchOverheadBytes =
+ (long) (bucketCacheSizeMB * 1024L * 1024L * (1 - acceptableFactor));
+ cacheSpaceTrackingEnabled = bucketCacheSizeMB > 0;
+ if (!cacheSpaceTrackingEnabled) {
+ LOG.warn("{} is not configured on the master. The free-space relocation
heuristic cannot "
+ + "account for the prefetch threshold and may move regions to servers
where prefetch "
+ + "will be blocked.", BUCKET_CACHE_SIZE_KEY);
+ }
}
@Override
@@ -145,11 +158,14 @@ public class CacheAwareLoadBalancer extends
StochasticLoadBalancer {
}
protected Map<ServerName, Long> getServerBlockCacheFreeBytes() {
- if (clusterStatus == null) {
+ if (clusterStatus == null || !cacheSpaceTrackingEnabled) {
return null;
}
Map<ServerName, Long> map = new HashMap<>();
- clusterStatus.getLiveServerMetrics().forEach((sn, sm) -> map.put(sn,
sm.getCacheFreeSize()));
+ clusterStatus.getLiveServerMetrics().forEach((sn, sm) -> {
+ long effectiveFree = Math.max(0, sm.getCacheFreeSize() -
cachePrefetchOverheadBytes);
+ map.put(sn, effectiveFree);
+ });
return map;
}
@@ -229,6 +245,19 @@ public class CacheAwareLoadBalancer extends
StochasticLoadBalancer {
return null;
}
+ private boolean serverHasCacheSpaceForRegion(BalancerClusterState cluster,
int region,
+ int server) {
+ if (cluster.serverBlockCacheFreeSize == null) {
+ return true;
+ }
+ int regionSizeMb = cluster.getRegionSizeMinusColdDataMB(region);
+ if (regionSizeMb <= 0) {
+ return true;
+ }
+ long bytesNeeded = (long) regionSizeMb * 1024L * 1024L;
+ return cluster.serverBlockCacheFreeSize[server] >= bytesNeeded;
+ }
+
@Override
public long getThrottleDurationMs(RegionPlan plan) {
Pair<ServerName, Float> rsRatio =
this.regionCacheRatioOnOldServerMap.get(plan.getRegionName());
@@ -239,20 +268,38 @@ public class CacheAwareLoadBalancer extends
StochasticLoadBalancer {
LOG.debug("Moving region {} to server {} with cache ratio {}. No
throttling needed.",
plan.getRegionInfo().getEncodedName(), plan.getDestination(),
rsRatio.getSecond());
return 0L;
+ }
+ // Skip throttling for regions with low cache ratio on their source server
— there is
+ // negligible cached data to lose, so no warm-up delay is needed on the
destination.
+ float cacheRatioOnSource = getRegionCacheRatioOnSource(plan);
+ if (cacheRatioOnSource < lowCacheRatioThreshold) {
+ LOG.debug(
+ "Moving region {} to server {} with low cache ratio {} on source. No
throttling needed.",
+ plan.getRegionInfo().getEncodedName(), plan.getDestination(),
cacheRatioOnSource);
+ return 0L;
+ }
+
+ if (rsRatio != null) {
+ LOG.debug("Moving region {} to server {} with cache ratio: {}.
Throttling move for {}ms.",
+ plan.getRegionInfo().getEncodedName(), plan.getDestination(),
+ plan.getDestination().equals(rsRatio.getFirst()) ? rsRatio.getSecond()
: "unknown",
+ sleepTime);
} else {
- if (rsRatio != null) {
- LOG.debug("Moving region {} to server {} with cache ratio: {}.
Throttling move for {}ms.",
- plan.getRegionInfo().getEncodedName(), plan.getDestination(),
- plan.getDestination().equals(rsRatio.getFirst()) ?
rsRatio.getSecond() : "unknown",
- sleepTime);
- } else {
- LOG.debug(
- "Moving region {} to server {} with no cache ratio info for the
region. "
- + "Throttling move for {}ms.",
- plan.getRegionInfo().getEncodedName(), plan.getDestination(),
sleepTime);
- }
- return sleepTime;
+ LOG.debug(
+ "Moving region {} to server {} with no cache ratio info for the
region. "
+ + "Throttling move for {}ms.",
+ plan.getRegionInfo().getEncodedName(), plan.getDestination(),
sleepTime);
+ }
+ return sleepTime;
+ }
+
+ private float getRegionCacheRatioOnSource(RegionPlan plan) {
+ Deque<BalancerRegionLoad> regionLoad = loads.get(plan.getRegionName());
+ if (regionLoad != null && !regionLoad.isEmpty()) {
+ return regionLoad.getFirst().getCurrentRegionCacheRatio();
}
+ // Unknown cache ratio — assume it may be cached and require throttling
+ return 1.0f;
}
@Override
@@ -384,6 +431,28 @@ public class CacheAwareLoadBalancer extends
StochasticLoadBalancer {
return false;
}
+ // If the region is already well-cached on its current server, don't
disrupt it.
+ // The old server's historical cache data may be stale, and moving a hot
region
+ // causes unnecessary cache churn.
+ if (cacheRatioOnCurrentServer >= ratioThreshold) {
+ if (LOG.isDebugEnabled()) {
+ LOG.debug(
+ "Region {} not moved from {} to {} as it is already well-cached
({}) on current server",
+ cluster.regions[regionIndex].getEncodedName(),
cluster.servers[currentServerIndex],
+ cluster.servers[oldServerIndex], cacheRatioOnCurrentServer);
+ }
+ return false;
+ }
+
+ if (!serverHasCacheSpaceForRegion(cluster, regionIndex, oldServerIndex))
{
+ if (LOG.isDebugEnabled()) {
+ LOG.debug("Region {} not moved from {} to {} as destination server
lacks cache space",
+ cluster.regions[regionIndex].getEncodedName(),
cluster.servers[currentServerIndex],
+ cluster.servers[oldServerIndex]);
+ }
+ return false;
+ }
+
DecimalFormat df = new DecimalFormat("#");
df.setMaximumFractionDigits(4);
@@ -444,60 +513,6 @@ public class CacheAwareLoadBalancer extends
StochasticLoadBalancer {
@Override
BalanceAction pickRandomRegions(BalancerClusterState cluster, int
thisServer, int otherServer) {
simulatedRatio = BigDecimal.ZERO;
- // First move all the regions which were hosted previously on some other
server back to their
- // old servers
- if (
- !regionCacheRatioOnOldServerMap.isEmpty()
- && regionCacheRatioOnOldServerMap.entrySet().iterator().hasNext()
- ) {
- // Get the first region index in the historical cache ratio list
- Map.Entry<String, Pair<ServerName, Float>> regionEntry =
- regionCacheRatioOnOldServerMap.entrySet().iterator().next();
- String regionEncodedName = regionEntry.getKey();
-
- RegionInfo regionInfo = getRegionInfoByEncodedName(cluster,
regionEncodedName);
- if (regionInfo == null) {
- LOG.warn("Region {} does not exist", regionEncodedName);
- regionCacheRatioOnOldServerMap.remove(regionEncodedName);
- return BalanceAction.NULL_ACTION;
- }
- if (regionInfo.isMetaRegion() ||
regionInfo.getTable().isSystemTable()) {
- regionCacheRatioOnOldServerMap.remove(regionEncodedName);
- return BalanceAction.NULL_ACTION;
- }
-
- int regionIndex = cluster.regionsToIndex.get(regionInfo);
-
- // Get the current host name for this region
- thisServer = cluster.regionIndexToServerIndex[regionIndex];
-
- // Get the old server index
- otherServer =
cluster.serversToIndex.get(regionEntry.getValue().getFirst().getAddress());
-
- regionCacheRatioOnOldServerMap.remove(regionEncodedName);
-
- if (otherServer < 0) {
- // The old server has been moved to other host and hence, the region
cannot be moved back
- // to the old server
- if (LOG.isDebugEnabled()) {
- LOG.debug(
- "CacheAwareSkewnessCandidateGenerator: Region {} not moved to
the old "
- + "server {} as the server does not exist",
- regionEncodedName,
regionEntry.getValue().getFirst().getHostname());
- }
- return BalanceAction.NULL_ACTION;
- }
-
- if (LOG.isDebugEnabled()) {
- LOG.debug(
- "CacheAwareSkewnessCandidateGenerator: Region {} moved from {} to
{} as it "
- + "was hosted their earlier",
- regionEncodedName, cluster.servers[thisServer].getHostname(),
- cluster.servers[otherServer].getHostname());
- }
-
- return getAction(thisServer, regionIndex, otherServer, -1);
- }
if (thisServer < 0 || otherServer < 0) {
return BalanceAction.NULL_ACTION;
@@ -574,7 +589,7 @@ public class CacheAwareLoadBalancer extends
StochasticLoadBalancer {
@Override
public final void updateWeight(Map<Class<? extends CandidateGenerator>,
Double> weights) {
- weights.merge(LoadCandidateGenerator.class, cost(), Double::sum);
+ weights.merge(CacheAwareSkewnessCandidateGenerator.class, cost(),
Double::sum);
}
}
@@ -671,6 +686,9 @@ public class CacheAwareLoadBalancer extends
StochasticLoadBalancer {
* from low ratios when capacity exists somewhere in the cluster.
*/
private double potentialBestWeightedFromFreeCache(BalancerClusterState
cluster, int region) {
+ if (cluster.serverBlockCacheFreeSize == null) {
+ return 0.0;
+ }
float observedRatio = cluster.getSumRegionCacheAndColdDataRatio(region);
if (observedRatio >= lowCacheRatioThreshold) {
return 0.0;
@@ -701,9 +719,11 @@ public class CacheAwareLoadBalancer extends
StochasticLoadBalancer {
if (simulatedRatio.equals(BigDecimal.ZERO)) {
double potentialCachedSizeOnNewServer =
cluster.getRegionSizeMinusColdDataMB(region) *
potentialCacheRatioAfterMove;
- boolean simulateCacheBasedOnFreeSpace =
- cluster.getOrComputeRegionCacheRatio(region, oldServer) <
lowCacheRatioThreshold
- && cluster.serverBlockCacheFreeSize[newServer] >=
potentialCachedSizeOnNewServer;
+ long potentialCachedBytesOnNewServer =
+ (long) (potentialCachedSizeOnNewServer * 1024L * 1024L);
+ boolean simulateCacheBasedOnFreeSpace =
cluster.serverBlockCacheFreeSize != null
+ && cluster.getOrComputeRegionCacheRatio(region, oldServer) <
lowCacheRatioThreshold
+ && cluster.serverBlockCacheFreeSize[newServer] >=
potentialCachedBytesOnNewServer;
double regionCacheRatioOnNewServer = simulateCacheBasedOnFreeSpace
? potentialCachedSizeOnNewServer
: cluster.getOrComputeWeightedRegionCacheRatio(region, newServer);
@@ -740,7 +760,7 @@ public class CacheAwareLoadBalancer extends
StochasticLoadBalancer {
@Override
public void updateWeight(Map<Class<? extends CandidateGenerator>, Double>
weights) {
- weights.merge(LoadCandidateGenerator.class, cost(), Double::sum);
+ weights.merge(CacheAwareCandidateGenerator.class, cost(), Double::sum);
}
}
}
diff --git
a/hbase-server/src/test/java/org/apache/hadoop/hbase/master/balancer/TestCacheAwareLoadBalancer.java
b/hbase-server/src/test/java/org/apache/hadoop/hbase/master/balancer/TestCacheAwareLoadBalancer.java
index 8b8fdbf11ad..04110c31498 100644
---
a/hbase-server/src/test/java/org/apache/hadoop/hbase/master/balancer/TestCacheAwareLoadBalancer.java
+++
b/hbase-server/src/test/java/org/apache/hadoop/hbase/master/balancer/TestCacheAwareLoadBalancer.java
@@ -17,16 +17,21 @@
*/
package org.apache.hadoop.hbase.master.balancer;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.junit.jupiter.api.Assertions.fail;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
import java.util.ArrayList;
import java.util.HashMap;
+import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Random;
+import java.util.Set;
import java.util.TreeMap;
import java.util.concurrent.ThreadLocalRandom;
import org.apache.hadoop.conf.Configuration;
@@ -44,6 +49,7 @@ import org.apache.hadoop.hbase.client.TableDescriptorBuilder;
import org.apache.hadoop.hbase.master.RegionPlan;
import org.apache.hadoop.hbase.testclassification.LargeTests;
import org.apache.hadoop.hbase.util.Bytes;
+import org.apache.hadoop.hbase.util.Pair;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Tag;
import org.junit.jupiter.api.Test;
@@ -86,9 +92,9 @@ public class TestCacheAwareLoadBalancer extends
BalancerTestBase {
return tds;
}
- private ServerMetrics mockServerMetricsWithRegionCacheInfo(ServerName server,
- List<RegionInfo> regionsOnServer, float currentCacheRatio,
List<RegionInfo> oldRegionCacheInfo,
- int oldRegionCachedSize, int regionSize) {
+ private ServerMetrics mockServerMetricsWithRegionCacheInfo(List<RegionInfo>
regionsOnServer,
+ float currentCacheRatio, List<RegionInfo> oldRegionCacheInfo, int
oldRegionCachedSize,
+ int regionSize) {
ServerMetrics serverMetrics = mock(ServerMetrics.class);
Map<byte[], RegionMetrics> regionLoadMap = new
TreeMap<>(Bytes.BYTES_COMPARATOR);
for (RegionInfo info : regionsOnServer) {
@@ -107,6 +113,7 @@ public class TestCacheAwareLoadBalancer extends
BalancerTestBase {
oldCacheRatioMap.put(info.getEncodedName(), oldRegionCachedSize);
}
when(serverMetrics.getRegionCachedInfo()).thenReturn(oldCacheRatioMap);
+ when(serverMetrics.getCacheFreeSize()).thenReturn(100L * 1024 * 1024 *
1024);
return serverMetrics;
}
@@ -116,11 +123,527 @@ public class TestCacheAwareLoadBalancer extends
BalancerTestBase {
tableDescs = constructTableDesc(false);
Configuration conf = HBaseConfiguration.create();
conf.set(HConstants.BUCKET_CACHE_PERSISTENT_PATH_KEY,
"prefetch_file_list");
+ conf.setFloat(HConstants.BUCKET_CACHE_SIZE_KEY, 10);
loadBalancer = new CacheAwareLoadBalancer();
loadBalancer.setClusterInfoProvider(new DummyClusterInfoProvider(conf));
loadBalancer.loadConf(conf);
}
+ @Test
+ public void testRegionsNotCachedOnOldServerAndCurrentServer() throws
Exception {
+ // The regions are not cached on old server as well as the current server.
This causes
+ // skewness in the region allocation which should be fixed by the balancer
+
+ Map<ServerName, List<RegionInfo>> clusterState = new HashMap<>();
+ ServerName server0 = servers.get(0);
+ ServerName server1 = servers.get(1);
+ ServerName server2 = servers.get(2);
+
+ // Simulate that the regions previously hosted by server1 are now hosted
on server0
+ List<RegionInfo> regionsOnServer0 = randomRegions(10);
+ List<RegionInfo> regionsOnServer1 = randomRegions(0);
+ List<RegionInfo> regionsOnServer2 = randomRegions(5);
+
+ clusterState.put(server0, regionsOnServer0);
+ clusterState.put(server1, regionsOnServer1);
+ clusterState.put(server2, regionsOnServer2);
+
+ // Mock cluster metrics — give only server1 free cache so moves are
directed there
+ Map<ServerName, ServerMetrics> serverMetricsMap = new TreeMap<>();
+ ServerMetrics sm0 =
+ mockServerMetricsWithRegionCacheInfo(regionsOnServer0, 0.0f, new
ArrayList<>(), 0, 10);
+ when(sm0.getCacheFreeSize()).thenReturn(0L);
+ ServerMetrics sm1 =
+ mockServerMetricsWithRegionCacheInfo(regionsOnServer1, 0.0f, new
ArrayList<>(), 0, 10);
+ ServerMetrics sm2 =
+ mockServerMetricsWithRegionCacheInfo(regionsOnServer2, 0.0f, new
ArrayList<>(), 0, 10);
+ when(sm2.getCacheFreeSize()).thenReturn(0L);
+ serverMetricsMap.put(server0, sm0);
+ serverMetricsMap.put(server1, sm1);
+ serverMetricsMap.put(server2, sm2);
+ ClusterMetrics clusterMetrics = mock(ClusterMetrics.class);
+ when(clusterMetrics.getLiveServerMetrics()).thenReturn(serverMetricsMap);
+ loadBalancer.updateClusterMetrics(clusterMetrics);
+
+ Map<TableName, Map<ServerName, List<RegionInfo>>> LoadOfAllTable =
+ (Map) mockClusterServersWithTables(clusterState);
+ List<RegionPlan> plans = loadBalancer.balanceCluster(LoadOfAllTable);
+ Set<RegionInfo> regionsMovedFromServer0 = new HashSet<>();
+ Map<ServerName, List<RegionInfo>> targetServers = new HashMap<>();
+ for (RegionPlan plan : plans) {
+ if (plan.getSource().equals(server0)) {
+ regionsMovedFromServer0.add(plan.getRegionInfo());
+ if (!targetServers.containsKey(plan.getDestination())) {
+ targetServers.put(plan.getDestination(), new ArrayList<>());
+ }
+ targetServers.get(plan.getDestination()).add(plan.getRegionInfo());
+ }
+ }
+ // should move at least 5 regions from server0 to balance cluster (10/0/5
-> ~5/5/5)
+ assertTrue(regionsMovedFromServer0.size() >= 5,
+ "Expected at least 5 moves from server0, got " +
regionsMovedFromServer0.size());
+ }
+
+ /**
+ * Regions on the overloaded RS report low block-cache ratio; no RS reports
prefetch/historical
+ * cache for those regions (so {@link
CacheAwareLoadBalancer.CacheAwareCandidateGenerator} has no
+ * "old server" to prefer). Another RS has ample free block cache. The
balancer should still emit
+ * plans that shed load from the hot RS onto the idle RS with spare cache
capacity.
+ */
+ @Test
+ public void
testLowCacheRatioNoHistoricalCacheRelocatesWhenTargetHasFreeBlockCache()
+ throws Exception {
+ Map<ServerName, List<RegionInfo>> clusterState = new HashMap<>();
+ ServerName server0 = servers.get(0);
+ ServerName server1 = servers.get(1);
+ ServerName server2 = servers.get(2);
+
+ List<RegionInfo> regionsOnServer0 = randomRegions(10);
+ List<RegionInfo> regionsOnServer1 = randomRegions(0);
+ List<RegionInfo> regionsOnServer2 = randomRegions(5);
+
+ clusterState.put(server0, regionsOnServer0);
+ clusterState.put(server1, regionsOnServer1);
+ clusterState.put(server2, regionsOnServer2);
+
+ // Below LOW_CACHE_RATIO_FOR_RELOCATION_DEFAULT (0.35);
+ ServerMetrics sm0 =
+ mockServerMetricsWithRegionCacheInfo(regionsOnServer0, 0.1f, new
ArrayList<>(), 0, 10);
+ when(sm0.getCacheFreeSize()).thenReturn(0L);
+ ServerMetrics sm1 =
+ mockServerMetricsWithRegionCacheInfo(regionsOnServer1, 0.0f, new
ArrayList<>(), 0, 10);
+ // Simulates 1GB free cache space on server1
+ when(sm1.getCacheFreeSize()).thenReturn(1024L * 1024 * 1024);
+ ServerMetrics sm2 =
+ mockServerMetricsWithRegionCacheInfo(regionsOnServer2, 1.0f, new
ArrayList<>(), 0, 10);
+ when(sm2.getCacheFreeSize()).thenReturn(0L);
+
+ Map<ServerName, ServerMetrics> serverMetricsMap = new TreeMap<>();
+ serverMetricsMap.put(server0, sm0);
+ serverMetricsMap.put(server1, sm1);
+ serverMetricsMap.put(server2, sm2);
+ ClusterMetrics clusterMetrics = mock(ClusterMetrics.class);
+ when(clusterMetrics.getLiveServerMetrics()).thenReturn(serverMetricsMap);
+ loadBalancer.updateClusterMetrics(clusterMetrics);
+
+ assertTrue(loadBalancer.regionCacheRatioOnOldServerMap.isEmpty());
+
+ Map<TableName, Map<ServerName, List<RegionInfo>>> loadOfAllTable =
+ (Map) mockClusterServersWithTables(clusterState);
+ List<RegionPlan> plans = loadBalancer.balanceCluster(loadOfAllTable);
+ assertNotNull(plans);
+
+ Set<RegionInfo> regionsMovedFromServer0 = new HashSet<>();
+ Map<ServerName, List<RegionInfo>> targetServers = new HashMap<>();
+ for (RegionPlan plan : plans) {
+ if (plan.getSource().equals(server0)) {
+ regionsMovedFromServer0.add(plan.getRegionInfo());
+ if (!targetServers.containsKey(plan.getDestination())) {
+ targetServers.put(plan.getDestination(), new ArrayList<>());
+ }
+ targetServers.get(plan.getDestination()).add(plan.getRegionInfo());
+ }
+ }
+ assertEquals(5, regionsMovedFromServer0.size());
+ assertNotNull(targetServers.get(server1));
+ assertEquals(5, targetServers.get(server1).size());
+ }
+
+ @Test
+ public void
testRegionsPartiallyCachedOnOldServerAndNotCachedOnCurrentServer() throws
Exception {
+ // The regions are partially cached on old server but not cached on the
current server
+
+ Map<ServerName, List<RegionInfo>> clusterState = new HashMap<>();
+ ServerName server0 = servers.get(0);
+ ServerName server1 = servers.get(1);
+ ServerName server2 = servers.get(2);
+
+ // Simulate that the regions previously hosted by server1 are now hosted
on server0
+ List<RegionInfo> regionsOnServer0 = randomRegions(10);
+ List<RegionInfo> regionsOnServer1 = randomRegions(0);
+ List<RegionInfo> regionsOnServer2 = randomRegions(5);
+
+ clusterState.put(server0, regionsOnServer0);
+ clusterState.put(server1, regionsOnServer1);
+ clusterState.put(server2, regionsOnServer2);
+
+ // Mock cluster metrics
+
+ // Mock 5 regions from server0 were previously hosted on server1
+ List<RegionInfo> oldCachedRegions = regionsOnServer0.subList(5,
regionsOnServer0.size());
+
+ Map<ServerName, ServerMetrics> serverMetricsMap = new TreeMap<>();
+ serverMetricsMap.put(server0,
+ mockServerMetricsWithRegionCacheInfo(regionsOnServer0, 0.0f, new
ArrayList<>(), 0, 10));
+ serverMetricsMap.put(server1,
+ mockServerMetricsWithRegionCacheInfo(regionsOnServer1, 0.0f,
oldCachedRegions, 6, 10));
+ serverMetricsMap.put(server2,
+ mockServerMetricsWithRegionCacheInfo(regionsOnServer2, 0.0f, new
ArrayList<>(), 0, 10));
+ ClusterMetrics clusterMetrics = mock(ClusterMetrics.class);
+ when(clusterMetrics.getLiveServerMetrics()).thenReturn(serverMetricsMap);
+ loadBalancer.updateClusterMetrics(clusterMetrics);
+
+ Map<TableName, Map<ServerName, List<RegionInfo>>> LoadOfAllTable =
+ (Map) mockClusterServersWithTables(clusterState);
+ List<RegionPlan> plans = loadBalancer.balanceCluster(LoadOfAllTable);
+ Set<RegionInfo> regionsMovedFromServer0 = new HashSet<>();
+ Map<ServerName, List<RegionInfo>> targetServers = new HashMap<>();
+ for (RegionPlan plan : plans) {
+ if (plan.getSource().equals(server0)) {
+ regionsMovedFromServer0.add(plan.getRegionInfo());
+ if (!targetServers.containsKey(plan.getDestination())) {
+ targetServers.put(plan.getDestination(), new ArrayList<>());
+ }
+ targetServers.get(plan.getDestination()).add(plan.getRegionInfo());
+ }
+ }
+ // should move regions from server0 to server1 (old-cached regions should
be among them)
+ assertEquals(5, regionsMovedFromServer0.size());
+ assertNotNull(targetServers.get(server1));
+ int oldCachedOnServer1 = 0;
+ for (RegionInfo ri : oldCachedRegions) {
+ if (targetServers.get(server1).contains(ri)) {
+ oldCachedOnServer1++;
+ }
+ }
+ assertTrue(oldCachedOnServer1 > 0,
+ "Expected old-cached regions to move to server1, got " +
oldCachedOnServer1);
+ }
+
+ @Test
+ public void testThrottlingRegionBeyondThreshold() throws Exception {
+ Configuration conf = HBaseConfiguration.create();
+ CacheAwareLoadBalancer balancer = new CacheAwareLoadBalancer();
+ balancer.loadConf(conf);
+ balancer.setClusterInfoProvider(new DummyClusterInfoProvider(conf));
+ balancer.initialize();
+ ServerName server0 = servers.get(0);
+ ServerName server1 = servers.get(1);
+ Pair<ServerName, Float> regionRatio = new Pair<>();
+ regionRatio.setFirst(server0);
+ regionRatio.setSecond(1.0f);
+ balancer.regionCacheRatioOnOldServerMap.put("region1", regionRatio);
+ RegionInfo mockedInfo = mock(RegionInfo.class);
+ when(mockedInfo.getEncodedName()).thenReturn("region1");
+ RegionPlan plan = new RegionPlan(mockedInfo, server1, server0);
+ assertEquals(0L, balancer.getThrottleDurationMs(plan));
+ }
+
+ @Test
+ public void testThrottlingRegionBelowThreshold() throws Exception {
+ Configuration conf = HBaseConfiguration.create();
+ conf.setLong(CacheAwareLoadBalancer.MOVE_THROTTLING, 100);
+ CacheAwareLoadBalancer balancer = new CacheAwareLoadBalancer();
+ balancer.loadConf(conf);
+ balancer.setClusterInfoProvider(new DummyClusterInfoProvider(conf));
+ balancer.initialize();
+ ServerName server0 = servers.get(0);
+ ServerName server1 = servers.get(1);
+ Pair<ServerName, Float> regionRatio = new Pair<>();
+ regionRatio.setFirst(server0);
+ regionRatio.setSecond(0.1f);
+ balancer.regionCacheRatioOnOldServerMap.put("region1", regionRatio);
+ RegionInfo mockedInfo = mock(RegionInfo.class);
+ when(mockedInfo.getEncodedName()).thenReturn("region1");
+ RegionPlan plan = new RegionPlan(mockedInfo, server1, server0);
+ assertEquals(100L, balancer.getThrottleDurationMs(plan));
+ }
+
+ @Test
+ public void testThrottlingCacheRatioUnknownOnTarget() throws Exception {
+ Configuration conf = HBaseConfiguration.create();
+ conf.setLong(CacheAwareLoadBalancer.MOVE_THROTTLING, 100);
+ CacheAwareLoadBalancer balancer = new CacheAwareLoadBalancer();
+ balancer.loadConf(conf);
+ balancer.setClusterInfoProvider(new DummyClusterInfoProvider(conf));
+ balancer.initialize();
+ ServerName server0 = servers.get(0);
+ ServerName server1 = servers.get(1);
+ ServerName server3 = servers.get(2);
+ // setting region cache ratio 100% on server 3, though this is not the
target in the region plan
+ Pair<ServerName, Float> regionRatio = new Pair<>();
+ regionRatio.setFirst(server3);
+ regionRatio.setSecond(1.0f);
+ balancer.regionCacheRatioOnOldServerMap.put("region1", regionRatio);
+ RegionInfo mockedInfo = mock(RegionInfo.class);
+ when(mockedInfo.getEncodedName()).thenReturn("region1");
+ RegionPlan plan = new RegionPlan(mockedInfo, server1, server0);
+ assertEquals(100L, balancer.getThrottleDurationMs(plan));
+ }
+
+ @Test
+ public void testThrottlingCacheRatioUnknownForRegion() throws Exception {
+ Configuration conf = HBaseConfiguration.create();
+ conf.setLong(CacheAwareLoadBalancer.MOVE_THROTTLING, 100);
+ CacheAwareLoadBalancer balancer = new CacheAwareLoadBalancer();
+ balancer.loadConf(conf);
+ balancer.setClusterInfoProvider(new DummyClusterInfoProvider(conf));
+ balancer.initialize();
+ ServerName server0 = servers.get(0);
+ ServerName server1 = servers.get(1);
+ ServerName server3 = servers.get(2);
+ // No cache ratio available for region1
+ RegionInfo mockedInfo = mock(RegionInfo.class);
+ when(mockedInfo.getEncodedName()).thenReturn("region1");
+ RegionPlan plan = new RegionPlan(mockedInfo, server1, server0);
+ assertEquals(100L, balancer.getThrottleDurationMs(plan));
+ }
+
+ @Test
+ public void testRegionPlansSortedByCacheRatioOnTarget() throws Exception {
+ // The regions are fully cached on old server
+
+ Map<ServerName, List<RegionInfo>> clusterState = new HashMap<>();
+ ServerName server0 = servers.get(0);
+ ServerName server1 = servers.get(1);
+ ServerName server2 = servers.get(2);
+
+ // Simulate on RS with all regions, and two RSes with no regions
+ List<RegionInfo> regionsOnServer0 = randomRegions(15);
+ List<RegionInfo> regionsOnServer1 = randomRegions(0);
+ List<RegionInfo> regionsOnServer2 = randomRegions(0);
+
+ clusterState.put(server0, regionsOnServer0);
+ clusterState.put(server1, regionsOnServer1);
+ clusterState.put(server2, regionsOnServer2);
+
+ // Mock cluster metrics
+ // Mock 5 regions from server0 were previously hosted on server1
+ List<RegionInfo> oldCachedRegions1 = regionsOnServer0.subList(5, 10);
+ List<RegionInfo> oldCachedRegions2 = regionsOnServer0.subList(10,
regionsOnServer0.size());
+ Map<ServerName, ServerMetrics> serverMetricsMap = new TreeMap<>();
+ // mock server metrics to set cache ratio as 0 in the RS 0
+ serverMetricsMap.put(server0,
+ mockServerMetricsWithRegionCacheInfo(regionsOnServer0, 0.0f, new
ArrayList<>(), 0, 10));
+ // mock server metrics to set cache ratio as 1 in the RS 1
+ serverMetricsMap.put(server1,
+ mockServerMetricsWithRegionCacheInfo(regionsOnServer1, 0.0f,
oldCachedRegions1, 10, 10));
+ // mock server metrics to set cache ratio as .8 in the RS 2
+ serverMetricsMap.put(server2,
+ mockServerMetricsWithRegionCacheInfo(regionsOnServer2, 0.0f,
oldCachedRegions2, 8, 10));
+ ClusterMetrics clusterMetrics = mock(ClusterMetrics.class);
+ when(clusterMetrics.getLiveServerMetrics()).thenReturn(serverMetricsMap);
+ loadBalancer.updateClusterMetrics(clusterMetrics);
+
+ Map<TableName, Map<ServerName, List<RegionInfo>>> LoadOfAllTable =
+ (Map) mockClusterServersWithTables(clusterState);
+ List<RegionPlan> plans = loadBalancer.balanceCluster(LoadOfAllTable);
+ LOG.debug("plans size: {}", plans.size());
+ LOG.debug("plans: {}", plans);
+ // Plans are sorted by cache ratio on destination (descending). Verify
that plans
+ // for old-cached regions going to their correct servers appear before
plans with no
+ // cache data on destination.
+ float prevRatio = Float.MAX_VALUE;
+ int oldCached1Count = 0;
+ int oldCached2Count = 0;
+ for (RegionPlan plan : plans) {
+ LOG.debug("plan region: {}, target server: {}",
plan.getRegionInfo().getEncodedName(),
+ plan.getDestination().getServerName());
+ float ratio = 0f;
+ if (
+ oldCachedRegions1.contains(plan.getRegionInfo()) &&
server1.equals(plan.getDestination())
+ ) {
+ ratio = 1.0f;
+ oldCached1Count++;
+ } else if (
+ oldCachedRegions2.contains(plan.getRegionInfo()) &&
server2.equals(plan.getDestination())
+ ) {
+ ratio = 0.8f;
+ oldCached2Count++;
+ }
+ assertTrue(ratio <= prevRatio,
+ "Plans should be sorted by cache ratio on destination (descending)");
+ prevRatio = ratio;
+ }
+ // The cache-aware generator should move at least some old-cached regions
to their
+ // cached servers. Exact count depends on stochastic walk order.
+ assertTrue(oldCached1Count > 0, "Some old-cached regions should move to
server1");
+ assertTrue(oldCached2Count > 0, "Some old-cached regions should move to
server2");
+
+ }
+
+ @Test
+ public void testRegionsFullyCachedOnOldServerAndNotCachedOnCurrentServers()
throws Exception {
+ // The regions are fully cached on old server
+
+ Map<ServerName, List<RegionInfo>> clusterState = new HashMap<>();
+ ServerName server0 = servers.get(0);
+ ServerName server1 = servers.get(1);
+ ServerName server2 = servers.get(2);
+
+ // Simulate that the regions previously hosted by server1 are now hosted
on server0
+ List<RegionInfo> regionsOnServer0 = randomRegions(10);
+ List<RegionInfo> regionsOnServer1 = randomRegions(0);
+ List<RegionInfo> regionsOnServer2 = randomRegions(5);
+
+ clusterState.put(server0, regionsOnServer0);
+ clusterState.put(server1, regionsOnServer1);
+ clusterState.put(server2, regionsOnServer2);
+
+ // Mock cluster metrics
+
+ // Mock 5 regions from server0 were previously hosted on server1
+ List<RegionInfo> oldCachedRegions = regionsOnServer0.subList(5,
regionsOnServer0.size() - 1);
+
+ Map<ServerName, ServerMetrics> serverMetricsMap = new TreeMap<>();
+ serverMetricsMap.put(server0,
+ mockServerMetricsWithRegionCacheInfo(regionsOnServer0, 0.0f, new
ArrayList<>(), 0, 10));
+ serverMetricsMap.put(server1,
+ mockServerMetricsWithRegionCacheInfo(regionsOnServer1, 0.0f,
oldCachedRegions, 10, 10));
+ serverMetricsMap.put(server2,
+ mockServerMetricsWithRegionCacheInfo(regionsOnServer2, 0.0f, new
ArrayList<>(), 0, 10));
+ ClusterMetrics clusterMetrics = mock(ClusterMetrics.class);
+ when(clusterMetrics.getLiveServerMetrics()).thenReturn(serverMetricsMap);
+ loadBalancer.updateClusterMetrics(clusterMetrics);
+
+ Map<TableName, Map<ServerName, List<RegionInfo>>> LoadOfAllTable =
+ (Map) mockClusterServersWithTables(clusterState);
+ List<RegionPlan> plans = loadBalancer.balanceCluster(LoadOfAllTable);
+ Set<RegionInfo> regionsMovedFromServer0 = new HashSet<>();
+ Map<ServerName, List<RegionInfo>> targetServers = new HashMap<>();
+ for (RegionPlan plan : plans) {
+ if (plan.getSource().equals(server0)) {
+ regionsMovedFromServer0.add(plan.getRegionInfo());
+ if (!targetServers.containsKey(plan.getDestination())) {
+ targetServers.put(plan.getDestination(), new ArrayList<>());
+ }
+ targetServers.get(plan.getDestination()).add(plan.getRegionInfo());
+ }
+ }
+ // should move regions from server0 to server1 (old-cached regions should
be among them)
+ assertTrue(regionsMovedFromServer0.size() >= 4);
+ assertNotNull(targetServers.get(server1));
+ assertTrue(targetServers.get(server1).size() >= 4);
+ int oldCachedOnServer1 = 0;
+ for (RegionInfo ri : oldCachedRegions) {
+ if (targetServers.get(server1).contains(ri)) {
+ oldCachedOnServer1++;
+ }
+ }
+ assertTrue(oldCachedOnServer1 > 0,
+ "Expected most old-cached regions to move to server1, got " +
oldCachedOnServer1);
+ }
+
+ @Test
+ public void testRegionsFullyCachedOnOldAndCurrentServers() throws Exception {
+ // When regions are fully cached on BOTH the old and current server, the
balancer should
+ // NOT disrupt them by moving them to the old server based on potentially
stale historical
+ // cache data. Instead, it should still rebalance for skew.
+
+ Map<ServerName, List<RegionInfo>> clusterState = new HashMap<>();
+ ServerName server0 = servers.get(0);
+ ServerName server1 = servers.get(1);
+ ServerName server2 = servers.get(2);
+
+ // Simulate that the regions previously hosted by server1 are now hosted
on server0
+ List<RegionInfo> regionsOnServer0 = randomRegions(10);
+ List<RegionInfo> regionsOnServer1 = randomRegions(0);
+ List<RegionInfo> regionsOnServer2 = randomRegions(5);
+
+ clusterState.put(server0, regionsOnServer0);
+ clusterState.put(server1, regionsOnServer1);
+ clusterState.put(server2, regionsOnServer2);
+
+ // Mock cluster metrics
+
+ // Mock 4 regions from server0 were previously hosted on server1
+ List<RegionInfo> oldCachedRegions = regionsOnServer0.subList(5,
regionsOnServer0.size() - 1);
+
+ Map<ServerName, ServerMetrics> serverMetricsMap = new TreeMap<>();
+ serverMetricsMap.put(server0,
+ mockServerMetricsWithRegionCacheInfo(regionsOnServer0, 1.0f, new
ArrayList<>(), 0, 10));
+ serverMetricsMap.put(server1,
+ mockServerMetricsWithRegionCacheInfo(regionsOnServer1, 1.0f,
oldCachedRegions, 10, 10));
+ serverMetricsMap.put(server2,
+ mockServerMetricsWithRegionCacheInfo(regionsOnServer2, 1.0f, new
ArrayList<>(), 0, 10));
+ ClusterMetrics clusterMetrics = mock(ClusterMetrics.class);
+ when(clusterMetrics.getLiveServerMetrics()).thenReturn(serverMetricsMap);
+ loadBalancer.updateClusterMetrics(clusterMetrics);
+
+ Map<TableName, Map<ServerName, List<RegionInfo>>> LoadOfAllTable =
+ (Map) mockClusterServersWithTables(clusterState);
+ List<RegionPlan> plans = loadBalancer.balanceCluster(LoadOfAllTable);
+ Set<RegionInfo> regionsMovedFromServer0 = new HashSet<>();
+ Map<ServerName, List<RegionInfo>> targetServers = new HashMap<>();
+ for (RegionPlan plan : plans) {
+ if (plan.getSource().equals(server0)) {
+ regionsMovedFromServer0.add(plan.getRegionInfo());
+ if (!targetServers.containsKey(plan.getDestination())) {
+ targetServers.put(plan.getDestination(), new ArrayList<>());
+ }
+ targetServers.get(plan.getDestination()).add(plan.getRegionInfo());
+ }
+ }
+ // Skew rebalancing should still move 5 regions from server0 to server1 to
balance the
+ // cluster (10 on server0, 0 on server1, 5 on server2 → target ~5 on
each). But the
+ // specific regions moved are not dictated by old cache data since all
regions are already
+ // well-cached on their current server.
+ assertEquals(5, regionsMovedFromServer0.size());
+ assertEquals(5, targetServers.get(server1).size());
+ }
+
+ @Test
+ public void testRegionsPartiallyCachedOnOldServerAndCurrentServer() throws
Exception {
+ // The regions are partially cached on old server (0.6) and have lower
cache on current (0.2).
+ // The balancer should move regions to server1 to fix skew, and the
cache-aware generator
+ // guides some of those moves to be the old-cached regions.
+
+ Map<ServerName, List<RegionInfo>> clusterState = new HashMap<>();
+ ServerName server0 = servers.get(0);
+ ServerName server1 = servers.get(1);
+ ServerName server2 = servers.get(2);
+
+ // Simulate that the regions previously hosted by server1 are now hosted
on server0
+ List<RegionInfo> regionsOnServer0 = randomRegions(10);
+ List<RegionInfo> regionsOnServer1 = randomRegions(0);
+ List<RegionInfo> regionsOnServer2 = randomRegions(5);
+
+ clusterState.put(server0, regionsOnServer0);
+ clusterState.put(server1, regionsOnServer1);
+ clusterState.put(server2, regionsOnServer2);
+
+ // Mock cluster metrics
+
+ // Mock 4 regions from server0 were previously hosted on server1
+ List<RegionInfo> oldCachedRegions = regionsOnServer0.subList(5,
regionsOnServer0.size() - 1);
+
+ Map<ServerName, ServerMetrics> serverMetricsMap = new TreeMap<>();
+ serverMetricsMap.put(server0,
+ mockServerMetricsWithRegionCacheInfo(regionsOnServer0, 0.2f, new
ArrayList<>(), 0, 10));
+ serverMetricsMap.put(server1,
+ mockServerMetricsWithRegionCacheInfo(regionsOnServer1, 0.0f,
oldCachedRegions, 6, 10));
+ serverMetricsMap.put(server2,
+ mockServerMetricsWithRegionCacheInfo(regionsOnServer2, 1.0f, new
ArrayList<>(), 0, 10));
+ ClusterMetrics clusterMetrics = mock(ClusterMetrics.class);
+ when(clusterMetrics.getLiveServerMetrics()).thenReturn(serverMetricsMap);
+ loadBalancer.updateClusterMetrics(clusterMetrics);
+
+ Map<TableName, Map<ServerName, List<RegionInfo>>> LoadOfAllTable =
+ (Map) mockClusterServersWithTables(clusterState);
+ List<RegionPlan> plans = loadBalancer.balanceCluster(LoadOfAllTable);
+ Set<RegionInfo> regionsMovedFromServer0 = new HashSet<>();
+ Map<ServerName, List<RegionInfo>> targetServers = new HashMap<>();
+ for (RegionPlan plan : plans) {
+ if (plan.getSource().equals(server0)) {
+ regionsMovedFromServer0.add(plan.getRegionInfo());
+ if (!targetServers.containsKey(plan.getDestination())) {
+ targetServers.put(plan.getDestination(), new ArrayList<>());
+ }
+ targetServers.get(plan.getDestination()).add(plan.getRegionInfo());
+ }
+ }
+ // Balanced state for 15 total regions on 3 servers = 5 each.
+ // server0(10) → server1(0): should move 5
+ assertEquals(5, regionsMovedFromServer0.size());
+ assertEquals(5, targetServers.get(server1).size());
+ // The cache-aware generator should move at least some old-cached regions
to server1
+ // (where they have better cache). Due to stochastic walk non-determinism,
not all 4
+ // are guaranteed to be picked over equally-viable alternatives.
+ long oldCachedOnServer1 =
+
targetServers.get(server1).stream().filter(oldCachedRegions::contains).count();
+ assertTrue(oldCachedOnServer1 > 0, "At least some old-cached regions
should move to server1");
+ }
+
@Test
public void testBalancerNotThrowNPEWhenBalancerPlansIsNull() throws
Exception {
Map<ServerName, List<RegionInfo>> clusterState = new HashMap<>();
@@ -138,12 +661,12 @@ public class TestCacheAwareLoadBalancer extends
BalancerTestBase {
// Mock cluster metrics
Map<ServerName, ServerMetrics> serverMetricsMap = new TreeMap<>();
- serverMetricsMap.put(server0,
mockServerMetricsWithRegionCacheInfo(server0, regionsOnServer0,
- 0.0f, new ArrayList<>(), 0, 10));
- serverMetricsMap.put(server1,
mockServerMetricsWithRegionCacheInfo(server1, regionsOnServer1,
- 0.0f, new ArrayList<>(), 0, 10));
- serverMetricsMap.put(server2,
mockServerMetricsWithRegionCacheInfo(server2, regionsOnServer2,
- 0.0f, new ArrayList<>(), 0, 10));
+ serverMetricsMap.put(server0,
+ mockServerMetricsWithRegionCacheInfo(regionsOnServer0, 0.0f, new
ArrayList<>(), 0, 10));
+ serverMetricsMap.put(server1,
+ mockServerMetricsWithRegionCacheInfo(regionsOnServer1, 0.0f, new
ArrayList<>(), 0, 10));
+ serverMetricsMap.put(server2,
+ mockServerMetricsWithRegionCacheInfo(regionsOnServer2, 0.0f, new
ArrayList<>(), 0, 10));
ClusterMetrics clusterMetrics = mock(ClusterMetrics.class);
when(clusterMetrics.getLiveServerMetrics()).thenReturn(serverMetricsMap);
@@ -158,4 +681,5 @@ public class TestCacheAwareLoadBalancer extends
BalancerTestBase {
fail("NPE should not be thrown");
}
}
+
}
diff --git
a/hbase-server/src/test/java/org/apache/hadoop/hbase/master/balancer/TestRSGroupBasedLoadBalancerWithCacheAwareLoadBalancerAsInternal.java
b/hbase-server/src/test/java/org/apache/hadoop/hbase/master/balancer/TestRSGroupBasedLoadBalancerWithCacheAwareLoadBalancerAsInternal.java
index 4ed6c84a026..c60290e0c04 100644
---
a/hbase-server/src/test/java/org/apache/hadoop/hbase/master/balancer/TestRSGroupBasedLoadBalancerWithCacheAwareLoadBalancerAsInternal.java
+++
b/hbase-server/src/test/java/org/apache/hadoop/hbase/master/balancer/TestRSGroupBasedLoadBalancerWithCacheAwareLoadBalancerAsInternal.java
@@ -18,7 +18,6 @@
package org.apache.hadoop.hbase.master.balancer;
import static org.junit.jupiter.api.Assertions.assertEquals;
-import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.Mockito.mock;
@@ -244,10 +243,17 @@ public class
TestRSGroupBasedLoadBalancerWithCacheAwareLoadBalancerAsInternal
targetServers.get(plan.getDestination()).add(plan.getRegionInfo());
}
}
- // should move 5 regions from server0 to server1
- assertEquals(5, regionsMovedFromServer0.size());
- assertEquals(5, targetServers.get(server1).size());
- assertTrue(targetServers.get(server1).containsAll(oldCachedRegions));
+ // should move regions from server0 to server1 (old-cached regions should
be among them)
+ assertTrue(regionsMovedFromServer0.size() >= 4);
+ assertNotNull(targetServers.get(server1));
+ int oldCachedOnServer1 = 0;
+ for (RegionInfo ri : oldCachedRegions) {
+ if (targetServers.get(server1).contains(ri)) {
+ oldCachedOnServer1++;
+ }
+ }
+ assertTrue(oldCachedOnServer1 > 0,
+ "Expected old-cached regions to move to server1, got " +
oldCachedOnServer1);
}
@Test
@@ -415,22 +421,32 @@ public class
TestRSGroupBasedLoadBalancerWithCacheAwareLoadBalancerAsInternal
List<RegionPlan> plans = loadBalancer.balanceCluster(LoadOfAllTable);
LOG.debug("plans size: {}", plans.size());
LOG.debug("plans: {}", plans);
- LOG.debug("server1 name: {}", server1.getServerName());
- // assert the plans are in descending order from the most cached to the
least cached
- int highCacheCount = 0;
+ // Plans are sorted by cache ratio on destination (descending). Verify
ordering
+ // and that at least some old-cached regions are moved to their cached
servers.
+ float prevRatio = Float.MAX_VALUE;
+ int oldCached1Count = 0;
+ int oldCached2Count = 0;
for (RegionPlan plan : plans) {
LOG.debug("plan region: {}, target server: {}",
plan.getRegionInfo().getEncodedName(),
plan.getDestination().getServerName());
- if (highCacheCount < 5) {
- LOG.debug("Count: {}", highCacheCount);
- assertTrue(oldCachedRegions1.contains(plan.getRegionInfo()));
- assertFalse(oldCachedRegions2.contains(plan.getRegionInfo()));
- highCacheCount++;
- } else {
- assertTrue(oldCachedRegions2.contains(plan.getRegionInfo()));
- assertFalse(oldCachedRegions1.contains(plan.getRegionInfo()));
+ float ratio = 0f;
+ if (
+ oldCachedRegions1.contains(plan.getRegionInfo()) &&
server1.equals(plan.getDestination())
+ ) {
+ ratio = 1.0f;
+ oldCached1Count++;
+ } else if (
+ oldCachedRegions2.contains(plan.getRegionInfo()) &&
server2.equals(plan.getDestination())
+ ) {
+ ratio = 0.8f;
+ oldCached2Count++;
}
+ assertTrue(ratio <= prevRatio,
+ "Plans should be sorted by cache ratio on destination (descending)");
+ prevRatio = ratio;
}
+ assertTrue(oldCached1Count > 0, "Some old-cached regions should move to
server1");
+ assertTrue(oldCached2Count > 0, "Some old-cached regions should move to
server2");
}
@Test
@@ -481,10 +497,18 @@ public class
TestRSGroupBasedLoadBalancerWithCacheAwareLoadBalancerAsInternal
targetServers.get(plan.getDestination()).add(plan.getRegionInfo());
}
}
- // should move 5 regions from server0 to server1
- assertEquals(5, regionsMovedFromServer0.size());
- assertEquals(5, targetServers.get(server1).size());
- assertTrue(targetServers.get(server1).containsAll(oldCachedRegions));
+ // should move regions from server0 to server1 (old-cached regions should
be among them)
+ assertTrue(regionsMovedFromServer0.size() >= 4);
+ assertNotNull(targetServers.get(server1));
+ assertTrue(targetServers.get(server1).size() >= 4);
+ int oldCachedOnServer1 = 0;
+ for (RegionInfo ri : oldCachedRegions) {
+ if (targetServers.get(server1).contains(ri)) {
+ oldCachedOnServer1++;
+ }
+ }
+ assertTrue(oldCachedOnServer1 > 0,
+ "Expected most old-cached regions to move to server1, got " +
oldCachedOnServer1);
}
@Test
@@ -535,10 +559,11 @@ public class
TestRSGroupBasedLoadBalancerWithCacheAwareLoadBalancerAsInternal
targetServers.get(plan.getDestination()).add(plan.getRegionInfo());
}
}
- // should move 5 regions from server0 to server1
+ // should move 5 regions from server0 to server1 to balance the cluster,
but the specific
+ // regions moved are not dictated by old cache data since all regions are
already well-cached
+ // on their current server (currentCacheRatio >= ratioThreshold)
assertEquals(5, regionsMovedFromServer0.size());
assertEquals(5, targetServers.get(server1).size());
- assertTrue(targetServers.get(server1).containsAll(oldCachedRegions));
}
@Test
@@ -589,9 +614,14 @@ public class
TestRSGroupBasedLoadBalancerWithCacheAwareLoadBalancerAsInternal
targetServers.get(plan.getDestination()).add(plan.getRegionInfo());
}
}
+ // server0(10) → server1(0): should move 5
assertEquals(5, regionsMovedFromServer0.size());
assertEquals(5, targetServers.get(server1).size());
- assertTrue(targetServers.get(server1).containsAll(oldCachedRegions));
+ // The cache-aware generator should move at least some old-cached regions
to server1.
+ // Due to stochastic walk non-determinism, not all are guaranteed.
+ long oldCachedOnServer1 =
+
targetServers.get(server1).stream().filter(oldCachedRegions::contains).count();
+ assertTrue(oldCachedOnServer1 > 0, "At least some old-cached regions
should move to server1");
}
@Timeout(60)
@@ -626,14 +656,14 @@ public class
TestRSGroupBasedLoadBalancerWithCacheAwareLoadBalancerAsInternal
clusterState.put(server1, regionsOnServer1);
clusterState.put(server2, regionsOnServer2);
- // Mock metrics: NO cache info for any region = all will be throttled
+ // Mock metrics: regions have moderate cache ratio so throttle applies on
move
Map<ServerName, ServerMetrics> serverMetricsMap = new TreeMap<>();
serverMetricsMap.put(server0,
mockServerMetricsWithRegionCacheInfo(server0, regionsOnServer0,
- 0.0f, new ArrayList<>(), 0, 10));
+ 0.5f, new ArrayList<>(), 0, 10));
serverMetricsMap.put(server1,
mockServerMetricsWithRegionCacheInfo(server1, regionsOnServer1,
- 0.0f, new ArrayList<>(), 0, 10));
+ 0.5f, new ArrayList<>(), 0, 10));
serverMetricsMap.put(server2,
mockServerMetricsWithRegionCacheInfo(server2, regionsOnServer2,
- 0.0f, new ArrayList<>(), 0, 10));
+ 0.5f, new ArrayList<>(), 0, 10));
ClusterMetrics clusterMetrics = mock(ClusterMetrics.class);
when(clusterMetrics.getLiveServerMetrics()).thenReturn(serverMetricsMap);