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.util;
019
020import static org.apache.hadoop.hbase.master.HMaster.HBASE_MASTER_RSPROC_DISPATCHER_CLASS;
021import static org.junit.jupiter.api.Assertions.assertEquals;
022
023import java.util.List;
024import java.util.stream.Collectors;
025import java.util.stream.IntStream;
026import org.apache.hadoop.hbase.HBaseTestingUtil;
027import org.apache.hadoop.hbase.ServerName;
028import org.apache.hadoop.hbase.SingleProcessHBaseCluster;
029import org.apache.hadoop.hbase.TableName;
030import org.apache.hadoop.hbase.client.Admin;
031import org.apache.hadoop.hbase.client.ColumnFamilyDescriptorBuilder;
032import org.apache.hadoop.hbase.client.Put;
033import org.apache.hadoop.hbase.client.Table;
034import org.apache.hadoop.hbase.client.TableDescriptor;
035import org.apache.hadoop.hbase.client.TableDescriptorBuilder;
036import org.apache.hadoop.hbase.master.HMaster;
037import org.apache.hadoop.hbase.master.hbck.HbckChore;
038import org.apache.hadoop.hbase.master.hbck.HbckReport;
039import org.apache.hadoop.hbase.master.procedure.ServerCrashProcedure;
040import org.apache.hadoop.hbase.regionserver.HRegion;
041import org.apache.hadoop.hbase.regionserver.HRegionServer;
042import org.apache.hadoop.hbase.testclassification.LargeTests;
043import org.apache.hadoop.hbase.testclassification.MiscTests;
044import org.junit.jupiter.api.AfterAll;
045import org.junit.jupiter.api.AfterEach;
046import org.junit.jupiter.api.BeforeAll;
047import org.junit.jupiter.api.BeforeEach;
048import org.junit.jupiter.api.Tag;
049import org.junit.jupiter.api.Test;
050import org.junit.jupiter.api.TestInfo;
051import org.slf4j.Logger;
052import org.slf4j.LoggerFactory;
053
054import org.apache.hadoop.hbase.shaded.protobuf.generated.ProcedureProtos;
055
056/**
057 * Testing custom RSProcedureDispatcher to ensure retry limit can be imposed on certain errors.
058 */
059@Tag(MiscTests.TAG)
060@Tag(LargeTests.TAG)
061public class TestProcDispatcher {
062
063  private static final Logger LOG = LoggerFactory.getLogger(TestProcDispatcher.class);
064
065  private static final HBaseTestingUtil TEST_UTIL = new HBaseTestingUtil();
066  private static ServerName rs0;
067
068  @BeforeAll
069  public static void setUpBeforeClass() throws Exception {
070    TEST_UTIL.getConfiguration().set(HBASE_MASTER_RSPROC_DISPATCHER_CLASS,
071      RSProcDispatcher.class.getName());
072    TEST_UTIL.getConfiguration().setInt(RSProcDispatcher.FAIL_FAST_LIMIT_KEY, 5);
073    TEST_UTIL.startMiniCluster(3);
074    SingleProcessHBaseCluster cluster = TEST_UTIL.getHBaseCluster();
075    rs0 = cluster.getRegionServer(0).getServerName();
076    TEST_UTIL.getAdmin().balancerSwitch(false, true);
077  }
078
079  @AfterAll
080  public static void tearDownAfterClass() throws Exception {
081    TEST_UTIL.shutdownMiniCluster();
082  }
083
084  @BeforeEach
085  public void setUp(TestInfo testInfo) throws Exception {
086    final TableName tableName = TableName.valueOf(testInfo.getTestMethod().get().getName());
087    TableDescriptor tableDesc = TableDescriptorBuilder.newBuilder(tableName)
088      .setColumnFamily(ColumnFamilyDescriptorBuilder.of("fam1")).build();
089    int startKey = 0;
090    int endKey = 80000;
091    TEST_UTIL.getAdmin().createTable(tableDesc, Bytes.toBytes(startKey), Bytes.toBytes(endKey), 9);
092  }
093
094  @AfterEach
095  public void tearDown() {
096    RSProcDispatcher.stopInjecting();
097  }
098
099  @Test
100  public void testRetryLimitOnConnClosedErrors(TestInfo testInfo) throws Exception {
101    HbckChore hbckChore = new HbckChore(TEST_UTIL.getHBaseCluster().getMaster());
102    final TableName tableName = TableName.valueOf(testInfo.getTestMethod().get().getName());
103    SingleProcessHBaseCluster cluster = TEST_UTIL.getHBaseCluster();
104    Admin admin = TEST_UTIL.getAdmin();
105    List<Put> puts = IntStream.range(10, 50000).mapToObj(i -> new Put(Bytes.toBytes(i))
106      .addColumn(Bytes.toBytes("fam1"), Bytes.toBytes("q1"), Bytes.toBytes("val_" + i)))
107      .collect(Collectors.toList());
108    try (Table table = TEST_UTIL.getConnection().getTable(tableName)) {
109      table.put(puts);
110    }
111    admin.flush(tableName);
112    admin.compact(tableName);
113    Thread.sleep(3000);
114    HRegionServer hRegionServer0 = cluster.getRegionServer(0);
115    HRegionServer hRegionServer1 = cluster.getRegionServer(1);
116    HRegionServer hRegionServer2 = cluster.getRegionServer(2);
117    int numRegions0 = hRegionServer0.getNumberOfOnlineRegions();
118    int numRegions1 = hRegionServer1.getNumberOfOnlineRegions();
119    int numRegions2 = hRegionServer2.getNumberOfOnlineRegions();
120
121    hbckChore.choreForTesting();
122    HbckReport hbckReport = hbckChore.getLastReport();
123    assertEquals(0, hbckReport.getInconsistentRegions().size());
124    assertEquals(0, hbckReport.getOrphanRegionsOnFS().size());
125    assertEquals(0, hbckReport.getOrphanRegionsOnRS().size());
126
127    HRegion region0 = hRegionServer0.getRegions().get(0);
128    // Fail the next two open/close-region requests for this table so the moves trigger SCP(s).
129    RSProcDispatcher.injectErrorsForNextRequests(tableName, 2);
130    // move all regions from server1 to server0
131    for (HRegion region : hRegionServer1.getRegions()) {
132      TEST_UTIL.getAdmin().move(region.getRegionInfo().getEncodedNameAsBytes(), rs0);
133    }
134    TEST_UTIL.getAdmin().move(region0.getRegionInfo().getEncodedNameAsBytes());
135    HMaster master = TEST_UTIL.getHBaseCluster().getMaster();
136
137    // Ensure, after the injected connection errors:
138    // 1. the total number of regions is unchanged before and after the SCP(s)
139    // 2. all procedures (including the SCP(s)) complete successfully
140    // 3. at least one ServerCrashProcedure was scheduled
141    TEST_UTIL.waitFor(60000, 1000, () -> {
142      LOG.info("numRegions0: {} , numRegions1: {} , numRegions2: {}", numRegions0, numRegions1,
143        numRegions2);
144      LOG.info("Online regions - server0 : {} , server1: {} , server2: {}",
145        cluster.getRegionServer(0).getNumberOfOnlineRegions(),
146        cluster.getRegionServer(1).getNumberOfOnlineRegions(),
147        cluster.getRegionServer(2).getNumberOfOnlineRegions());
148      LOG.info("Num of successfully completed procedures: {} , num of all procedures: {}",
149        master.getMasterProcedureExecutor().getProcedures().stream()
150          .filter(masterProcedureEnvProcedure -> masterProcedureEnvProcedure.getState()
151              == ProcedureProtos.ProcedureState.SUCCESS)
152          .count(),
153        master.getMasterProcedureExecutor().getProcedures().size());
154      LOG.info("Num of SCPs: {}", master.getMasterProcedureExecutor().getProcedures().stream()
155        .filter(proc -> proc instanceof ServerCrashProcedure).count());
156      return (numRegions0 + numRegions1 + numRegions2)
157          == (cluster.getRegionServer(0).getNumberOfOnlineRegions()
158            + cluster.getRegionServer(1).getNumberOfOnlineRegions()
159            + cluster.getRegionServer(2).getNumberOfOnlineRegions())
160        && master.getMasterProcedureExecutor().getProcedures().stream()
161          .filter(masterProcedureEnvProcedure -> masterProcedureEnvProcedure.getState()
162              == ProcedureProtos.ProcedureState.SUCCESS)
163          .count() == master.getMasterProcedureExecutor().getProcedures().size()
164        && master.getMasterProcedureExecutor().getProcedures().stream()
165          .anyMatch(proc -> proc instanceof ServerCrashProcedure);
166    });
167
168    // Ensure we have no inconsistent regions
169    TEST_UTIL.waitFor(60000, 1000, () -> {
170      hbckChore.choreForTesting();
171      HbckReport report = hbckChore.getLastReport();
172      return report.getInconsistentRegions().isEmpty() && report.getOrphanRegionsOnFS().isEmpty()
173        && report.getOrphanRegionsOnRS().isEmpty();
174    });
175  }
176}