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 3d8e6e003b85..c0debe259959 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 @@ -251,35 +251,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; + }); } } } @@ -1061,6 +1076,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 482e77321c72..fb7596930c79 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 @@ -2244,6 +2244,10 @@ void regionOpenedWithoutPersistingToMeta(RegionStateNode regionNode) throws IOEx 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 f05aa23613aa..1173851b0266 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 @@ -46,6 +46,7 @@ import org.apache.hadoop.hbase.client.Table; import org.apache.hadoop.hbase.ipc.ServerNotRunningYetException; 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; @@ -351,6 +352,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 b72876f44f43..2bd2bec96e34 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() .getLastSequenceId(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() @@ -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 { + MiniHBaseCluster 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..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); + } + } +}