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.assertTrue; 023import static org.mockito.Mockito.mock; 024import static org.mockito.Mockito.when; 025 026import java.util.ArrayList; 027import java.util.HashMap; 028import java.util.HashSet; 029import java.util.List; 030import java.util.Map; 031import java.util.Set; 032import java.util.TreeMap; 033import java.util.concurrent.CountDownLatch; 034import java.util.concurrent.ExecutorService; 035import java.util.concurrent.Executors; 036import java.util.concurrent.Future; 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.ServerMetrics; 042import org.apache.hadoop.hbase.ServerName; 043import org.apache.hadoop.hbase.TableName; 044import org.apache.hadoop.hbase.client.RegionInfo; 045import org.apache.hadoop.hbase.master.RegionPlan; 046import org.apache.hadoop.hbase.rsgroup.RSGroupBasedLoadBalancer; 047import org.apache.hadoop.hbase.rsgroup.RSGroupInfo; 048import org.apache.hadoop.hbase.testclassification.LargeTests; 049import org.apache.hadoop.hbase.util.EnvironmentEdgeManager; 050import org.apache.hadoop.hbase.util.Pair; 051import org.junit.jupiter.api.BeforeAll; 052import org.junit.jupiter.api.Tag; 053import org.junit.jupiter.api.Test; 054import org.junit.jupiter.api.Timeout; 055import org.slf4j.Logger; 056import org.slf4j.LoggerFactory; 057 058@Tag(LargeTests.TAG) 059public class TestRSGroupBasedLoadBalancerWithCacheAwareLoadBalancerAsInternal 060 extends RSGroupableBalancerTestBase { 061 062 private static final Logger LOG = 063 LoggerFactory.getLogger(TestRSGroupBasedLoadBalancerWithCacheAwareLoadBalancerAsInternal.class); 064 065 private static RSGroupBasedLoadBalancer loadBalancer; 066 067 @BeforeAll 068 public static void beforeAllTests() throws Exception { 069 groups = new String[] { RSGroupInfo.DEFAULT_GROUP }; 070 servers = generateServers(3); 071 groupMap = constructGroupInfo(servers, groups); 072 tableDescs = constructTableDesc(false); 073 Configuration cong = HBaseConfiguration.create(); 074 conf.set(HConstants.BUCKET_CACHE_PERSISTENT_PATH_KEY, "prefetch_file_list"); 075 conf.set("hbase.rsgroup.grouploadbalancer.class", 076 CacheAwareLoadBalancer.class.getCanonicalName()); 077 loadBalancer = new RSGroupBasedLoadBalancer(); 078 loadBalancer.setMasterServices(getMockedMaster()); 079 loadBalancer.initialize(); 080 } 081 082 @Test 083 public void testRegionsNotCachedOnOldServerAndCurrentServer() throws Exception { 084 // The regions are not cached on old server as well as the current server. This causes 085 // skewness in the region allocation which should be fixed by the balancer 086 087 Map<ServerName, List<RegionInfo>> clusterState = new HashMap<>(); 088 ServerName server0 = servers.get(0); 089 ServerName server1 = servers.get(1); 090 ServerName server2 = servers.get(2); 091 092 // Simulate that the regions previously hosted by server1 are now hosted on server0 093 List<RegionInfo> regionsOnServer0 = randomRegions(10); 094 List<RegionInfo> regionsOnServer1 = randomRegions(0); 095 List<RegionInfo> regionsOnServer2 = randomRegions(5); 096 097 clusterState.put(server0, regionsOnServer0); 098 clusterState.put(server1, regionsOnServer1); 099 clusterState.put(server2, regionsOnServer2); 100 101 // Mock cluster metrics 102 Map<ServerName, ServerMetrics> serverMetricsMap = new TreeMap<>(); 103 serverMetricsMap.put(server0, mockServerMetricsWithRegionCacheInfo(server0, regionsOnServer0, 104 0.0f, new ArrayList<>(), 0, 10)); 105 serverMetricsMap.put(server1, mockServerMetricsWithRegionCacheInfo(server1, regionsOnServer1, 106 0.0f, new ArrayList<>(), 0, 10)); 107 serverMetricsMap.put(server2, mockServerMetricsWithRegionCacheInfo(server2, regionsOnServer2, 108 0.0f, new ArrayList<>(), 0, 10)); 109 ClusterMetrics clusterMetrics = mock(ClusterMetrics.class); 110 when(clusterMetrics.getLiveServerMetrics()).thenReturn(serverMetricsMap); 111 loadBalancer.updateClusterMetrics(clusterMetrics); 112 113 Map<TableName, Map<ServerName, List<RegionInfo>>> LoadOfAllTable = 114 (Map) mockClusterServersWithTables(clusterState); 115 List<RegionPlan> plans = loadBalancer.balanceCluster(LoadOfAllTable); 116 Set<RegionInfo> regionsMovedFromServer0 = new HashSet<>(); 117 Map<ServerName, List<RegionInfo>> targetServers = new HashMap<>(); 118 for (RegionPlan plan : plans) { 119 if (plan.getSource().equals(server0)) { 120 regionsMovedFromServer0.add(plan.getRegionInfo()); 121 if (!targetServers.containsKey(plan.getDestination())) { 122 targetServers.put(plan.getDestination(), new ArrayList<>()); 123 } 124 targetServers.get(plan.getDestination()).add(plan.getRegionInfo()); 125 } 126 } 127 // should move 5 regions from server0 to server 1 128 assertEquals(5, regionsMovedFromServer0.size()); 129 assertEquals(5, targetServers.get(server1).size()); 130 } 131 132 /** 133 * Regions on the overloaded RS report low block-cache ratio; no RS reports prefetch/historical 134 * cache for those regions (so {@link CacheAwareLoadBalancer.CacheAwareCandidateGenerator} has no 135 * "old server" to prefer). Another RS has ample free block cache. The balancer should still emit 136 * plans that shed load from the hot RS onto the idle RS with spare cache capacity. 137 */ 138 @Test 139 public void testLowCacheRatioNoHistoricalCacheRelocatesWhenTargetHasFreeBlockCache() 140 throws Exception { 141 Map<ServerName, List<RegionInfo>> clusterState = new HashMap<>(); 142 ServerName server0 = servers.get(0); 143 ServerName server1 = servers.get(1); 144 ServerName server2 = servers.get(2); 145 146 List<RegionInfo> regionsOnServer0 = randomRegions(10); 147 List<RegionInfo> regionsOnServer1 = randomRegions(0); 148 List<RegionInfo> regionsOnServer2 = randomRegions(5); 149 150 clusterState.put(server0, regionsOnServer0); 151 clusterState.put(server1, regionsOnServer1); 152 clusterState.put(server2, regionsOnServer2); 153 154 // Below LOW_CACHE_RATIO_FOR_RELOCATION_DEFAULT (0.35); 155 ServerMetrics sm0 = mockServerMetricsWithRegionCacheInfo(server0, regionsOnServer0, 0.1f, 156 new ArrayList<>(), 0, 10); 157 when(sm0.getCacheFreeSize()).thenReturn(0L); 158 ServerMetrics sm1 = mockServerMetricsWithRegionCacheInfo(server1, regionsOnServer1, 0.0f, 159 new ArrayList<>(), 0, 10); 160 // Simulates 1GB free cache space on server1 161 when(sm1.getCacheFreeSize()).thenReturn(1024L * 1024 * 1024); 162 ServerMetrics sm2 = mockServerMetricsWithRegionCacheInfo(server2, regionsOnServer2, 1.0f, 163 new ArrayList<>(), 0, 10); 164 when(sm2.getCacheFreeSize()).thenReturn(0L); 165 166 Map<ServerName, ServerMetrics> serverMetricsMap = new TreeMap<>(); 167 serverMetricsMap.put(server0, sm0); 168 serverMetricsMap.put(server1, sm1); 169 serverMetricsMap.put(server2, sm2); 170 ClusterMetrics clusterMetrics = mock(ClusterMetrics.class); 171 when(clusterMetrics.getLiveServerMetrics()).thenReturn(serverMetricsMap); 172 loadBalancer.updateClusterMetrics(clusterMetrics); 173 174 CacheAwareLoadBalancer internalBalancer = 175 (CacheAwareLoadBalancer) loadBalancer.getInternalBalancer(); 176 assertNotNull(internalBalancer); 177 assertTrue(internalBalancer.regionCacheRatioOnOldServerMap.isEmpty()); 178 179 Map<TableName, Map<ServerName, List<RegionInfo>>> loadOfAllTable = 180 (Map) mockClusterServersWithTables(clusterState); 181 List<RegionPlan> plans = loadBalancer.balanceCluster(loadOfAllTable); 182 assertNotNull(plans); 183 184 Set<RegionInfo> regionsMovedFromServer0 = new HashSet<>(); 185 Map<ServerName, List<RegionInfo>> targetServers = new HashMap<>(); 186 for (RegionPlan plan : plans) { 187 if (plan.getSource().equals(server0)) { 188 regionsMovedFromServer0.add(plan.getRegionInfo()); 189 if (!targetServers.containsKey(plan.getDestination())) { 190 targetServers.put(plan.getDestination(), new ArrayList<>()); 191 } 192 targetServers.get(plan.getDestination()).add(plan.getRegionInfo()); 193 } 194 } 195 assertEquals(5, regionsMovedFromServer0.size()); 196 assertNotNull(targetServers.get(server1)); 197 assertEquals(5, targetServers.get(server1).size()); 198 } 199 200 @Test 201 public void testRegionsPartiallyCachedOnOldServerAndNotCachedOnCurrentServer() throws Exception { 202 Map<ServerName, List<RegionInfo>> clusterState = new HashMap<>(); 203 ServerName server0 = servers.get(0); 204 ServerName server1 = servers.get(1); 205 ServerName server2 = servers.get(2); 206 207 // Simulate that the regions previously hosted by server1 are now hosted on server0 208 List<RegionInfo> regionsOnServer0 = randomRegions(10); 209 List<RegionInfo> regionsOnServer1 = randomRegions(0); 210 List<RegionInfo> regionsOnServer2 = randomRegions(5); 211 212 clusterState.put(server0, regionsOnServer0); 213 clusterState.put(server1, regionsOnServer1); 214 clusterState.put(server2, regionsOnServer2); 215 216 // Mock cluster metrics 217 218 // Mock 5 regions from server0 were previously hosted on server1 219 List<RegionInfo> oldCachedRegions = regionsOnServer0.subList(5, regionsOnServer0.size()); 220 221 Map<ServerName, ServerMetrics> serverMetricsMap = new TreeMap<>(); 222 serverMetricsMap.put(server0, mockServerMetricsWithRegionCacheInfo(server0, regionsOnServer0, 223 0.0f, new ArrayList<>(), 0, 10)); 224 serverMetricsMap.put(server1, mockServerMetricsWithRegionCacheInfo(server1, regionsOnServer1, 225 0.0f, oldCachedRegions, 6, 10)); 226 serverMetricsMap.put(server2, mockServerMetricsWithRegionCacheInfo(server2, regionsOnServer2, 227 0.0f, new ArrayList<>(), 0, 10)); 228 ClusterMetrics clusterMetrics = mock(ClusterMetrics.class); 229 when(clusterMetrics.getLiveServerMetrics()).thenReturn(serverMetricsMap); 230 loadBalancer.updateClusterMetrics(clusterMetrics); 231 232 Map<TableName, Map<ServerName, List<RegionInfo>>> LoadOfAllTable = 233 (Map) mockClusterServersWithTables(clusterState); 234 List<RegionPlan> plans = loadBalancer.balanceCluster(LoadOfAllTable); 235 Set<RegionInfo> regionsMovedFromServer0 = new HashSet<>(); 236 Map<ServerName, List<RegionInfo>> targetServers = new HashMap<>(); 237 for (RegionPlan plan : plans) { 238 if (plan.getSource().equals(server0)) { 239 regionsMovedFromServer0.add(plan.getRegionInfo()); 240 if (!targetServers.containsKey(plan.getDestination())) { 241 targetServers.put(plan.getDestination(), new ArrayList<>()); 242 } 243 targetServers.get(plan.getDestination()).add(plan.getRegionInfo()); 244 } 245 } 246 // should move regions from server0 to server1 (old-cached regions should be among them) 247 assertTrue(regionsMovedFromServer0.size() >= 4); 248 assertNotNull(targetServers.get(server1)); 249 int oldCachedOnServer1 = 0; 250 for (RegionInfo ri : oldCachedRegions) { 251 if (targetServers.get(server1).contains(ri)) { 252 oldCachedOnServer1++; 253 } 254 } 255 assertTrue(oldCachedOnServer1 > 0, 256 "Expected old-cached regions to move to server1, got " + oldCachedOnServer1); 257 } 258 259 @Test 260 public void testThrottlingRegionBeyondThreshold() throws Exception { 261 Configuration conf = HBaseConfiguration.create(); 262 conf.set(HConstants.BUCKET_CACHE_PERSISTENT_PATH_KEY, "prefetch_file_list"); 263 conf.set("hbase.rsgroup.grouploadbalancer.class", 264 CacheAwareLoadBalancer.class.getCanonicalName()); 265 RSGroupBasedLoadBalancer balancer = new RSGroupBasedLoadBalancer(); 266 balancer.setMasterServices(getMockedMaster()); 267 balancer.initialize(); 268 269 ServerName server0 = servers.get(0); 270 ServerName server1 = servers.get(1); 271 Pair<ServerName, Float> regionRatio = new Pair<>(); 272 regionRatio.setFirst(server0); 273 regionRatio.setSecond(1.0f); 274 CacheAwareLoadBalancer internalBalancer = 275 (CacheAwareLoadBalancer) balancer.getInternalBalancer(); 276 internalBalancer.regionCacheRatioOnOldServerMap.put("region1", regionRatio); 277 RegionInfo mockedInfo = mock(RegionInfo.class); 278 when(mockedInfo.getEncodedName()).thenReturn("region1"); 279 RegionPlan plan = new RegionPlan(mockedInfo, server1, server0); 280 long startTime = EnvironmentEdgeManager.currentTime(); 281 synchronized (balancer) { 282 balancer.throttle(plan); 283 } 284 long endTime = EnvironmentEdgeManager.currentTime(); 285 assertTrue((endTime - startTime) < 10); 286 } 287 288 @Test 289 public void testThrottlingRegionBelowThreshold() throws Exception { 290 Configuration conf = HBaseConfiguration.create(); 291 conf.set(HConstants.BUCKET_CACHE_PERSISTENT_PATH_KEY, "prefetch_file_list"); 292 conf.setLong(CacheAwareLoadBalancer.MOVE_THROTTLING, 100); 293 conf.set("hbase.rsgroup.grouploadbalancer.class", 294 CacheAwareLoadBalancer.class.getCanonicalName()); 295 RSGroupBasedLoadBalancer loadBalancer = new RSGroupBasedLoadBalancer(); 296 loadBalancer.setMasterServices(getMockedMaster()); 297 loadBalancer.initialize(); 298 CacheAwareLoadBalancer internalBalancer = 299 (CacheAwareLoadBalancer) loadBalancer.getInternalBalancer(); 300 internalBalancer.loadConf(conf); 301 302 ServerName server0 = servers.get(0); 303 ServerName server1 = servers.get(1); 304 Pair<ServerName, Float> regionRatio = new Pair<>(); 305 regionRatio.setFirst(server0); 306 regionRatio.setSecond(0.1f); 307 internalBalancer = (CacheAwareLoadBalancer) loadBalancer.getInternalBalancer(); 308 internalBalancer.regionCacheRatioOnOldServerMap.put("region1", regionRatio); 309 RegionInfo mockedInfo = mock(RegionInfo.class); 310 when(mockedInfo.getEncodedName()).thenReturn("region1"); 311 RegionPlan plan = new RegionPlan(mockedInfo, server1, server0); 312 long startTime = EnvironmentEdgeManager.currentTime(); 313 synchronized (loadBalancer) { 314 loadBalancer.throttle(plan); 315 } 316 long endTime = EnvironmentEdgeManager.currentTime(); 317 assertTrue((endTime - startTime) >= 100); 318 } 319 320 @Test 321 public void testThrottlingCacheRatioUnknownOnTarget() throws Exception { 322 Configuration conf = HBaseConfiguration.create(); 323 conf.set(HConstants.BUCKET_CACHE_PERSISTENT_PATH_KEY, "prefetch_file_list"); 324 conf.setLong(CacheAwareLoadBalancer.MOVE_THROTTLING, 100); 325 conf.set("hbase.rsgroup.grouploadbalancer.class", 326 CacheAwareLoadBalancer.class.getCanonicalName()); 327 RSGroupBasedLoadBalancer loadBalancer = new RSGroupBasedLoadBalancer(); 328 loadBalancer.setMasterServices(getMockedMaster()); 329 loadBalancer.initialize(); 330 CacheAwareLoadBalancer internalBalancer = 331 (CacheAwareLoadBalancer) loadBalancer.getInternalBalancer(); 332 internalBalancer.loadConf(conf); 333 334 ServerName server0 = servers.get(0); 335 ServerName server1 = servers.get(1); 336 ServerName server3 = servers.get(2); 337 // setting region cache ratio 100% on server 3, though this is not the target in the region plan 338 Pair<ServerName, Float> regionRatio = new Pair<>(); 339 regionRatio.setFirst(server3); 340 regionRatio.setSecond(1.0f); 341 internalBalancer = (CacheAwareLoadBalancer) loadBalancer.getInternalBalancer(); 342 internalBalancer.regionCacheRatioOnOldServerMap.put("region1", regionRatio); 343 RegionInfo mockedInfo = mock(RegionInfo.class); 344 when(mockedInfo.getEncodedName()).thenReturn("region1"); 345 RegionPlan plan = new RegionPlan(mockedInfo, server1, server0); 346 long startTime = EnvironmentEdgeManager.currentTime(); 347 synchronized (loadBalancer) { 348 loadBalancer.throttle(plan); 349 } 350 long endTime = EnvironmentEdgeManager.currentTime(); 351 assertTrue((endTime - startTime) >= 100); 352 } 353 354 @Test 355 public void testThrottlingCacheRatioUnknownForRegion() throws Exception { 356 Configuration conf = HBaseConfiguration.create(); 357 conf.set(HConstants.BUCKET_CACHE_PERSISTENT_PATH_KEY, "prefetch_file_list"); 358 conf.setLong(CacheAwareLoadBalancer.MOVE_THROTTLING, 100); 359 conf.set("hbase.rsgroup.grouploadbalancer.class", 360 CacheAwareLoadBalancer.class.getCanonicalName()); 361 RSGroupBasedLoadBalancer loadBalancer = new RSGroupBasedLoadBalancer(); 362 loadBalancer.setMasterServices(getMockedMaster()); 363 loadBalancer.initialize(); 364 CacheAwareLoadBalancer internalBalancer = 365 (CacheAwareLoadBalancer) loadBalancer.getInternalBalancer(); 366 internalBalancer.loadConf(conf); 367 368 ServerName server0 = servers.get(0); 369 ServerName server1 = servers.get(1); 370 ServerName server3 = servers.get(2); 371 // No cache ratio available for region1 372 RegionInfo mockedInfo = mock(RegionInfo.class); 373 when(mockedInfo.getEncodedName()).thenReturn("region1"); 374 RegionPlan plan = new RegionPlan(mockedInfo, server1, server0); 375 long startTime = EnvironmentEdgeManager.currentTime(); 376 synchronized (loadBalancer) { 377 loadBalancer.throttle(plan); 378 } 379 long endTime = EnvironmentEdgeManager.currentTime(); 380 assertTrue((endTime - startTime) >= 100); 381 } 382 383 @Test 384 public void testRegionPlansSortedByCacheRatioOnTarget() throws Exception { 385 // The regions are fully cached on old server 386 387 Map<ServerName, List<RegionInfo>> clusterState = new HashMap<>(); 388 ServerName server0 = servers.get(0); 389 ServerName server1 = servers.get(1); 390 ServerName server2 = servers.get(2); 391 392 // Simulate on RS with all regions, and two RSes with no regions 393 List<RegionInfo> regionsOnServer0 = randomRegions(15); 394 List<RegionInfo> regionsOnServer1 = randomRegions(0); 395 List<RegionInfo> regionsOnServer2 = randomRegions(0); 396 397 clusterState.put(server0, regionsOnServer0); 398 clusterState.put(server1, regionsOnServer1); 399 clusterState.put(server2, regionsOnServer2); 400 401 // Mock cluster metrics 402 // Mock 5 regions from server0 were previously hosted on server1 403 List<RegionInfo> oldCachedRegions1 = regionsOnServer0.subList(5, 10); 404 List<RegionInfo> oldCachedRegions2 = regionsOnServer0.subList(10, regionsOnServer0.size()); 405 Map<ServerName, ServerMetrics> serverMetricsMap = new TreeMap<>(); 406 // mock server metrics to set cache ratio as 0 in the RS 0 407 serverMetricsMap.put(server0, mockServerMetricsWithRegionCacheInfo(server0, regionsOnServer0, 408 0.0f, new ArrayList<>(), 0, 10)); 409 // mock server metrics to set cache ratio as 1 in the RS 1 410 serverMetricsMap.put(server1, mockServerMetricsWithRegionCacheInfo(server1, regionsOnServer1, 411 0.0f, oldCachedRegions1, 10, 10)); 412 // mock server metrics to set cache ratio as .8 in the RS 2 413 serverMetricsMap.put(server2, mockServerMetricsWithRegionCacheInfo(server2, regionsOnServer2, 414 0.0f, oldCachedRegions2, 8, 10)); 415 ClusterMetrics clusterMetrics = mock(ClusterMetrics.class); 416 when(clusterMetrics.getLiveServerMetrics()).thenReturn(serverMetricsMap); 417 loadBalancer.updateClusterMetrics(clusterMetrics); 418 419 Map<TableName, Map<ServerName, List<RegionInfo>>> LoadOfAllTable = 420 (Map) mockClusterServersWithTables(clusterState); 421 List<RegionPlan> plans = loadBalancer.balanceCluster(LoadOfAllTable); 422 LOG.debug("plans size: {}", plans.size()); 423 LOG.debug("plans: {}", plans); 424 // Plans are sorted by cache ratio on destination (descending). Verify ordering 425 // and that at least some old-cached regions are moved to their cached servers. 426 float prevRatio = Float.MAX_VALUE; 427 int oldCached1Count = 0; 428 int oldCached2Count = 0; 429 for (RegionPlan plan : plans) { 430 LOG.debug("plan region: {}, target server: {}", plan.getRegionInfo().getEncodedName(), 431 plan.getDestination().getServerName()); 432 float ratio = 0f; 433 if ( 434 oldCachedRegions1.contains(plan.getRegionInfo()) && server1.equals(plan.getDestination()) 435 ) { 436 ratio = 1.0f; 437 oldCached1Count++; 438 } else if ( 439 oldCachedRegions2.contains(plan.getRegionInfo()) && server2.equals(plan.getDestination()) 440 ) { 441 ratio = 0.8f; 442 oldCached2Count++; 443 } 444 assertTrue(ratio <= prevRatio, 445 "Plans should be sorted by cache ratio on destination (descending)"); 446 prevRatio = ratio; 447 } 448 assertTrue(oldCached1Count > 0, "Some old-cached regions should move to server1"); 449 assertTrue(oldCached2Count > 0, "Some old-cached regions should move to server2"); 450 } 451 452 @Test 453 public void testRegionsFullyCachedOnOldServerAndNotCachedOnCurrentServers() throws Exception { 454 // The regions are fully cached on old server 455 456 Map<ServerName, List<RegionInfo>> clusterState = new HashMap<>(); 457 ServerName server0 = servers.get(0); 458 ServerName server1 = servers.get(1); 459 ServerName server2 = servers.get(2); 460 461 // Simulate that the regions previously hosted by server1 are now hosted on server0 462 List<RegionInfo> regionsOnServer0 = randomRegions(10); 463 List<RegionInfo> regionsOnServer1 = randomRegions(0); 464 List<RegionInfo> regionsOnServer2 = randomRegions(5); 465 466 clusterState.put(server0, regionsOnServer0); 467 clusterState.put(server1, regionsOnServer1); 468 clusterState.put(server2, regionsOnServer2); 469 470 // Mock cluster metrics 471 472 // Mock 5 regions from server0 were previously hosted on server1 473 List<RegionInfo> oldCachedRegions = regionsOnServer0.subList(5, regionsOnServer0.size()); 474 475 Map<ServerName, ServerMetrics> serverMetricsMap = new TreeMap<>(); 476 serverMetricsMap.put(server0, mockServerMetricsWithRegionCacheInfo(server0, regionsOnServer0, 477 0.0f, new ArrayList<>(), 0, 10)); 478 serverMetricsMap.put(server1, mockServerMetricsWithRegionCacheInfo(server1, regionsOnServer1, 479 0.0f, oldCachedRegions, 10, 10)); 480 serverMetricsMap.put(server2, mockServerMetricsWithRegionCacheInfo(server2, regionsOnServer2, 481 0.0f, new ArrayList<>(), 0, 10)); 482 ClusterMetrics clusterMetrics = mock(ClusterMetrics.class); 483 when(clusterMetrics.getLiveServerMetrics()).thenReturn(serverMetricsMap); 484 loadBalancer.updateClusterMetrics(clusterMetrics); 485 486 Map<TableName, Map<ServerName, List<RegionInfo>>> LoadOfAllTable = 487 (Map) mockClusterServersWithTables(clusterState); 488 List<RegionPlan> plans = loadBalancer.balanceCluster(LoadOfAllTable); 489 Set<RegionInfo> regionsMovedFromServer0 = new HashSet<>(); 490 Map<ServerName, List<RegionInfo>> targetServers = new HashMap<>(); 491 for (RegionPlan plan : plans) { 492 if (plan.getSource().equals(server0)) { 493 regionsMovedFromServer0.add(plan.getRegionInfo()); 494 if (!targetServers.containsKey(plan.getDestination())) { 495 targetServers.put(plan.getDestination(), new ArrayList<>()); 496 } 497 targetServers.get(plan.getDestination()).add(plan.getRegionInfo()); 498 } 499 } 500 // should move regions from server0 to server1 (old-cached regions should be among them) 501 assertTrue(regionsMovedFromServer0.size() >= 4); 502 assertNotNull(targetServers.get(server1)); 503 assertTrue(targetServers.get(server1).size() >= 4); 504 int oldCachedOnServer1 = 0; 505 for (RegionInfo ri : oldCachedRegions) { 506 if (targetServers.get(server1).contains(ri)) { 507 oldCachedOnServer1++; 508 } 509 } 510 assertTrue(oldCachedOnServer1 > 0, 511 "Expected most old-cached regions to move to server1, got " + oldCachedOnServer1); 512 } 513 514 @Test 515 public void testRegionsFullyCachedOnOldAndCurrentServers() throws Exception { 516 // The regions are fully cached on old server 517 518 Map<ServerName, List<RegionInfo>> clusterState = new HashMap<>(); 519 ServerName server0 = servers.get(0); 520 ServerName server1 = servers.get(1); 521 ServerName server2 = servers.get(2); 522 523 // Simulate that the regions previously hosted by server1 are now hosted on server0 524 List<RegionInfo> regionsOnServer0 = randomRegions(10); 525 List<RegionInfo> regionsOnServer1 = randomRegions(0); 526 List<RegionInfo> regionsOnServer2 = randomRegions(5); 527 528 clusterState.put(server0, regionsOnServer0); 529 clusterState.put(server1, regionsOnServer1); 530 clusterState.put(server2, regionsOnServer2); 531 532 // Mock cluster metrics 533 534 // Mock 5 regions from server0 were previously hosted on server1 535 List<RegionInfo> oldCachedRegions = regionsOnServer0.subList(5, regionsOnServer0.size()); 536 537 Map<ServerName, ServerMetrics> serverMetricsMap = new TreeMap<>(); 538 serverMetricsMap.put(server0, mockServerMetricsWithRegionCacheInfo(server0, regionsOnServer0, 539 1.0f, new ArrayList<>(), 0, 10)); 540 serverMetricsMap.put(server1, mockServerMetricsWithRegionCacheInfo(server1, regionsOnServer1, 541 1.0f, oldCachedRegions, 10, 10)); 542 serverMetricsMap.put(server2, mockServerMetricsWithRegionCacheInfo(server2, regionsOnServer2, 543 1.0f, new ArrayList<>(), 0, 10)); 544 ClusterMetrics clusterMetrics = mock(ClusterMetrics.class); 545 when(clusterMetrics.getLiveServerMetrics()).thenReturn(serverMetricsMap); 546 loadBalancer.updateClusterMetrics(clusterMetrics); 547 548 Map<TableName, Map<ServerName, List<RegionInfo>>> LoadOfAllTable = 549 (Map) mockClusterServersWithTables(clusterState); 550 List<RegionPlan> plans = loadBalancer.balanceCluster(LoadOfAllTable); 551 Set<RegionInfo> regionsMovedFromServer0 = new HashSet<>(); 552 Map<ServerName, List<RegionInfo>> targetServers = new HashMap<>(); 553 for (RegionPlan plan : plans) { 554 if (plan.getSource().equals(server0)) { 555 regionsMovedFromServer0.add(plan.getRegionInfo()); 556 if (!targetServers.containsKey(plan.getDestination())) { 557 targetServers.put(plan.getDestination(), new ArrayList<>()); 558 } 559 targetServers.get(plan.getDestination()).add(plan.getRegionInfo()); 560 } 561 } 562 // should move 5 regions from server0 to server1 to balance the cluster, but the specific 563 // regions moved are not dictated by old cache data since all regions are already well-cached 564 // on their current server (currentCacheRatio >= ratioThreshold) 565 assertEquals(5, regionsMovedFromServer0.size()); 566 assertEquals(5, targetServers.get(server1).size()); 567 } 568 569 @Test 570 public void testRegionsPartiallyCachedOnOldServerAndCurrentServer() throws Exception { 571 // The regions are partially cached on old server 572 573 Map<ServerName, List<RegionInfo>> clusterState = new HashMap<>(); 574 ServerName server0 = servers.get(0); 575 ServerName server1 = servers.get(1); 576 ServerName server2 = servers.get(2); 577 578 // Simulate that the regions previously hosted by server1 are now hosted on server0 579 List<RegionInfo> regionsOnServer0 = randomRegions(10); 580 List<RegionInfo> regionsOnServer1 = randomRegions(0); 581 List<RegionInfo> regionsOnServer2 = randomRegions(5); 582 583 clusterState.put(server0, regionsOnServer0); 584 clusterState.put(server1, regionsOnServer1); 585 clusterState.put(server2, regionsOnServer2); 586 587 // Mock cluster metrics 588 589 // Mock 5 regions from server0 were previously hosted on server1 590 List<RegionInfo> oldCachedRegions = regionsOnServer0.subList(5, regionsOnServer0.size()); 591 592 Map<ServerName, ServerMetrics> serverMetricsMap = new TreeMap<>(); 593 serverMetricsMap.put(server0, mockServerMetricsWithRegionCacheInfo(server0, regionsOnServer0, 594 0.2f, new ArrayList<>(), 0, 10)); 595 serverMetricsMap.put(server1, mockServerMetricsWithRegionCacheInfo(server1, regionsOnServer1, 596 0.0f, oldCachedRegions, 6, 10)); 597 serverMetricsMap.put(server2, mockServerMetricsWithRegionCacheInfo(server2, regionsOnServer2, 598 1.0f, new ArrayList<>(), 0, 10)); 599 ClusterMetrics clusterMetrics = mock(ClusterMetrics.class); 600 when(clusterMetrics.getLiveServerMetrics()).thenReturn(serverMetricsMap); 601 loadBalancer.updateClusterMetrics(clusterMetrics); 602 603 Map<TableName, Map<ServerName, List<RegionInfo>>> LoadOfAllTable = 604 (Map) mockClusterServersWithTables(clusterState); 605 List<RegionPlan> plans = loadBalancer.balanceCluster(LoadOfAllTable); 606 Set<RegionInfo> regionsMovedFromServer0 = new HashSet<>(); 607 Map<ServerName, List<RegionInfo>> targetServers = new HashMap<>(); 608 for (RegionPlan plan : plans) { 609 if (plan.getSource().equals(server0)) { 610 regionsMovedFromServer0.add(plan.getRegionInfo()); 611 if (!targetServers.containsKey(plan.getDestination())) { 612 targetServers.put(plan.getDestination(), new ArrayList<>()); 613 } 614 targetServers.get(plan.getDestination()).add(plan.getRegionInfo()); 615 } 616 } 617 // server0(10) → server1(0): should move 5 618 assertEquals(5, regionsMovedFromServer0.size()); 619 assertEquals(5, targetServers.get(server1).size()); 620 // The cache-aware generator should move at least some old-cached regions to server1. 621 // Due to stochastic walk non-determinism, not all are guaranteed. 622 long oldCachedOnServer1 = 623 targetServers.get(server1).stream().filter(oldCachedRegions::contains).count(); 624 assertTrue(oldCachedOnServer1 > 0, "At least some old-cached regions should move to server1"); 625 } 626 627 @Timeout(60) 628 @Test 629 public void testConfigUpdateDuringBalance() throws Exception { 630 float expectedOldRatioThreshold = 0.8f; 631 float expectedNewRatioThreshold = 0.95f; 632 long throttlingTimeMillis = 10000; 633 634 conf = HBaseConfiguration.create(); 635 conf.setLong(CacheAwareLoadBalancer.MOVE_THROTTLING, throttlingTimeMillis); 636 conf.setFloat(CacheAwareLoadBalancer.CACHE_RATIO_THRESHOLD, expectedOldRatioThreshold); 637 conf.set(HConstants.BUCKET_CACHE_PERSISTENT_PATH_KEY, "prefetch_file_list"); 638 conf.set("hbase.rsgroup.grouploadbalancer.class", 639 CacheAwareLoadBalancer.class.getCanonicalName()); 640 641 RSGroupBasedLoadBalancer balancer = new RSGroupBasedLoadBalancer(); 642 balancer.setMasterServices(getMockedMaster()); 643 balancer.initialize(); 644 645 Map<ServerName, List<RegionInfo>> clusterState = new HashMap<>(); 646 ServerName server0 = servers.get(0); 647 ServerName server1 = servers.get(1); 648 ServerName server2 = servers.get(2); 649 650 // Setup cluster: all 3 regions on server0 (unbalanced) 651 List<RegionInfo> regionsOnServer0 = randomRegions(3); 652 List<RegionInfo> regionsOnServer1 = randomRegions(0); 653 List<RegionInfo> regionsOnServer2 = randomRegions(0); 654 655 clusterState.put(server0, regionsOnServer0); 656 clusterState.put(server1, regionsOnServer1); 657 clusterState.put(server2, regionsOnServer2); 658 659 // Mock metrics: regions have moderate cache ratio so throttle applies on move 660 Map<ServerName, ServerMetrics> serverMetricsMap = new TreeMap<>(); 661 serverMetricsMap.put(server0, mockServerMetricsWithRegionCacheInfo(server0, regionsOnServer0, 662 0.5f, new ArrayList<>(), 0, 10)); 663 serverMetricsMap.put(server1, mockServerMetricsWithRegionCacheInfo(server1, regionsOnServer1, 664 0.5f, new ArrayList<>(), 0, 10)); 665 serverMetricsMap.put(server2, mockServerMetricsWithRegionCacheInfo(server2, regionsOnServer2, 666 0.5f, new ArrayList<>(), 0, 10)); 667 668 ClusterMetrics clusterMetrics = mock(ClusterMetrics.class); 669 when(clusterMetrics.getLiveServerMetrics()).thenReturn(serverMetricsMap); 670 balancer.updateClusterMetrics(clusterMetrics); 671 672 final Map<TableName, Map<ServerName, List<RegionInfo>>> loadOfAllTable = 673 (Map) mockClusterServersWithTables(clusterState); 674 675 // Verify initial configuration is set correctly 676 CacheAwareLoadBalancer internalBalancer = 677 (CacheAwareLoadBalancer) balancer.getInternalBalancer(); 678 assertEquals(expectedOldRatioThreshold, internalBalancer.ratioThreshold, 0.001f); 679 680 CountDownLatch balanceStarted = new CountDownLatch(1); 681 CountDownLatch configUpdateInitiated = new CountDownLatch(1); 682 long[] configUpdateDuration = new long[1]; 683 long[] balanceDuration = new long[1]; 684 685 // Actual old ratio threshold used during balance 686 float[] actualOldRatioThresholdDuringBalance = new float[1]; 687 688 ExecutorService executor = Executors.newFixedThreadPool(2); 689 690 try { 691 // Thread 1 Simulate similar flow to HMaster.balance() which holds synchronized(balancer) for 692 // the duration of balance 693 Future<Long> balanceFuture = executor.submit(() -> { 694 try { 695 long start = EnvironmentEdgeManager.currentTime(); 696 synchronized (balancer) { 697 try { 698 // Simulate beginning of HMaster.balance() mark balancing window open 699 balancer.onBalancingStart(); 700 balanceStarted.countDown(); 701 702 List<RegionPlan> plans = balancer.balanceCluster(loadOfAllTable); 703 if (plans != null) { 704 for (RegionPlan plan : plans) { 705 balancer.throttle(plan); 706 } 707 } 708 // Wait until config update is initiated while balance is still in progress 709 configUpdateInitiated.await(); 710 711 // Old config should still be visible during current balance run 712 CacheAwareLoadBalancer currentInternal = 713 (CacheAwareLoadBalancer) balancer.getInternalBalancer(); 714 actualOldRatioThresholdDuringBalance[0] = currentInternal.ratioThreshold; 715 716 } finally { 717 balancer.onBalancingComplete(); 718 } 719 } 720 return EnvironmentEdgeManager.currentTime() - start; 721 } catch (Exception e) { 722 throw new RuntimeException(e); 723 } 724 }); 725 726 // Thread 2: Simulate update_all_config / onConfigurationChange 727 Future<Long> configUpdateFuture = executor.submit(() -> { 728 try { 729 long startTime = System.currentTimeMillis(); 730 // Wait for balance to start 731 balanceStarted.await(); 732 733 // Call onConfigurationChange - should NOT hang 734 Configuration newConf = new Configuration(conf); 735 newConf.set(HConstants.BUCKET_CACHE_PERSISTENT_PATH_KEY, "prefetch_file_list"); 736 newConf.setLong(CacheAwareLoadBalancer.MOVE_THROTTLING, 10000); 737 newConf.setFloat(CacheAwareLoadBalancer.CACHE_RATIO_THRESHOLD, expectedNewRatioThreshold); 738 balancer.onConfigurationChange(newConf); 739 configUpdateInitiated.countDown(); 740 741 return System.currentTimeMillis() - startTime; 742 } catch (Exception e) { 743 throw new RuntimeException(e); 744 } 745 }); 746 747 // Wait for both threads to complete 748 configUpdateDuration[0] = configUpdateFuture.get(); 749 balanceDuration[0] = balanceFuture.get(); 750 751 // Verify that config update didn't hang/timeout waiting for balance 752 assertTrue(configUpdateDuration[0] < balanceDuration[0]); 753 754 // Verify that ratio threshold used during balance is stll the old 755 assertEquals(expectedOldRatioThreshold, actualOldRatioThresholdDuringBalance[0], 0.001f); 756 757 // Verify that config updated successfully after balance completed 758 internalBalancer = (CacheAwareLoadBalancer) balancer.getInternalBalancer(); 759 assertEquals(expectedNewRatioThreshold, internalBalancer.ratioThreshold, 0.001f); 760 761 } finally { 762 executor.shutdownNow(); 763 } 764 } 765}