From ab5a3e96edd5f8e701b5e743d3a192644e71a888 Mon Sep 17 00:00:00 2001 From: "nirdosh.yadav" <> Date: Fri, 28 Aug 2026 15:11:03 +0530 Subject: [PATCH 1/6] 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/6] 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/6] 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/6] 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/6] 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/6] 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();