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}