From 015574f3952e6efe8b423f756eafe3ef6ab09f22 Mon Sep 17 00:00:00 2001 From: Xinyi Yan Date: Tue, 17 Aug 2021 17:40:01 -0700 Subject: [PATCH] PHOENIX-6472 In case of region inconsistency phoenix should stop gracefully --- .../query/ConnectionQueryServicesImpl.java | 17 +++++- .../ConnectionQueryServicesImplTest.java | 59 +++++++++++++++++++ 2 files changed, 75 insertions(+), 1 deletion(-) diff --git a/phoenix-core/src/main/java/org/apache/phoenix/query/ConnectionQueryServicesImpl.java b/phoenix-core/src/main/java/org/apache/phoenix/query/ConnectionQueryServicesImpl.java index e0c252ed8b6..dca8b7441f4 100644 --- a/phoenix-core/src/main/java/org/apache/phoenix/query/ConnectionQueryServicesImpl.java +++ b/phoenix-core/src/main/java/org/apache/phoenix/query/ConnectionQueryServicesImpl.java @@ -81,6 +81,7 @@ import java.io.IOException; import java.lang.management.ManagementFactory; import java.lang.ref.WeakReference; +import java.nio.charset.StandardCharsets; import java.sql.PreparedStatement; import java.sql.ResultSetMetaData; import java.sql.SQLException; @@ -648,6 +649,20 @@ public void clearTableRegionCache(byte[] tableName) throws SQLException { connection.clearRegionCache(TableName.valueOf(tableName)); } + public byte[] getNextRegionStartKey(HRegionLocation regionLocation, byte[] currentKey) throws IOException { + // in order to check the overlap/inconsistencies bad region info, we have to make sure + // the current endKey always increasing(compare the previous endKey) + if (Bytes.compareTo(regionLocation.getRegionInfo().getEndKey(), currentKey) <= 0 + && !Bytes.equals(currentKey, HConstants.EMPTY_START_ROW) + && !Bytes.equals(regionLocation.getRegionInfo().getEndKey(), HConstants.EMPTY_END_ROW)) { + String regionNameString = + new String(regionLocation.getRegionInfo().getRegionName(), StandardCharsets.UTF_8); + throw new IOException(String.format( + "HBase region information overlap/inconsistencies on region %s", regionNameString)); + } + return regionLocation.getRegionInfo().getEndKey(); + } + @Override public List getAllTableRegions(byte[] tableName) throws SQLException { /* @@ -666,8 +681,8 @@ public List getAllTableRegions(byte[] tableName) throws SQLExce do { HRegionLocation regionLocation = connection.getRegionLocation( TableName.valueOf(tableName), currentKey, reload); + currentKey = getNextRegionStartKey(regionLocation, currentKey); locations.add(regionLocation); - currentKey = regionLocation.getRegionInfo().getEndKey(); } while (!Bytes.equals(currentKey, HConstants.EMPTY_END_ROW)); return locations; } catch (org.apache.hadoop.hbase.TableNotFoundException e) { diff --git a/phoenix-core/src/test/java/org/apache/phoenix/query/ConnectionQueryServicesImplTest.java b/phoenix-core/src/test/java/org/apache/phoenix/query/ConnectionQueryServicesImplTest.java index 1b94f94c30e..a845ca919d9 100644 --- a/phoenix-core/src/test/java/org/apache/phoenix/query/ConnectionQueryServicesImplTest.java +++ b/phoenix-core/src/test/java/org/apache/phoenix/query/ConnectionQueryServicesImplTest.java @@ -46,6 +46,8 @@ import org.apache.hadoop.hbase.HColumnDescriptor; import org.apache.hadoop.hbase.HConstants; import org.apache.hadoop.hbase.HTableDescriptor; +import org.apache.hadoop.hbase.HRegionInfo; +import org.apache.hadoop.hbase.HRegionLocation; import org.apache.hadoop.hbase.TableName; import org.apache.hadoop.hbase.TableNotFoundException; import org.apache.hadoop.hbase.client.ClusterConnection; @@ -157,6 +159,63 @@ public void testExceptionHandlingOnSystemNamespaceCreation() throws Exception { } } + @Test + public void testGetNextRegionStartKey() { + HRegionInfo mockHRegionInfo = org.mockito.Mockito.mock(HRegionInfo.class); + HRegionLocation mockRegionLocation = org.mockito.Mockito.mock(HRegionLocation.class); + ConnectionQueryServicesImpl mockCqsi = org.mockito.Mockito.mock(ConnectionQueryServicesImpl.class, + org.mockito.Mockito.CALLS_REAL_METHODS); + byte[] corruptedStartAndEndKey = "0x3000".getBytes(); + byte[] corruptedDecreasingKey = "0x2999".getBytes(); + byte[] notCorruptedStartKey = "0x2999".getBytes(); + byte[] notCorruptedEndKey = "0x3000".getBytes(); + byte[] notCorruptedNewKey = "0x3001".getBytes(); + byte[] mockTableName = "dummyTable".getBytes(); + when(mockRegionLocation.getRegionInfo()).thenReturn(mockHRegionInfo); + when(mockHRegionInfo.getRegionName()).thenReturn(mockTableName); + + // comparing the current regionInfo endKey is equal to the previous endKey + // [0x3000, Ox3000) vs 0x3000 + when(mockHRegionInfo.getStartKey()).thenReturn(corruptedStartAndEndKey); + when(mockHRegionInfo.getEndKey()).thenReturn(corruptedStartAndEndKey); + testGetNextRegionStartKey(mockCqsi, mockRegionLocation, corruptedStartAndEndKey, true); + + // comparing the current regionInfo endKey is less than previous endKey + // [0x3000,0x2999) vs 0x3000 + when(mockHRegionInfo.getStartKey()).thenReturn(corruptedStartAndEndKey); + when(mockHRegionInfo.getEndKey()).thenReturn(corruptedDecreasingKey); + testGetNextRegionStartKey(mockCqsi, mockRegionLocation, corruptedStartAndEndKey, true); + + // comparing the current regionInfo endKey is greater than the previous endKey + // [0x3000,0x3000) vs 0x3001 + when(mockHRegionInfo.getStartKey()).thenReturn(notCorruptedStartKey); + when(mockHRegionInfo.getEndKey()).thenReturn(notCorruptedNewKey); + testGetNextRegionStartKey(mockCqsi, mockRegionLocation, notCorruptedEndKey, false); + + // test EMPTY_START_ROW + when(mockHRegionInfo.getStartKey()).thenReturn(HConstants.EMPTY_START_ROW); + when(mockHRegionInfo.getEndKey()).thenReturn(notCorruptedEndKey); + testGetNextRegionStartKey(mockCqsi, mockRegionLocation, HConstants.EMPTY_START_ROW, false); + + //test EMPTY_END_ROW + when(mockHRegionInfo.getStartKey()).thenReturn(notCorruptedStartKey); + when(mockHRegionInfo.getEndKey()).thenReturn(HConstants.EMPTY_END_ROW); + testGetNextRegionStartKey(mockCqsi, mockRegionLocation, notCorruptedStartKey, false); + } + + private void testGetNextRegionStartKey(ConnectionQueryServicesImpl mockCqsi, + HRegionLocation mockRegionLocation, byte[] key, boolean isCorrupted) { + try { + mockCqsi.getNextRegionStartKey(mockRegionLocation, key); + if (isCorrupted) { + fail(); + } + } catch (IOException e) { + if (!isCorrupted) { + fail(); + } + } + } @Test public void testSysMutexCheckReturnsFalseWhenTableAbsent() throws Exception { // Override the getTableDescriptor() call to throw instead