001/* 002 * Licensed to the Apache Software Foundation (ASF) under one 003 * or more contributor license agreements. See the NOTICE file 004 * distributed with this work for additional information 005 * regarding copyright ownership. The ASF licenses this file 006 * to you under the Apache License, Version 2.0 (the 007 * "License"); you may not use this file except in compliance 008 * with the License. You may obtain a copy of the License at 009 * 010 * http://www.apache.org/licenses/LICENSE-2.0 011 * 012 * Unless required by applicable law or agreed to in writing, software 013 * distributed under the License is distributed on an "AS IS" BASIS, 014 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. 015 * See the License for the specific language governing permissions and 016 * limitations under the License. 017 */ 018package org.apache.hadoop.hbase.master.balancer; 019 020import static org.junit.jupiter.api.Assertions.assertEquals; 021import static org.junit.jupiter.api.Assertions.assertNotNull; 022import static org.junit.jupiter.api.Assertions.assertNull; 023import static org.junit.jupiter.api.Assertions.assertTrue; 024import static org.junit.jupiter.api.Assertions.fail; 025import static org.mockito.Mockito.mock; 026import static org.mockito.Mockito.when; 027 028import java.util.ArrayList; 029import java.util.HashMap; 030import java.util.HashSet; 031import java.util.List; 032import java.util.Map; 033import java.util.Random; 034import java.util.Set; 035import java.util.TreeMap; 036import java.util.concurrent.ThreadLocalRandom; 037import org.apache.hadoop.conf.Configuration; 038import org.apache.hadoop.hbase.ClusterMetrics; 039import org.apache.hadoop.hbase.HBaseConfiguration; 040import org.apache.hadoop.hbase.HConstants; 041import org.apache.hadoop.hbase.RegionMetrics; 042import org.apache.hadoop.hbase.ServerMetrics; 043import org.apache.hadoop.hbase.ServerName; 044import org.apache.hadoop.hbase.Size; 045import org.apache.hadoop.hbase.TableName; 046import org.apache.hadoop.hbase.client.RegionInfo; 047import org.apache.hadoop.hbase.client.TableDescriptor; 048import org.apache.hadoop.hbase.client.TableDescriptorBuilder; 049import org.apache.hadoop.hbase.master.RegionPlan; 050import org.apache.hadoop.hbase.testclassification.LargeTests; 051import org.apache.hadoop.hbase.util.Bytes; 052import org.apache.hadoop.hbase.util.Pair; 053import org.junit.jupiter.api.BeforeAll; 054import org.junit.jupiter.api.Tag; 055import org.junit.jupiter.api.Test; 056import org.slf4j.Logger; 057import org.slf4j.LoggerFactory; 058 059import org.apache.hbase.thirdparty.com.google.common.collect.Lists; 060 061@Tag(LargeTests.TAG) 062public class TestCacheAwareLoadBalancer extends BalancerTestBase { 063 064 private static final Logger LOG = LoggerFactory.getLogger(TestCacheAwareLoadBalancer.class); 065 066 private static CacheAwareLoadBalancer loadBalancer; 067 068 static List<ServerName> servers; 069 070 static List<TableDescriptor> tableDescs; 071 072 static TableName[] tables = new TableName[] { TableName.valueOf("dt1"), TableName.valueOf("dt2"), 073 TableName.valueOf("dt3"), TableName.valueOf("dt4") }; 074 075 private static List<ServerName> generateServers(int numServers) { 076 List<ServerName> servers = new ArrayList<>(numServers); 077 Random rand = ThreadLocalRandom.current(); 078 for (int i = 0; i < numServers; i++) { 079 String host = "server" + rand.nextInt(100000); 080 int port = rand.nextInt(60000); 081 servers.add(ServerName.valueOf(host, port, -1)); 082 } 083 return servers; 084 } 085 086 private static List<TableDescriptor> constructTableDesc(boolean hasBogusTable) { 087 List<TableDescriptor> tds = Lists.newArrayList(); 088 for (int i = 0; i < tables.length; i++) { 089 TableDescriptor htd = TableDescriptorBuilder.newBuilder(tables[i]).build(); 090 tds.add(htd); 091 } 092 return tds; 093 } 094 095 private ServerMetrics mockServerMetricsWithRegionCacheInfo(List<RegionInfo> regionsOnServer, 096 float currentCacheRatio, List<RegionInfo> oldRegionCacheInfo, int oldRegionCachedSize, 097 int regionSize) { 098 ServerMetrics serverMetrics = mock(ServerMetrics.class); 099 Map<byte[], RegionMetrics> regionLoadMap = new TreeMap<>(Bytes.BYTES_COMPARATOR); 100 for (RegionInfo info : regionsOnServer) { 101 RegionMetrics rl = mock(RegionMetrics.class); 102 when(rl.getReadRequestCount()).thenReturn(0L); 103 when(rl.getWriteRequestCount()).thenReturn(0L); 104 when(rl.getMemStoreSize()).thenReturn(Size.ZERO); 105 when(rl.getStoreFileSize()).thenReturn(Size.ZERO); 106 when(rl.getCurrentRegionCachedRatio()).thenReturn(currentCacheRatio); 107 when(rl.getRegionSizeMB()).thenReturn(new Size(regionSize, Size.Unit.MEGABYTE)); 108 regionLoadMap.put(info.getRegionName(), rl); 109 } 110 when(serverMetrics.getRegionMetrics()).thenReturn(regionLoadMap); 111 Map<String, Integer> oldCacheRatioMap = new HashMap<>(); 112 for (RegionInfo info : oldRegionCacheInfo) { 113 oldCacheRatioMap.put(info.getEncodedName(), oldRegionCachedSize); 114 } 115 when(serverMetrics.getRegionCachedInfo()).thenReturn(oldCacheRatioMap); 116 when(serverMetrics.getCacheFreeSize()).thenReturn(100L * 1024 * 1024 * 1024); 117 return serverMetrics; 118 } 119 120 @BeforeAll 121 public static void beforeAllTests() throws Exception { 122 servers = generateServers(3); 123 tableDescs = constructTableDesc(false); 124 Configuration conf = HBaseConfiguration.create(); 125 conf.set(HConstants.BUCKET_CACHE_PERSISTENT_PATH_KEY, "prefetch_file_list"); 126 conf.setFloat(HConstants.BUCKET_CACHE_SIZE_KEY, 10); 127 loadBalancer = new CacheAwareLoadBalancer(); 128 loadBalancer.setClusterInfoProvider(new DummyClusterInfoProvider(conf)); 129 loadBalancer.loadConf(conf); 130 } 131 132 @Test 133 public void testRegionsNotCachedOnOldServerAndCurrentServer() throws Exception { 134 // The regions are not cached on old server as well as the current server. This causes 135 // skewness in the region allocation which should be fixed by the balancer 136 137 Map<ServerName, List<RegionInfo>> clusterState = new HashMap<>(); 138 ServerName server0 = servers.get(0); 139 ServerName server1 = servers.get(1); 140 ServerName server2 = servers.get(2); 141 142 // Simulate that the regions previously hosted by server1 are now hosted on server0 143 List<RegionInfo> regionsOnServer0 = randomRegions(10); 144 List<RegionInfo> regionsOnServer1 = randomRegions(0); 145 List<RegionInfo> regionsOnServer2 = randomRegions(5); 146 147 clusterState.put(server0, regionsOnServer0); 148 clusterState.put(server1, regionsOnServer1); 149 clusterState.put(server2, regionsOnServer2); 150 151 // Mock cluster metrics — give only server1 free cache so moves are directed there 152 Map<ServerName, ServerMetrics> serverMetricsMap = new TreeMap<>(); 153 ServerMetrics sm0 = 154 mockServerMetricsWithRegionCacheInfo(regionsOnServer0, 0.0f, new ArrayList<>(), 0, 10); 155 when(sm0.getCacheFreeSize()).thenReturn(0L); 156 ServerMetrics sm1 = 157 mockServerMetricsWithRegionCacheInfo(regionsOnServer1, 0.0f, new ArrayList<>(), 0, 10); 158 ServerMetrics sm2 = 159 mockServerMetricsWithRegionCacheInfo(regionsOnServer2, 0.0f, new ArrayList<>(), 0, 10); 160 when(sm2.getCacheFreeSize()).thenReturn(0L); 161 serverMetricsMap.put(server0, sm0); 162 serverMetricsMap.put(server1, sm1); 163 serverMetricsMap.put(server2, sm2); 164 ClusterMetrics clusterMetrics = mock(ClusterMetrics.class); 165 when(clusterMetrics.getLiveServerMetrics()).thenReturn(serverMetricsMap); 166 loadBalancer.updateClusterMetrics(clusterMetrics); 167 168 Map<TableName, Map<ServerName, List<RegionInfo>>> LoadOfAllTable = 169 (Map) mockClusterServersWithTables(clusterState); 170 List<RegionPlan> plans = loadBalancer.balanceCluster(LoadOfAllTable); 171 Set<RegionInfo> regionsMovedFromServer0 = new HashSet<>(); 172 Map<ServerName, List<RegionInfo>> targetServers = new HashMap<>(); 173 for (RegionPlan plan : plans) { 174 if (plan.getSource().equals(server0)) { 175 regionsMovedFromServer0.add(plan.getRegionInfo()); 176 if (!targetServers.containsKey(plan.getDestination())) { 177 targetServers.put(plan.getDestination(), new ArrayList<>()); 178 } 179 targetServers.get(plan.getDestination()).add(plan.getRegionInfo()); 180 } 181 } 182 // should move at least 5 regions from server0 to balance cluster (10/0/5 -> ~5/5/5) 183 assertTrue(regionsMovedFromServer0.size() >= 5, 184 "Expected at least 5 moves from server0, got " + regionsMovedFromServer0.size()); 185 } 186 187 /** 188 * Regions on the overloaded RS report low block-cache ratio; no RS reports prefetch/historical 189 * cache for those regions (so {@link CacheAwareLoadBalancer.CacheAwareCandidateGenerator} has no 190 * "old server" to prefer). Another RS has ample free block cache. The balancer should still emit 191 * plans that shed load from the hot RS onto the idle RS with spare cache capacity. 192 */ 193 @Test 194 public void testLowCacheRatioNoHistoricalCacheRelocatesWhenTargetHasFreeBlockCache() 195 throws Exception { 196 Map<ServerName, List<RegionInfo>> clusterState = new HashMap<>(); 197 ServerName server0 = servers.get(0); 198 ServerName server1 = servers.get(1); 199 ServerName server2 = servers.get(2); 200 201 List<RegionInfo> regionsOnServer0 = randomRegions(10); 202 List<RegionInfo> regionsOnServer1 = randomRegions(0); 203 List<RegionInfo> regionsOnServer2 = randomRegions(5); 204 205 clusterState.put(server0, regionsOnServer0); 206 clusterState.put(server1, regionsOnServer1); 207 clusterState.put(server2, regionsOnServer2); 208 209 // Below LOW_CACHE_RATIO_FOR_RELOCATION_DEFAULT (0.35); 210 ServerMetrics sm0 = 211 mockServerMetricsWithRegionCacheInfo(regionsOnServer0, 0.1f, new ArrayList<>(), 0, 10); 212 when(sm0.getCacheFreeSize()).thenReturn(0L); 213 ServerMetrics sm1 = 214 mockServerMetricsWithRegionCacheInfo(regionsOnServer1, 0.0f, new ArrayList<>(), 0, 10); 215 // Simulates 1GB free cache space on server1 216 when(sm1.getCacheFreeSize()).thenReturn(1024L * 1024 * 1024); 217 ServerMetrics sm2 = 218 mockServerMetricsWithRegionCacheInfo(regionsOnServer2, 1.0f, new ArrayList<>(), 0, 10); 219 when(sm2.getCacheFreeSize()).thenReturn(0L); 220 221 Map<ServerName, ServerMetrics> serverMetricsMap = new TreeMap<>(); 222 serverMetricsMap.put(server0, sm0); 223 serverMetricsMap.put(server1, sm1); 224 serverMetricsMap.put(server2, sm2); 225 ClusterMetrics clusterMetrics = mock(ClusterMetrics.class); 226 when(clusterMetrics.getLiveServerMetrics()).thenReturn(serverMetricsMap); 227 loadBalancer.updateClusterMetrics(clusterMetrics); 228 229 assertTrue(loadBalancer.regionCacheRatioOnOldServerMap.isEmpty()); 230 231 Map<TableName, Map<ServerName, List<RegionInfo>>> loadOfAllTable = 232 (Map) mockClusterServersWithTables(clusterState); 233 List<RegionPlan> plans = loadBalancer.balanceCluster(loadOfAllTable); 234 assertNotNull(plans); 235 236 Set<RegionInfo> regionsMovedFromServer0 = new HashSet<>(); 237 Map<ServerName, List<RegionInfo>> targetServers = new HashMap<>(); 238 for (RegionPlan plan : plans) { 239 if (plan.getSource().equals(server0)) { 240 regionsMovedFromServer0.add(plan.getRegionInfo()); 241 if (!targetServers.containsKey(plan.getDestination())) { 242 targetServers.put(plan.getDestination(), new ArrayList<>()); 243 } 244 targetServers.get(plan.getDestination()).add(plan.getRegionInfo()); 245 } 246 } 247 assertEquals(5, regionsMovedFromServer0.size()); 248 assertNotNull(targetServers.get(server1)); 249 assertEquals(5, targetServers.get(server1).size()); 250 } 251 252 @Test 253 public void testRegionsPartiallyCachedOnOldServerAndNotCachedOnCurrentServer() throws Exception { 254 // The regions are partially cached on old server but not cached on the current server 255 256 Map<ServerName, List<RegionInfo>> clusterState = new HashMap<>(); 257 ServerName server0 = servers.get(0); 258 ServerName server1 = servers.get(1); 259 ServerName server2 = servers.get(2); 260 261 // Simulate that the regions previously hosted by server1 are now hosted on server0 262 List<RegionInfo> regionsOnServer0 = randomRegions(10); 263 List<RegionInfo> regionsOnServer1 = randomRegions(0); 264 List<RegionInfo> regionsOnServer2 = randomRegions(5); 265 266 clusterState.put(server0, regionsOnServer0); 267 clusterState.put(server1, regionsOnServer1); 268 clusterState.put(server2, regionsOnServer2); 269 270 // Mock cluster metrics 271 272 // Mock 5 regions from server0 were previously hosted on server1 273 List<RegionInfo> oldCachedRegions = regionsOnServer0.subList(5, regionsOnServer0.size()); 274 275 Map<ServerName, ServerMetrics> serverMetricsMap = new TreeMap<>(); 276 serverMetricsMap.put(server0, 277 mockServerMetricsWithRegionCacheInfo(regionsOnServer0, 0.0f, new ArrayList<>(), 0, 10)); 278 serverMetricsMap.put(server1, 279 mockServerMetricsWithRegionCacheInfo(regionsOnServer1, 0.0f, oldCachedRegions, 6, 10)); 280 serverMetricsMap.put(server2, 281 mockServerMetricsWithRegionCacheInfo(regionsOnServer2, 0.0f, new ArrayList<>(), 0, 10)); 282 ClusterMetrics clusterMetrics = mock(ClusterMetrics.class); 283 when(clusterMetrics.getLiveServerMetrics()).thenReturn(serverMetricsMap); 284 loadBalancer.updateClusterMetrics(clusterMetrics); 285 286 Map<TableName, Map<ServerName, List<RegionInfo>>> LoadOfAllTable = 287 (Map) mockClusterServersWithTables(clusterState); 288 List<RegionPlan> plans = loadBalancer.balanceCluster(LoadOfAllTable); 289 Set<RegionInfo> regionsMovedFromServer0 = new HashSet<>(); 290 Map<ServerName, List<RegionInfo>> targetServers = new HashMap<>(); 291 for (RegionPlan plan : plans) { 292 if (plan.getSource().equals(server0)) { 293 regionsMovedFromServer0.add(plan.getRegionInfo()); 294 if (!targetServers.containsKey(plan.getDestination())) { 295 targetServers.put(plan.getDestination(), new ArrayList<>()); 296 } 297 targetServers.get(plan.getDestination()).add(plan.getRegionInfo()); 298 } 299 } 300 // should move regions from server0 to server1 (old-cached regions should be among them) 301 assertEquals(5, regionsMovedFromServer0.size()); 302 assertNotNull(targetServers.get(server1)); 303 int oldCachedOnServer1 = 0; 304 for (RegionInfo ri : oldCachedRegions) { 305 if (targetServers.get(server1).contains(ri)) { 306 oldCachedOnServer1++; 307 } 308 } 309 assertTrue(oldCachedOnServer1 > 0, 310 "Expected old-cached regions to move to server1, got " + oldCachedOnServer1); 311 } 312 313 @Test 314 public void testThrottlingRegionBeyondThreshold() throws Exception { 315 Configuration conf = HBaseConfiguration.create(); 316 CacheAwareLoadBalancer balancer = new CacheAwareLoadBalancer(); 317 balancer.loadConf(conf); 318 balancer.setClusterInfoProvider(new DummyClusterInfoProvider(conf)); 319 balancer.initialize(); 320 ServerName server0 = servers.get(0); 321 ServerName server1 = servers.get(1); 322 Pair<ServerName, Float> regionRatio = new Pair<>(); 323 regionRatio.setFirst(server0); 324 regionRatio.setSecond(1.0f); 325 balancer.regionCacheRatioOnOldServerMap.put("region1", regionRatio); 326 RegionInfo mockedInfo = mock(RegionInfo.class); 327 when(mockedInfo.getEncodedName()).thenReturn("region1"); 328 RegionPlan plan = new RegionPlan(mockedInfo, server1, server0); 329 assertEquals(0L, balancer.getThrottleDurationMs(plan)); 330 } 331 332 @Test 333 public void testThrottlingRegionBelowThreshold() throws Exception { 334 Configuration conf = HBaseConfiguration.create(); 335 conf.setLong(CacheAwareLoadBalancer.MOVE_THROTTLING, 100); 336 CacheAwareLoadBalancer balancer = new CacheAwareLoadBalancer(); 337 balancer.loadConf(conf); 338 balancer.setClusterInfoProvider(new DummyClusterInfoProvider(conf)); 339 balancer.initialize(); 340 ServerName server0 = servers.get(0); 341 ServerName server1 = servers.get(1); 342 Pair<ServerName, Float> regionRatio = new Pair<>(); 343 regionRatio.setFirst(server0); 344 regionRatio.setSecond(0.1f); 345 balancer.regionCacheRatioOnOldServerMap.put("region1", regionRatio); 346 RegionInfo mockedInfo = mock(RegionInfo.class); 347 when(mockedInfo.getEncodedName()).thenReturn("region1"); 348 RegionPlan plan = new RegionPlan(mockedInfo, server1, server0); 349 assertEquals(100L, balancer.getThrottleDurationMs(plan)); 350 } 351 352 @Test 353 public void testThrottlingCacheRatioUnknownOnTarget() throws Exception { 354 Configuration conf = HBaseConfiguration.create(); 355 conf.setLong(CacheAwareLoadBalancer.MOVE_THROTTLING, 100); 356 CacheAwareLoadBalancer balancer = new CacheAwareLoadBalancer(); 357 balancer.loadConf(conf); 358 balancer.setClusterInfoProvider(new DummyClusterInfoProvider(conf)); 359 balancer.initialize(); 360 ServerName server0 = servers.get(0); 361 ServerName server1 = servers.get(1); 362 ServerName server3 = servers.get(2); 363 // setting region cache ratio 100% on server 3, though this is not the target in the region plan 364 Pair<ServerName, Float> regionRatio = new Pair<>(); 365 regionRatio.setFirst(server3); 366 regionRatio.setSecond(1.0f); 367 balancer.regionCacheRatioOnOldServerMap.put("region1", regionRatio); 368 RegionInfo mockedInfo = mock(RegionInfo.class); 369 when(mockedInfo.getEncodedName()).thenReturn("region1"); 370 RegionPlan plan = new RegionPlan(mockedInfo, server1, server0); 371 assertEquals(100L, balancer.getThrottleDurationMs(plan)); 372 } 373 374 @Test 375 public void testThrottlingCacheRatioUnknownForRegion() throws Exception { 376 Configuration conf = HBaseConfiguration.create(); 377 conf.setLong(CacheAwareLoadBalancer.MOVE_THROTTLING, 100); 378 CacheAwareLoadBalancer balancer = new CacheAwareLoadBalancer(); 379 balancer.loadConf(conf); 380 balancer.setClusterInfoProvider(new DummyClusterInfoProvider(conf)); 381 balancer.initialize(); 382 ServerName server0 = servers.get(0); 383 ServerName server1 = servers.get(1); 384 ServerName server3 = servers.get(2); 385 // No cache ratio available for region1 386 RegionInfo mockedInfo = mock(RegionInfo.class); 387 when(mockedInfo.getEncodedName()).thenReturn("region1"); 388 RegionPlan plan = new RegionPlan(mockedInfo, server1, server0); 389 assertEquals(100L, balancer.getThrottleDurationMs(plan)); 390 } 391 392 @Test 393 public void testRegionPlansSortedByCacheRatioOnTarget() throws Exception { 394 // The regions are fully cached on old server 395 396 Map<ServerName, List<RegionInfo>> clusterState = new HashMap<>(); 397 ServerName server0 = servers.get(0); 398 ServerName server1 = servers.get(1); 399 ServerName server2 = servers.get(2); 400 401 // Simulate on RS with all regions, and two RSes with no regions 402 List<RegionInfo> regionsOnServer0 = randomRegions(15); 403 List<RegionInfo> regionsOnServer1 = randomRegions(0); 404 List<RegionInfo> regionsOnServer2 = randomRegions(0); 405 406 clusterState.put(server0, regionsOnServer0); 407 clusterState.put(server1, regionsOnServer1); 408 clusterState.put(server2, regionsOnServer2); 409 410 // Mock cluster metrics 411 // Mock 5 regions from server0 were previously hosted on server1 412 List<RegionInfo> oldCachedRegions1 = regionsOnServer0.subList(5, 10); 413 List<RegionInfo> oldCachedRegions2 = regionsOnServer0.subList(10, regionsOnServer0.size()); 414 Map<ServerName, ServerMetrics> serverMetricsMap = new TreeMap<>(); 415 // mock server metrics to set cache ratio as 0 in the RS 0 416 serverMetricsMap.put(server0, 417 mockServerMetricsWithRegionCacheInfo(regionsOnServer0, 0.0f, new ArrayList<>(), 0, 10)); 418 // mock server metrics to set cache ratio as 1 in the RS 1 419 serverMetricsMap.put(server1, 420 mockServerMetricsWithRegionCacheInfo(regionsOnServer1, 0.0f, oldCachedRegions1, 10, 10)); 421 // mock server metrics to set cache ratio as .8 in the RS 2 422 serverMetricsMap.put(server2, 423 mockServerMetricsWithRegionCacheInfo(regionsOnServer2, 0.0f, oldCachedRegions2, 8, 10)); 424 ClusterMetrics clusterMetrics = mock(ClusterMetrics.class); 425 when(clusterMetrics.getLiveServerMetrics()).thenReturn(serverMetricsMap); 426 loadBalancer.updateClusterMetrics(clusterMetrics); 427 428 Map<TableName, Map<ServerName, List<RegionInfo>>> LoadOfAllTable = 429 (Map) mockClusterServersWithTables(clusterState); 430 List<RegionPlan> plans = loadBalancer.balanceCluster(LoadOfAllTable); 431 LOG.debug("plans size: {}", plans.size()); 432 LOG.debug("plans: {}", plans); 433 // Plans are sorted by cache ratio on destination (descending). Verify that plans 434 // for old-cached regions going to their correct servers appear before plans with no 435 // cache data on destination. 436 float prevRatio = Float.MAX_VALUE; 437 int oldCached1Count = 0; 438 int oldCached2Count = 0; 439 for (RegionPlan plan : plans) { 440 LOG.debug("plan region: {}, target server: {}", plan.getRegionInfo().getEncodedName(), 441 plan.getDestination().getServerName()); 442 float ratio = 0f; 443 if ( 444 oldCachedRegions1.contains(plan.getRegionInfo()) && server1.equals(plan.getDestination()) 445 ) { 446 ratio = 1.0f; 447 oldCached1Count++; 448 } else if ( 449 oldCachedRegions2.contains(plan.getRegionInfo()) && server2.equals(plan.getDestination()) 450 ) { 451 ratio = 0.8f; 452 oldCached2Count++; 453 } 454 assertTrue(ratio <= prevRatio, 455 "Plans should be sorted by cache ratio on destination (descending)"); 456 prevRatio = ratio; 457 } 458 // The cache-aware generator should move at least some old-cached regions to their 459 // cached servers. Exact count depends on stochastic walk order. 460 assertTrue(oldCached1Count > 0, "Some old-cached regions should move to server1"); 461 assertTrue(oldCached2Count > 0, "Some old-cached regions should move to server2"); 462 463 } 464 465 @Test 466 public void testRegionsFullyCachedOnOldServerAndNotCachedOnCurrentServers() throws Exception { 467 // The regions are fully cached on old server 468 469 Map<ServerName, List<RegionInfo>> clusterState = new HashMap<>(); 470 ServerName server0 = servers.get(0); 471 ServerName server1 = servers.get(1); 472 ServerName server2 = servers.get(2); 473 474 // Simulate that the regions previously hosted by server1 are now hosted on server0 475 List<RegionInfo> regionsOnServer0 = randomRegions(10); 476 List<RegionInfo> regionsOnServer1 = randomRegions(0); 477 List<RegionInfo> regionsOnServer2 = randomRegions(5); 478 479 clusterState.put(server0, regionsOnServer0); 480 clusterState.put(server1, regionsOnServer1); 481 clusterState.put(server2, regionsOnServer2); 482 483 // Mock cluster metrics 484 485 // Mock 5 regions from server0 were previously hosted on server1 486 List<RegionInfo> oldCachedRegions = regionsOnServer0.subList(5, regionsOnServer0.size() - 1); 487 488 Map<ServerName, ServerMetrics> serverMetricsMap = new TreeMap<>(); 489 serverMetricsMap.put(server0, 490 mockServerMetricsWithRegionCacheInfo(regionsOnServer0, 0.0f, new ArrayList<>(), 0, 10)); 491 serverMetricsMap.put(server1, 492 mockServerMetricsWithRegionCacheInfo(regionsOnServer1, 0.0f, oldCachedRegions, 10, 10)); 493 serverMetricsMap.put(server2, 494 mockServerMetricsWithRegionCacheInfo(regionsOnServer2, 0.0f, new ArrayList<>(), 0, 10)); 495 ClusterMetrics clusterMetrics = mock(ClusterMetrics.class); 496 when(clusterMetrics.getLiveServerMetrics()).thenReturn(serverMetricsMap); 497 loadBalancer.updateClusterMetrics(clusterMetrics); 498 499 Map<TableName, Map<ServerName, List<RegionInfo>>> LoadOfAllTable = 500 (Map) mockClusterServersWithTables(clusterState); 501 List<RegionPlan> plans = loadBalancer.balanceCluster(LoadOfAllTable); 502 Set<RegionInfo> regionsMovedFromServer0 = new HashSet<>(); 503 Map<ServerName, List<RegionInfo>> targetServers = new HashMap<>(); 504 for (RegionPlan plan : plans) { 505 if (plan.getSource().equals(server0)) { 506 regionsMovedFromServer0.add(plan.getRegionInfo()); 507 if (!targetServers.containsKey(plan.getDestination())) { 508 targetServers.put(plan.getDestination(), new ArrayList<>()); 509 } 510 targetServers.get(plan.getDestination()).add(plan.getRegionInfo()); 511 } 512 } 513 // should move regions from server0 to server1 (old-cached regions should be among them) 514 assertTrue(regionsMovedFromServer0.size() >= 4); 515 assertNotNull(targetServers.get(server1)); 516 assertTrue(targetServers.get(server1).size() >= 4); 517 int oldCachedOnServer1 = 0; 518 for (RegionInfo ri : oldCachedRegions) { 519 if (targetServers.get(server1).contains(ri)) { 520 oldCachedOnServer1++; 521 } 522 } 523 assertTrue(oldCachedOnServer1 > 0, 524 "Expected most old-cached regions to move to server1, got " + oldCachedOnServer1); 525 } 526 527 @Test 528 public void testRegionsFullyCachedOnOldAndCurrentServers() throws Exception { 529 // When regions are fully cached on BOTH the old and current server, the balancer should 530 // NOT disrupt them by moving them to the old server based on potentially stale historical 531 // cache data. Instead, it should still rebalance for skew. 532 533 Map<ServerName, List<RegionInfo>> clusterState = new HashMap<>(); 534 ServerName server0 = servers.get(0); 535 ServerName server1 = servers.get(1); 536 ServerName server2 = servers.get(2); 537 538 // Simulate that the regions previously hosted by server1 are now hosted on server0 539 List<RegionInfo> regionsOnServer0 = randomRegions(10); 540 List<RegionInfo> regionsOnServer1 = randomRegions(0); 541 List<RegionInfo> regionsOnServer2 = randomRegions(5); 542 543 clusterState.put(server0, regionsOnServer0); 544 clusterState.put(server1, regionsOnServer1); 545 clusterState.put(server2, regionsOnServer2); 546 547 // Mock cluster metrics 548 549 // Mock 4 regions from server0 were previously hosted on server1 550 List<RegionInfo> oldCachedRegions = regionsOnServer0.subList(5, regionsOnServer0.size() - 1); 551 552 Map<ServerName, ServerMetrics> serverMetricsMap = new TreeMap<>(); 553 serverMetricsMap.put(server0, 554 mockServerMetricsWithRegionCacheInfo(regionsOnServer0, 1.0f, new ArrayList<>(), 0, 10)); 555 serverMetricsMap.put(server1, 556 mockServerMetricsWithRegionCacheInfo(regionsOnServer1, 1.0f, oldCachedRegions, 10, 10)); 557 serverMetricsMap.put(server2, 558 mockServerMetricsWithRegionCacheInfo(regionsOnServer2, 1.0f, new ArrayList<>(), 0, 10)); 559 ClusterMetrics clusterMetrics = mock(ClusterMetrics.class); 560 when(clusterMetrics.getLiveServerMetrics()).thenReturn(serverMetricsMap); 561 loadBalancer.updateClusterMetrics(clusterMetrics); 562 563 Map<TableName, Map<ServerName, List<RegionInfo>>> LoadOfAllTable = 564 (Map) mockClusterServersWithTables(clusterState); 565 List<RegionPlan> plans = loadBalancer.balanceCluster(LoadOfAllTable); 566 Set<RegionInfo> regionsMovedFromServer0 = new HashSet<>(); 567 Map<ServerName, List<RegionInfo>> targetServers = new HashMap<>(); 568 for (RegionPlan plan : plans) { 569 if (plan.getSource().equals(server0)) { 570 regionsMovedFromServer0.add(plan.getRegionInfo()); 571 if (!targetServers.containsKey(plan.getDestination())) { 572 targetServers.put(plan.getDestination(), new ArrayList<>()); 573 } 574 targetServers.get(plan.getDestination()).add(plan.getRegionInfo()); 575 } 576 } 577 // Skew rebalancing should still move 5 regions from server0 to server1 to balance the 578 // cluster (10 on server0, 0 on server1, 5 on server2 → target ~5 on each). But the 579 // specific regions moved are not dictated by old cache data since all regions are already 580 // well-cached on their current server. 581 assertEquals(5, regionsMovedFromServer0.size()); 582 assertEquals(5, targetServers.get(server1).size()); 583 } 584 585 @Test 586 public void testRegionsPartiallyCachedOnOldServerAndCurrentServer() throws Exception { 587 // The regions are partially cached on old server (0.6) and have lower cache on current (0.2). 588 // The balancer should move regions to server1 to fix skew, and the cache-aware generator 589 // guides some of those moves to be the old-cached regions. 590 591 Map<ServerName, List<RegionInfo>> clusterState = new HashMap<>(); 592 ServerName server0 = servers.get(0); 593 ServerName server1 = servers.get(1); 594 ServerName server2 = servers.get(2); 595 596 // Simulate that the regions previously hosted by server1 are now hosted on server0 597 List<RegionInfo> regionsOnServer0 = randomRegions(10); 598 List<RegionInfo> regionsOnServer1 = randomRegions(0); 599 List<RegionInfo> regionsOnServer2 = randomRegions(5); 600 601 clusterState.put(server0, regionsOnServer0); 602 clusterState.put(server1, regionsOnServer1); 603 clusterState.put(server2, regionsOnServer2); 604 605 // Mock cluster metrics 606 607 // Mock 4 regions from server0 were previously hosted on server1 608 List<RegionInfo> oldCachedRegions = regionsOnServer0.subList(5, regionsOnServer0.size() - 1); 609 610 Map<ServerName, ServerMetrics> serverMetricsMap = new TreeMap<>(); 611 serverMetricsMap.put(server0, 612 mockServerMetricsWithRegionCacheInfo(regionsOnServer0, 0.2f, new ArrayList<>(), 0, 10)); 613 serverMetricsMap.put(server1, 614 mockServerMetricsWithRegionCacheInfo(regionsOnServer1, 0.0f, oldCachedRegions, 6, 10)); 615 serverMetricsMap.put(server2, 616 mockServerMetricsWithRegionCacheInfo(regionsOnServer2, 1.0f, new ArrayList<>(), 0, 10)); 617 ClusterMetrics clusterMetrics = mock(ClusterMetrics.class); 618 when(clusterMetrics.getLiveServerMetrics()).thenReturn(serverMetricsMap); 619 loadBalancer.updateClusterMetrics(clusterMetrics); 620 621 Map<TableName, Map<ServerName, List<RegionInfo>>> LoadOfAllTable = 622 (Map) mockClusterServersWithTables(clusterState); 623 List<RegionPlan> plans = loadBalancer.balanceCluster(LoadOfAllTable); 624 Set<RegionInfo> regionsMovedFromServer0 = new HashSet<>(); 625 Map<ServerName, List<RegionInfo>> targetServers = new HashMap<>(); 626 for (RegionPlan plan : plans) { 627 if (plan.getSource().equals(server0)) { 628 regionsMovedFromServer0.add(plan.getRegionInfo()); 629 if (!targetServers.containsKey(plan.getDestination())) { 630 targetServers.put(plan.getDestination(), new ArrayList<>()); 631 } 632 targetServers.get(plan.getDestination()).add(plan.getRegionInfo()); 633 } 634 } 635 // Balanced state for 15 total regions on 3 servers = 5 each. 636 // server0(10) → server1(0): should move 5 637 assertEquals(5, regionsMovedFromServer0.size()); 638 assertEquals(5, targetServers.get(server1).size()); 639 // The cache-aware generator should move at least some old-cached regions to server1 640 // (where they have better cache). Due to stochastic walk non-determinism, not all 4 641 // are guaranteed to be picked over equally-viable alternatives. 642 long oldCachedOnServer1 = 643 targetServers.get(server1).stream().filter(oldCachedRegions::contains).count(); 644 assertTrue(oldCachedOnServer1 > 0, "At least some old-cached regions should move to server1"); 645 } 646 647 @Test 648 public void testBalancerNotThrowNPEWhenBalancerPlansIsNull() throws Exception { 649 Map<ServerName, List<RegionInfo>> clusterState = new HashMap<>(); 650 ServerName server0 = servers.get(0); 651 ServerName server1 = servers.get(1); 652 ServerName server2 = servers.get(2); 653 654 List<RegionInfo> regionsOnServer0 = randomRegions(5); 655 List<RegionInfo> regionsOnServer1 = randomRegions(5); 656 List<RegionInfo> regionsOnServer2 = randomRegions(5); 657 658 clusterState.put(server0, regionsOnServer0); 659 clusterState.put(server1, regionsOnServer1); 660 clusterState.put(server2, regionsOnServer2); 661 662 // Mock cluster metrics 663 Map<ServerName, ServerMetrics> serverMetricsMap = new TreeMap<>(); 664 serverMetricsMap.put(server0, 665 mockServerMetricsWithRegionCacheInfo(regionsOnServer0, 0.0f, new ArrayList<>(), 0, 10)); 666 serverMetricsMap.put(server1, 667 mockServerMetricsWithRegionCacheInfo(regionsOnServer1, 0.0f, new ArrayList<>(), 0, 10)); 668 serverMetricsMap.put(server2, 669 mockServerMetricsWithRegionCacheInfo(regionsOnServer2, 0.0f, new ArrayList<>(), 0, 10)); 670 671 ClusterMetrics clusterMetrics = mock(ClusterMetrics.class); 672 when(clusterMetrics.getLiveServerMetrics()).thenReturn(serverMetricsMap); 673 loadBalancer.updateClusterMetrics(clusterMetrics); 674 675 Map<TableName, Map<ServerName, List<RegionInfo>>> LoadOfAllTable = 676 (Map) mockClusterServersWithTables(clusterState); 677 try { 678 List<RegionPlan> plans = loadBalancer.balanceCluster(LoadOfAllTable); 679 assertNull(plans); 680 } catch (NullPointerException npe) { 681 fail("NPE should not be thrown"); 682 } 683 } 684 685}