From ab5a3e96edd5f8e701b5e743d3a192644e71a888 Mon Sep 17 00:00:00 2001 From: "nirdosh.yadav" <> Date: Fri, 28 Aug 2026 15:11:03 +0530 Subject: [PATCH 1/8] HBASE-30335 Seed master flushedSequenceIdByRegion with openSeqNum on region OPEN When a region is opened, the master does not populate flushedSequenceIdByRegion until the hosting RegionServer's next heartbeat delivers a flush report. If the source RegionServer of a drain-move crashes before that heartbeat, ServerManager. getLastFlushedSequenceId returns NO_SEQNUM (-1) for the region, and WALSplitter conservatively writes already-durable edits into recovered.edits. Those orphaned edits then trigger false-positive "data loss" warnings during subsequent merge/split operations and leave regions stuck in RIT. Add ServerManager.reportRegionOpen(regionInfo, openSeqNum) and call it from AssignmentManager.regionOpenedWithoutPersistingToMeta so the watermark is established synchronously at OPEN time. putIfAbsent is used so a subsequent heartbeat with a higher completedSequenceId is never regressed by a stale open value. Co-Authored-By: Claude Opus 4.7 (1M context) --- .../hadoop/hbase/master/ServerManager.java | 16 ++++++++++++++++ .../master/assignment/AssignmentManager.java | 4 ++++ 2 files changed, 20 insertions(+) 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..93a8349c3095 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,22 @@ 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 putIfAbsent} so a heartbeat-supplied value (which may reflect flushes after open) is + * never regressed. See HBASE-30335. + */ + public void reportRegionOpen(final RegionInfo regionInfo, final long openSeqNum) { + if (openSeqNum == HConstants.NO_SEQNUM || openSeqNum < 0) { + return; + } + flushedSequenceIdByRegion.putIfAbsent(regionInfo.getEncodedNameAsBytes(), openSeqNum); + } + 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 From 731f6c9708d06e7f2e8aa187092f7af1924994b9 Mon Sep 17 00:00:00 2001 From: "nirdosh.yadav" <> Date: Mon, 31 Aug 2026 09:38:22 +0530 Subject: [PATCH 2/8] HBASE-30335 Add tests for ServerManager.reportRegionOpen and update TestGetLastFlushedSequenceId New unit test TestServerManager covers reportRegionOpen behavior: - seeds flushedSequenceIdByRegion with the supplied openSeqNum; - putIfAbsent semantics prevent regressing a higher watermark that was already established (by an earlier open or a heartbeat); - NO_SEQNUM and negative openSeqNum are ignored (no-op). TestGetLastFlushedSequenceId previously asserted the pre-flush lastFlushedSequenceId was NO_SEQNUM. That assumption is invalidated by the fix (openSeqNum is now seeded synchronously on OPEN); the assertion is updated to require the watermark be present but strictly less than the memstore's earliest unflushed edit. Co-Authored-By: Claude Opus 4.7 (1M context) --- .../master/TestGetLastFlushedSequenceId.java | 7 +- .../hbase/master/TestServerManager.java | 99 +++++++++++++++++++ 2 files changed, 105 insertions(+), 1 deletion(-) create mode 100644 hbase-server/src/test/java/org/apache/hadoop/hbase/master/TestServerManager.java 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..ceeddf297b15 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; @@ -89,10 +90,14 @@ 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 - it is the region's openSeqNum, which must + // still be strictly less than the memstore's earliest unflushed edit. + assertNotEquals(HConstants.NO_SEQNUM, ids.getLastFlushedSequenceId()); + assertTrue(ids.getLastFlushedSequenceId() < storeSequenceId); testUtil.getAdmin().flush(tableName); Thread.sleep(2000); ids = testUtil.getHBaseCluster().getMaster().getServerManager() 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)); + } +} From cf1120c0461650c006a23ea3630328ef54a3d6fc Mon Sep 17 00:00:00 2001 From: "nirdosh.yadav" <> Date: Mon, 31 Aug 2026 10:18:52 +0530 Subject: [PATCH 3/8] Added unit test --- .../master/TestGetLastFlushedSequenceId.java | 38 +++++++++++++++++++ 1 file changed, 38 insertions(+) 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 ceeddf297b15..4400f03e4804 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 @@ -31,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; @@ -107,4 +108,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(); + } + } } From 20432cb539a36162f069efd1174db0fbefbbf3dd Mon Sep 17 00:00:00 2001 From: "nirdosh.yadav" <> Date: Tue, 1 Sep 2026 06:11:37 +0530 Subject: [PATCH 4/8] HBASE-30335 Address review: merge/Math::max seed and drop fragile test assertion Per review from @apurtell: - ServerManager.reportRegionOpen: switch putIfAbsent to merge with Math::max so a stale-low prior heartbeat value is lifted to openSeqNum rather than ignored. Safe because at OPEN a region cannot have flushed past its own openSeqNum. Guard simplified to openSeqNum < 0 (NO_SEQNUM == -1, so the disjunct was redundant). - TestGetLastFlushedSequenceId: drop the strict assertTrue(lastFlushed < storeSequenceId) - it holds only because the region-open marker consumes one seqId, so the assertion is coupled to an incidental accounting detail rather than the contract being tested. The assertNotEquals(NO_SEQNUM, ...) above captures the load-bearing invariant. Reflow adjacent javadoc block for spotless. --- .../org/apache/hadoop/hbase/master/ServerManager.java | 10 ++++++---- .../hbase/master/TestGetLastFlushedSequenceId.java | 10 ++++------ 2 files changed, 10 insertions(+), 10 deletions(-) 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 93a8349c3095..2bf2c4941dfd 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 @@ -1098,14 +1098,16 @@ public void removeRegion(final RegionInfo regionInfo) { * 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 putIfAbsent} so a heartbeat-supplied value (which may reflect flushes after open) is - * never regressed. See HBASE-30335. + * {@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 == HConstants.NO_SEQNUM || openSeqNum < 0) { + if (openSeqNum < 0) { // NO_SEQNUM == -1 return; } - flushedSequenceIdByRegion.putIfAbsent(regionInfo.getEncodedNameAsBytes(), openSeqNum); + flushedSequenceIdByRegion.merge(regionInfo.getEncodedNameAsBytes(), openSeqNum, Math::max); } public boolean isRegionInServerManagerStates(final RegionInfo hri) { 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 4400f03e4804..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 @@ -95,10 +95,8 @@ public void test() throws IOException, InterruptedException { 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 - it is the region's openSeqNum, which must - // still be strictly less than the memstore's earliest unflushed edit. + // longer NO_SEQNUM before the first flush. assertNotEquals(HConstants.NO_SEQNUM, ids.getLastFlushedSequenceId()); - assertTrue(ids.getLastFlushedSequenceId() < storeSequenceId); testUtil.getAdmin().flush(tableName); Thread.sleep(2000); ids = testUtil.getHBaseCluster().getMaster().getServerManager() @@ -111,9 +109,9 @@ public void test() throws IOException, InterruptedException { /** * 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. + * 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 { From 095ed06c025451b34e62f8d84f5c86ea1eb62235 Mon Sep 17 00:00:00 2001 From: "nirdosh.yadav" <> Date: Tue, 1 Sep 2026 06:14:30 +0530 Subject: [PATCH 5/8] HBASE-30335 spotless reflow of reportRegionOpen javadoc --- .../org/apache/hadoop/hbase/master/ServerManager.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) 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 2bf2c4941dfd..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 @@ -1097,10 +1097,10 @@ public void removeRegion(final RegionInfo regionInfo) { * {@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 + * 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) { From a59f82857c9fe7a97e5af208b7ef5a32fff6f69c Mon Sep 17 00:00:00 2001 From: "nirdosh.yadav" <> Date: Tue, 1 Sep 2026 06:45:53 +0530 Subject: [PATCH 6/8] HBASE-30335 Fix TestDLS makeWAL to respect openSeqNum invariant The seed added in reportRegionOpen (openSeqNum via Math::max) exposed a long-standing test-fixture issue: AbstractTestDLS.makeWAL uses a fresh MultiVersionConcurrencyControl that stamps WAL edits starting at seqid 1, inconsistent with the seqid sequence a real WAL preserves for the region. With the seed in place, the splitter correctly filters those low-seqid edits as already-durable and testMasterStartsUpWithLogSplittingWork loses 5/1000 rows. Advance the local MVCC past the max openSeqNum of the target regions before stamping edits, so the injected WAL entries get seqids a real region would have assigned. --- .../apache/hadoop/hbase/master/AbstractTestDLS.java | 13 +++++++++++++ 1 file changed, 13 insertions(+) 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(); From 7761ebfd59d3e7e7c684b4cb5c94d2adde8257b6 Mon Sep 17 00:00:00 2001 From: "nirdosh.yadav" <> Date: Wed, 2 Sep 2026 12:11:04 +0530 Subject: [PATCH 7/8] HBASE-30335 Fix TestMaster.testFlushedSequenceIdPersistLoad for openSeqNum seed The test previously asserted flushedSequenceIdByRegion is byte-for-byte identical across cluster shutdown+restart. After HBASE-30335 the master seeds this map on region OPEN via merge(openSeqNum, Math::max). openSeqNum is monotonic across close/open cycles, so a region reopened after restart can carry a strictly higher value than what was persisted at shutdown, and the equality assertion no longer holds. Assert the preserved invariant instead: every region persisted at shutdown is loaded on restart (same keyset) and no watermark regresses (after[r] >= before[r]). This validates persist/load correctness without conflicting with the new seed-on-open semantic. --- .../apache/hadoop/hbase/master/TestMaster.java | 18 ++++++++++++++++-- 1 file changed, 16 insertions(+), 2 deletions(-) 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 From 3827001c668be843e5bbf5f3ef9283b38f102c07 Mon Sep 17 00:00:00 2001 From: "nirdosh.yadav" <> Date: Thu, 17 Sep 2026 03:39:07 +0530 Subject: [PATCH 8/8] HBASE-30335 Make heartbeat flushed-seqid update atomic so it cannot clobber the OPEN seed updateLastFlushedSequenceIds did a non-atomic get-then-put on flushedSequenceIdByRegion. A stale in-flight heartbeat from the soon-to-be-dead source RS could read null/a low value, race past the reportRegionOpen seed (merge/Math::max), and then put its lower value on top - reintroducing the stale-fence bug this PR fixes, in a narrower window (flagged by Copilot's review). Replace both the region- and store-level updates with an atomic compute that keeps the same "never lower the watermark" rule. Every writer on the live serving path is now an atomic max-merge, so the stored value is monotonic non-decreasing and a concurrent seed can no longer be clobbered. Add TestServerManager.testConcurrentStaleHeartbeatDoesNotClobberOpenSeed, which drives the OPEN seed and a stale heartbeat concurrently over 500 rounds. It fails on the pre-fix code and passes with the fix. Co-Authored-By: Claude Opus 4.8 (1M context) --- .../hadoop/hbase/master/ServerManager.java | 63 ++++++++++------ .../hbase/master/TestServerManager.java | 73 +++++++++++++++++++ 2 files changed, 112 insertions(+), 24 deletions(-) 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 344ffa396b96..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; + }); } } } 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 index 3cbeaf265f9c..9917dea9c67f 100644 --- 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 @@ -18,12 +18,23 @@ 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; @@ -31,6 +42,7 @@ 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; @@ -96,4 +108,65 @@ 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); + } + } }