diff --git a/hbase-server/src/main/java/org/apache/hadoop/hbase/master/ServerManager.java b/hbase-server/src/main/java/org/apache/hadoop/hbase/master/ServerManager.java index 1ea500fe36a3..344ffa396b96 100644 --- a/hbase-server/src/main/java/org/apache/hadoop/hbase/master/ServerManager.java +++ b/hbase-server/src/main/java/org/apache/hadoop/hbase/master/ServerManager.java @@ -1092,6 +1092,24 @@ public void removeRegion(final RegionInfo regionInfo) { flushedSequenceIdByRegion.remove(encodedName); } + /** + * Called on region OPEN to seed {@link #flushedSequenceIdByRegion} with the region's + * {@code openSeqNum}. Without this, the entry stays absent until the hosting server's next + * heartbeat, so {@link #getLastFlushedSequenceId} returns {@link HConstants#NO_SEQNUM} and + * WALSplitter conservatively treats already-durable edits as unflushed - producing orphaned + * recovered.edits when the source server crashes soon after a drain-move. Uses {@code merge} with + * {@link Math#max} so a heartbeat-supplied value (which may reflect flushes after open) is never + * regressed - and, unlike {@code putIfAbsent}, a stale-low prior value is lifted to + * {@code openSeqNum}. Safe because at OPEN a region cannot have flushed past its own + * {@code openSeqNum}. See HBASE-30335. + */ + public void reportRegionOpen(final RegionInfo regionInfo, final long openSeqNum) { + if (openSeqNum < 0) { // NO_SEQNUM == -1 + return; + } + flushedSequenceIdByRegion.merge(regionInfo.getEncodedNameAsBytes(), openSeqNum, Math::max); + } + public boolean isRegionInServerManagerStates(final RegionInfo hri) { final byte[] encodedName = hri.getEncodedNameAsBytes(); return (storeFlushedSequenceIdsByRegion.containsKey(encodedName) diff --git a/hbase-server/src/main/java/org/apache/hadoop/hbase/master/assignment/AssignmentManager.java b/hbase-server/src/main/java/org/apache/hadoop/hbase/master/assignment/AssignmentManager.java index 5baf30846e08..021f8cdbb1f1 100644 --- a/hbase-server/src/main/java/org/apache/hadoop/hbase/master/assignment/AssignmentManager.java +++ b/hbase-server/src/main/java/org/apache/hadoop/hbase/master/assignment/AssignmentManager.java @@ -2287,6 +2287,10 @@ void regionOpenedWithoutPersistingToMeta(RegionStateNode regionNode) RegionInfo regionInfo = regionNode.getRegionInfo(); regionStates.addRegionToServer(regionNode); regionStates.removeFromFailedOpen(regionInfo); + // HBASE-30335: seed the master's flushed sequence cache with openSeqNum so a subsequent + // WAL split (e.g. source RS crashes after drain-move) recognizes already-durable edits + // instead of writing orphaned recovered.edits. + master.getServerManager().reportRegionOpen(regionInfo, regionNode.getOpenSeqNum()); } // should be called under the RegionStateNode lock diff --git a/hbase-server/src/test/java/org/apache/hadoop/hbase/master/AbstractTestDLS.java b/hbase-server/src/test/java/org/apache/hadoop/hbase/master/AbstractTestDLS.java index c63642863ed1..b50a3b218574 100644 --- a/hbase-server/src/test/java/org/apache/hadoop/hbase/master/AbstractTestDLS.java +++ b/hbase-server/src/test/java/org/apache/hadoop/hbase/master/AbstractTestDLS.java @@ -55,6 +55,7 @@ import org.apache.hadoop.hbase.client.Table; import org.apache.hadoop.hbase.coordination.ZKSplitLogManagerCoordination; import org.apache.hadoop.hbase.master.assignment.RegionStates; +import org.apache.hadoop.hbase.regionserver.HRegion; import org.apache.hadoop.hbase.regionserver.HRegionServer; import org.apache.hadoop.hbase.regionserver.MultiVersionConcurrencyControl; import org.apache.hadoop.hbase.regionserver.Region; @@ -415,6 +416,18 @@ public void makeWAL(HRegionServer hrs, List regions, int numEdits, i // sync every ~30k to line up with desired wal rolls final int syncEvery = 30 * 1024 / editSize; MultiVersionConcurrencyControl mvcc = new MultiVersionConcurrencyControl(); + // HBASE-30335: match the per-region seqid invariant a real WAL preserves so the splitter's + // openSeqNum-seeded filter doesn't drop our injected edits as already-flushed. + long maxOpen = 0L; + for (RegionInfo info : hris) { + HRegion r = hrs.getRegion(info.getEncodedName()); + if (r != null) { + maxOpen = Math.max(maxOpen, r.getOpenSeqNum()); + } + } + if (maxOpen > 0L) { + mvcc.advanceTo(maxOpen); + } if (n > 0) { for (int i = 0; i < numEdits; i += 1) { WALEdit e = new WALEdit(); diff --git a/hbase-server/src/test/java/org/apache/hadoop/hbase/master/TestGetLastFlushedSequenceId.java b/hbase-server/src/test/java/org/apache/hadoop/hbase/master/TestGetLastFlushedSequenceId.java index b7c533a859ca..f2b76a07065e 100644 --- a/hbase-server/src/test/java/org/apache/hadoop/hbase/master/TestGetLastFlushedSequenceId.java +++ b/hbase-server/src/test/java/org/apache/hadoop/hbase/master/TestGetLastFlushedSequenceId.java @@ -18,6 +18,7 @@ package org.apache.hadoop.hbase.master; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotEquals; import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -30,6 +31,7 @@ import org.apache.hadoop.hbase.TableName; import org.apache.hadoop.hbase.client.Put; import org.apache.hadoop.hbase.client.Table; +import org.apache.hadoop.hbase.regionserver.HRegion; import org.apache.hadoop.hbase.regionserver.HRegionServer; import org.apache.hadoop.hbase.regionserver.Region; import org.apache.hadoop.hbase.testclassification.MediumTests; @@ -89,10 +91,12 @@ public void test() throws IOException, InterruptedException { Thread.sleep(2000); RegionStoreSequenceIds ids = testUtil.getHBaseCluster().getMaster().getServerManager() .getLastFlushedSequenceId(region.getRegionInfo().getEncodedNameAsBytes()); - assertEquals(HConstants.NO_SEQNUM, ids.getLastFlushedSequenceId()); // This will be the sequenceid just before that of the earliest edit in memstore. long storeSequenceId = ids.getStoreSequenceId(0).getSequenceId(); assertTrue(storeSequenceId > 0); + // HBASE-30335: openSeqNum is now seeded on region OPEN, so lastFlushedSequenceId is no + // longer NO_SEQNUM before the first flush. + assertNotEquals(HConstants.NO_SEQNUM, ids.getLastFlushedSequenceId()); testUtil.getAdmin().flush(tableName); Thread.sleep(2000); ids = testUtil.getHBaseCluster().getMaster().getServerManager() @@ -102,4 +106,41 @@ public void test() throws IOException, InterruptedException { assertEquals(ids.getLastFlushedSequenceId(), ids.getStoreSequenceId(0).getSequenceId()); table.close(); } + + /** + * HBASE-30335: after a region is opened - and before any user write or flush - the master's + * flushedSequenceIdByRegion must already contain the region's openSeqNum. Otherwise a subsequent + * WAL split (e.g. the hosting RS crashes before its first flush heartbeat) would treat + * already-durable edits as unflushed and produce orphaned recovered.edits. + */ + @Test + public void testFlushedSequenceIdSeededOnRegionOpen() throws IOException, InterruptedException { + TableName freshTable = TableName.valueOf(getClass().getSimpleName(), "openseed"); + testUtil.getAdmin() + .createNamespace(NamespaceDescriptor.create(freshTable.getNamespaceAsString()).build()); + Table table = testUtil.createTable(freshTable, families); + try { + SingleProcessHBaseCluster cluster = testUtil.getMiniHBaseCluster(); + HRegion region = null; + for (JVMClusterUtil.RegionServerThread rst : cluster.getRegionServerThreads()) { + for (HRegion r : rst.getRegionServer().getRegions(freshTable)) { + region = r; + break; + } + if (region != null) { + break; + } + } + assertNotNull(region); + long openSeqNum = region.getOpenSeqNum(); + RegionStoreSequenceIds ids = testUtil.getHBaseCluster().getMaster().getServerManager() + .getLastFlushedSequenceId(region.getRegionInfo().getEncodedNameAsBytes()); + assertNotEquals(HConstants.NO_SEQNUM, ids.getLastFlushedSequenceId(), + "flushedSequenceIdByRegion should be seeded on region OPEN (HBASE-30335)"); + assertEquals(openSeqNum, ids.getLastFlushedSequenceId(), + "seeded value must equal the region's openSeqNum"); + } finally { + table.close(); + } + } } diff --git a/hbase-server/src/test/java/org/apache/hadoop/hbase/master/TestServerManager.java b/hbase-server/src/test/java/org/apache/hadoop/hbase/master/TestServerManager.java new file mode 100644 index 000000000000..3cbeaf265f9c --- /dev/null +++ b/hbase-server/src/test/java/org/apache/hadoop/hbase/master/TestServerManager.java @@ -0,0 +1,99 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.hadoop.hbase.master; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.hbase.HBaseConfiguration; +import org.apache.hadoop.hbase.HConstants; +import org.apache.hadoop.hbase.TableName; +import org.apache.hadoop.hbase.client.RegionInfo; +import org.apache.hadoop.hbase.client.RegionInfoBuilder; +import org.apache.hadoop.hbase.master.assignment.AssignmentManager; +import org.apache.hadoop.hbase.master.assignment.RegionStates; +import org.apache.hadoop.hbase.testclassification.MasterTests; +import org.apache.hadoop.hbase.testclassification.SmallTests; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Tag; +import org.junit.jupiter.api.Test; + +@Tag(MasterTests.TAG) +@Tag(SmallTests.TAG) +public class TestServerManager { + + private static final class DummyMasterServices extends MockNoopMasterServices { + private final AssignmentManager am; + + DummyMasterServices(Configuration conf) { + super(conf); + am = mock(AssignmentManager.class); + RegionStates rss = mock(RegionStates.class); + when(am.getRegionStates()).thenReturn(rss); + } + + @Override + public AssignmentManager getAssignmentManager() { + return am; + } + } + + private ServerManager sm; + private RegionInfo region; + + @BeforeEach + public void setUp() { + Configuration conf = HBaseConfiguration.create(); + sm = new ServerManager(new DummyMasterServices(conf), new DummyRegionServerList()); + region = RegionInfoBuilder.newBuilder(TableName.valueOf("t")).build(); + } + + private long lastFlushed(RegionInfo ri) { + return sm.getLastFlushedSequenceId(ri.getEncodedNameAsBytes()).getLastFlushedSequenceId(); + } + + @Test + public void testReportRegionOpenSeedsFlushedSequenceId() { + assertEquals(HConstants.NO_SEQNUM, lastFlushed(region)); + sm.reportRegionOpen(region, 42L); + assertEquals(42L, lastFlushed(region)); + } + + @Test + public void testReportRegionOpenDoesNotRegressExistingValue() { + sm.reportRegionOpen(region, 100L); + // A later OPEN carrying a smaller openSeqNum (e.g. after a restart replayed less) must not + // clobber a higher watermark already seeded here or supplied by a heartbeat. + sm.reportRegionOpen(region, 50L); + assertEquals(100L, lastFlushed(region)); + } + + @Test + public void testReportRegionOpenIgnoresNoSeqNum() { + sm.reportRegionOpen(region, HConstants.NO_SEQNUM); + assertEquals(HConstants.NO_SEQNUM, lastFlushed(region)); + } + + @Test + public void testReportRegionOpenIgnoresNegativeSeqNum() { + sm.reportRegionOpen(region, -5L); + assertEquals(HConstants.NO_SEQNUM, lastFlushed(region)); + } +}