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..fced737bca90 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 @@ -289,35 +289,50 @@ ServerName regionServerStartup(RegionServerStartupRequest request, int versionNu private void updateLastFlushedSequenceIds(ServerName sn, ServerMetrics hsl) { for (Entry entry : hsl.getRegionMetrics().entrySet()) { byte[] encodedRegionName = Bytes.toBytes(RegionInfo.encodeRegionName(entry.getKey())); - Long existingValue = flushedSequenceIdByRegion.get(encodedRegionName); - long l = entry.getValue().getCompletedSequenceId(); - // Don't let smaller sequence ids override greater sequence ids. - if (LOG.isTraceEnabled()) { - LOG.trace(Bytes.toString(encodedRegionName) + ", existingValue=" + existingValue - + ", completeSequenceId=" + l); - } - if (existingValue == null || (l != HConstants.NO_SEQNUM && l > existingValue)) { - flushedSequenceIdByRegion.put(encodedRegionName, l); - } else if (l != HConstants.NO_SEQNUM && l < existingValue) { - LOG.warn("RegionServer " + sn + " indicates a last flushed sequence id (" + l - + ") that is less than the previous last flushed sequence id (" + existingValue - + ") for region " + Bytes.toString(entry.getKey()) + " Ignoring."); - } + final long completedSeqId = entry.getValue().getCompletedSequenceId(); + // Atomic read-modify-write so a concurrent reportRegionOpen seed (which uses + // merge(Math::max)) cannot be clobbered by a stale in-flight heartbeat carrying a + // lower completedSequenceId. Don't let smaller sequence ids override greater ones. + flushedSequenceIdByRegion.compute(encodedRegionName, (k, existingValue) -> { + if (LOG.isTraceEnabled()) { + LOG.trace(Bytes.toString(k) + ", existingValue=" + existingValue + ", completeSequenceId=" + + completedSeqId); + } + if (existingValue == null) { + return completedSeqId; + } + if (completedSeqId != HConstants.NO_SEQNUM && completedSeqId > existingValue) { + return completedSeqId; + } + if (completedSeqId != HConstants.NO_SEQNUM && completedSeqId < existingValue) { + LOG.warn("RegionServer " + sn + " indicates a last flushed sequence id (" + completedSeqId + + ") that is less than the previous last flushed sequence id (" + existingValue + + ") for region " + Bytes.toString(entry.getKey()) + " Ignoring."); + } + return existingValue; + }); ConcurrentNavigableMap storeFlushedSequenceId = computeIfAbsent(storeFlushedSequenceIdsByRegion, encodedRegionName, () -> new ConcurrentSkipListMap<>(Bytes.BYTES_COMPARATOR)); for (Entry storeSeqId : entry.getValue().getStoreSequenceId().entrySet()) { byte[] family = storeSeqId.getKey(); - existingValue = storeFlushedSequenceId.get(family); - l = storeSeqId.getValue(); - if (LOG.isTraceEnabled()) { - LOG.trace(Bytes.toString(encodedRegionName) + ", family=" + Bytes.toString(family) - + ", existingValue=" + existingValue + ", completeSequenceId=" + l); - } - // Don't let smaller sequence ids override greater sequence ids. - if (existingValue == null || (l != HConstants.NO_SEQNUM && l > existingValue.longValue())) { - storeFlushedSequenceId.put(family, l); - } + final long storeCompletedSeqId = storeSeqId.getValue(); + storeFlushedSequenceId.compute(family, (k, existingValue) -> { + if (LOG.isTraceEnabled()) { + LOG.trace(Bytes.toString(encodedRegionName) + ", family=" + Bytes.toString(k) + + ", existingValue=" + existingValue + ", completeSequenceId=" + storeCompletedSeqId); + } + if (existingValue == null) { + return storeCompletedSeqId; + } + if ( + storeCompletedSeqId != HConstants.NO_SEQNUM + && storeCompletedSeqId > existingValue.longValue() + ) { + return storeCompletedSeqId; + } + return existingValue; + }); } } } @@ -1092,6 +1107,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 a9ae6f72a29b..9df9053d1892 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 @@ -1230,8 +1230,17 @@ private void reportRegionStateTransition(ReportRegionStateTransitionResponse.Bui final RegionInfo hri = ProtobufUtil.toRegionInfo(transition.getRegionInfo(0)); long procId = transition.getProcIdCount() > 0 ? transition.getProcId(0) : Procedure.NO_PROC_ID; - updateRegionTransition(serverNode, transition.getTransitionCode(), hri, - transition.hasOpenSeqNum() ? transition.getOpenSeqNum() : HConstants.NO_SEQNUM, procId); + long seqId = + transition.hasOpenSeqNum() ? transition.getOpenSeqNum() : HConstants.NO_SEQNUM; + // On CLOSE, seed the master's flushed-seqid watermark with the region's final flushed + // seqid (reported by the RS). This complements the OPEN-time seed (HBASE-30335) and + // narrows the graceful-close window where a WAL split of a crashed source RS could + // write orphaned recovered.edits for already-durable edits. Uses the same monotonic + // merge(Math::max) seed, so it never regresses a higher value. + if (transition.getTransitionCode() == TransitionCode.CLOSED && seqId >= 0) { + master.getServerManager().reportRegionOpen(hri, seqId); + } + updateRegionTransition(serverNode, transition.getTransitionCode(), hri, seqId, procId); break; case READY_TO_SPLIT: case SPLIT: @@ -2292,6 +2301,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/main/java/org/apache/hadoop/hbase/regionserver/HRegionServer.java b/hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/HRegionServer.java index e11b334442f1..30105cd82b17 100644 --- a/hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/HRegionServer.java +++ b/hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/HRegionServer.java @@ -2348,7 +2348,10 @@ private boolean skipReportingTransition(final RegionStateTransitionContext conte builder.setServer(ProtobufUtil.toServerName(serverName)); RegionStateTransition.Builder transition = builder.addTransitionBuilder(); transition.setTransitionCode(code); - if (code == TransitionCode.OPENED && openSeqNum >= 0) { + // Carry the seqid for OPENED (openSeqNum) and CLOSED (durable flushed seqid). The master uses + // both to seed its flushedSequenceIdByRegion watermark; see HBASE-30335 and the CLOSE-time + // follow-up. + if ((code == TransitionCode.OPENED || code == TransitionCode.CLOSED) && openSeqNum >= 0) { transition.setOpenSeqNum(openSeqNum); } for (RegionInfo hri : hris) { diff --git a/hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/handler/CloseRegionHandler.java b/hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/handler/CloseRegionHandler.java index f18e7d9ba635..c1ddd6ba5d8c 100644 --- a/hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/handler/CloseRegionHandler.java +++ b/hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/handler/CloseRegionHandler.java @@ -18,7 +18,6 @@ package org.apache.hadoop.hbase.regionserver.handler; import java.io.IOException; -import org.apache.hadoop.hbase.HConstants; import org.apache.hadoop.hbase.Server; import org.apache.hadoop.hbase.ServerName; import org.apache.hadoop.hbase.client.RegionInfo; @@ -110,8 +109,12 @@ public void process() throws IOException { } this.rsServices.removeRegion(region, destination); + // Report the region's durable flushed seqid on CLOSE so the master can seed its + // flushedSequenceIdByRegion watermark. This complements the OPEN-time seed (HBASE-30335) + // and narrows the graceful-close window where a subsequent WAL split of a crashed source + // RS could otherwise write orphaned recovered.edits for already-durable edits. rsServices.reportRegionStateTransition(new RegionStateTransitionContext(TransitionCode.CLOSED, - HConstants.NO_SEQNUM, Procedure.NO_PROC_ID, -1, regionInfo, -1)); + region.getMaxFlushedSeqId(), Procedure.NO_PROC_ID, -1, regionInfo, -1)); // Done! Region is closed on this RS LOG.debug("Closed {}", region.getRegionInfo().getRegionNameAsString()); diff --git a/hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/handler/UnassignRegionHandler.java b/hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/handler/UnassignRegionHandler.java index 815d8bed3322..542e04d117e5 100644 --- a/hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/handler/UnassignRegionHandler.java +++ b/hbase-server/src/main/java/org/apache/hadoop/hbase/regionserver/handler/UnassignRegionHandler.java @@ -25,7 +25,6 @@ import edu.umd.cs.findbugs.annotations.Nullable; import java.io.IOException; import java.util.concurrent.TimeUnit; -import org.apache.hadoop.hbase.HConstants; import org.apache.hadoop.hbase.ServerName; import org.apache.hadoop.hbase.executor.EventHandler; import org.apache.hadoop.hbase.executor.EventType; @@ -145,9 +144,14 @@ public void process() throws IOException { } rs.removeRegion(region, destination); + // Report the region's durable flushed seqid on CLOSE so the master can seed its + // flushedSequenceIdByRegion watermark. This complements the OPEN-time seed (HBASE-30335) + // and narrows the graceful-close window where a subsequent WAL split of a crashed source + // RS could otherwise write orphaned recovered.edits for already-durable edits. if ( - !rs.reportRegionStateTransition(new RegionStateTransitionContext(TransitionCode.CLOSED, - HConstants.NO_SEQNUM, closeProcId, -1, region.getRegionInfo(), initiatingMasterActiveTime)) + !rs.reportRegionStateTransition( + new RegionStateTransitionContext(TransitionCode.CLOSED, region.getMaxFlushedSeqId(), + closeProcId, -1, region.getRegionInfo(), initiatingMasterActiveTime)) ) { throw new IOException("Failed to report close to master: " + regionName); } 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..5bef3e4da540 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,97 @@ 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(); + } + } + + /** + * HBASE-30335 follow-up: on graceful region CLOSE the RS reports the region's final durable + * flushed seqid, and the master must lift its flushedSequenceIdByRegion watermark to that value. + * This complements the OPEN-time seed and covers the drain-move / disable window: if the source + * RS crashes after the close, a subsequent WAL split must see the already-durable edits as + * flushed and not resurrect them as orphaned recovered.edits. Writes and flushes so the region's + * maxFlushedSeqId advances well past its openSeqNum, then disables the table (a graceful close + * that keeps the region - unlike delete/split/merge, disable does not call + * ServerManager#removeRegion) and asserts the watermark reflects the post-close flushed seqid. + */ + @Test + public void testFlushedSequenceIdSeededOnRegionClose() throws IOException, InterruptedException { + TableName freshTable = TableName.valueOf(getClass().getSimpleName(), "closeseed"); + 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(); + // Write and flush a few times so the region's durable flushed seqid advances past openSeqNum; + // this makes the CLOSE-time seed distinguishable from the OPEN-time seed. + for (int i = 0; i < 3; i++) { + table.put(new Put(Bytes.toBytes("k" + i)).addColumn(family, Bytes.toBytes("q"), + Bytes.toBytes("v" + i))); + testUtil.getAdmin().flush(freshTable); + } + long flushedSeqId = region.getMaxFlushedSeqId(); + assertTrue(flushedSeqId > openSeqNum, + "test setup: flushed seqid " + flushedSeqId + " must exceed openSeqNum " + openSeqNum); + byte[] encodedName = region.getRegionInfo().getEncodedNameAsBytes(); + // Gracefully close the region (disable keeps the region entry - it is not removed like a + // delete/split/merge would), driving a CLOSED transition that carries the flushed seqid. + testUtil.getAdmin().disableTable(freshTable); + RegionStoreSequenceIds ids = testUtil.getHBaseCluster().getMaster().getServerManager() + .getLastFlushedSequenceId(encodedName); + assertNotEquals(HConstants.NO_SEQNUM, ids.getLastFlushedSequenceId(), + "flushedSequenceIdByRegion should be seeded on region CLOSE"); + assertTrue(ids.getLastFlushedSequenceId() >= flushedSeqId, + "CLOSE seed " + ids.getLastFlushedSequenceId() + " must be >= the region's flushed seqid " + + flushedSeqId); + } finally { + table.close(); + } + } } diff --git a/hbase-server/src/test/java/org/apache/hadoop/hbase/master/TestMaster.java b/hbase-server/src/test/java/org/apache/hadoop/hbase/master/TestMaster.java index 16e1bf3a2e9e..3d9f058ad26d 100644 --- a/hbase-server/src/test/java/org/apache/hadoop/hbase/master/TestMaster.java +++ b/hbase-server/src/test/java/org/apache/hadoop/hbase/master/TestMaster.java @@ -20,6 +20,7 @@ import static org.junit.jupiter.api.Assertions.assertArrayEquals; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.junit.jupiter.api.Assertions.fail; @@ -252,10 +253,23 @@ public void testFlushedSequenceIdPersistLoad() throws Exception { TEST_UTIL.getMiniHBaseCluster().shutdown(); TEST_UTIL.restartHBaseCluster(2); TEST_UTIL.waitUntilNoRegionsInTransition(); - // check equality after reloading flushed sequence id map + // Post HBASE-30335 the master seeds flushedSequenceIdByRegion on OPEN via + // merge(openSeqNum, Math::max). openSeqNum is monotonic across close/open cycles, so a region + // reopened after cluster restart may carry a strictly higher value than what was persisted. + // The preserved invariant is: every region persisted at shutdown is loaded on restart, and no + // watermark regresses. Map regionMapAfter = TEST_UTIL.getHBaseCluster().getMaster().getServerManager().getFlushedSequenceIdByRegion(); - assertTrue(regionMapBefore.equals(regionMapAfter)); + assertEquals(regionMapBefore.size(), regionMapAfter.size()); + for (Map.Entry before : regionMapBefore.entrySet()) { + Long after = regionMapAfter.get(before.getKey()); + assertNotNull(after, + "region missing after restart: " + Bytes.toStringBinary(before.getKey())); + assertTrue(after >= before.getValue(), + "flushedSequenceId regressed across restart for region " + + Bytes.toStringBinary(before.getKey()) + " before=" + before.getValue() + " after=" + + after); + } } @Test 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..9917dea9c67f --- /dev/null +++ b/hbase-server/src/test/java/org/apache/hadoop/hbase/master/TestServerManager.java @@ -0,0 +1,172 @@ +/* + * 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.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +import java.util.Collections; +import java.util.concurrent.CyclicBarrier; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.hbase.HBaseConfiguration; +import org.apache.hadoop.hbase.HConstants; +import org.apache.hadoop.hbase.RegionMetricsBuilder; +import org.apache.hadoop.hbase.ServerMetrics; +import org.apache.hadoop.hbase.ServerMetricsBuilder; +import org.apache.hadoop.hbase.ServerName; +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.apache.hadoop.hbase.util.Bytes; +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)); + } + + /** + * HBASE-30335: the OPEN-time seed ({@link ServerManager#reportRegionOpen}, which uses + * {@code merge(Math::max)}) and the heartbeat handler ({@link ServerManager#regionServerReport}) + * both write {@code flushedSequenceIdByRegion}. A stale in-flight heartbeat from the + * soon-to-be-dead source RS carries a lower {@code completedSequenceId}. Whatever the + * interleaving, the watermark must never regress below the seed: if the heartbeat lands first the + * seed lifts it to {@code openSeqNum}; if it lands after, the heartbeat's read-modify-write must + * refuse to lower it. Only a non-atomic check-then-put in the heartbeat path (the pre-fix bug) + * could let the stale value clobber the seed. This drives both writers concurrently over many + * rounds to catch that race. + */ + @Test + public void testConcurrentStaleHeartbeatDoesNotClobberOpenSeed() throws Exception { + final long seedSeqId = 200L; + final long staleSeqId = 100L; + // Register the server so regionServerReport takes the heartbeat (updateLastFlushedSequenceIds) + // path instead of the new-server-registration path. + ServerName sn = ServerName.valueOf("rs.example.org", 16020, 1L); + sm.recordNewServerWithLock(sn, ServerMetricsBuilder.of(sn)); + + ExecutorService pool = Executors.newFixedThreadPool(2); + try { + for (int i = 0; i < 500; i++) { + // Fresh region per round so no round is masked by a prior round's watermark. + final RegionInfo ri = RegionInfoBuilder.newBuilder(TableName.valueOf("concurrentSeed")) + .setStartKey(Bytes.toBytes(i)).setEndKey(Bytes.toBytes(i + 1)).build(); + final ServerMetrics staleReport = ServerMetricsBuilder.newBuilder(sn) + .setRegionMetrics(Collections.singletonList(RegionMetricsBuilder + .newBuilder(ri.getRegionName()).setCompletedSequenceId(staleSeqId).build())) + .build(); + final CyclicBarrier barrier = new CyclicBarrier(2); + Future seedTask = pool.submit(() -> { + await(barrier); + sm.reportRegionOpen(ri, seedSeqId); + }); + Future heartbeatTask = pool.submit(() -> { + await(barrier); + try { + sm.regionServerReport(sn, staleReport); + } catch (Exception e) { + throw new RuntimeException(e); + } + }); + seedTask.get(30, TimeUnit.SECONDS); + heartbeatTask.get(30, TimeUnit.SECONDS); + assertTrue(lastFlushed(ri) >= seedSeqId, "round " + i + ": watermark regressed to " + + lastFlushed(ri) + ", stale heartbeat clobbered the openSeqNum seed"); + } + } finally { + pool.shutdownNow(); + } + } + + private static void await(CyclicBarrier barrier) { + try { + barrier.await(30, TimeUnit.SECONDS); + } catch (Exception e) { + throw new RuntimeException(e); + } + } +} diff --git a/hbase-server/src/test/java/org/apache/hadoop/hbase/master/assignment/TestMergeTableRegionsProcedure.java b/hbase-server/src/test/java/org/apache/hadoop/hbase/master/assignment/TestMergeTableRegionsProcedure.java index 44bd5208977b..b7ce024881aa 100644 --- a/hbase-server/src/test/java/org/apache/hadoop/hbase/master/assignment/TestMergeTableRegionsProcedure.java +++ b/hbase-server/src/test/java/org/apache/hadoop/hbase/master/assignment/TestMergeTableRegionsProcedure.java @@ -355,7 +355,10 @@ public void testRollbackReopensParentsAfterCheckClosedRegionsFailure() throws Ex Path recoveredEditsDir = WALSplitUtil.getRegionDirRecoveredEditsDir(regionDir); FileSystem fs = CommonFSUtils.getRootDirFileSystem(conf); fs.mkdirs(recoveredEditsDir); - Path staleFile = new Path(recoveredEditsDir, "0000000000000000001"); + // Use a seqid above any watermark the master can hold: since HBASE-30335 the master seeds + // flushedSequenceIdByRegion with openSeqNum, so a low-seqid file would be dropped as stale by + // HBASE-30352 and CHECK_CLOSED_REGIONS would no longer fail. + Path staleFile = new Path(recoveredEditsDir, String.valueOf(Long.MAX_VALUE)); fs.createNewFile(staleFile); assertTrue(WALSplitUtil.hasRecoveredEdits(conf, regionsToMerge[0]), "stale recovered.edits file must be visible");