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}