From d81068f7477b1a332a5fffa277bbbb547694faa1 Mon Sep 17 00:00:00 2001 From: "nirdosh.yadav" <> Date: Fri, 28 Aug 2026 15:11:03 +0530 Subject: [PATCH 01/10] 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 a9ae6f72a29b..0fe327fb83da 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 @@ -2292,6 +2292,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 6619d4096342becd655b2eba49cd33cd32868994 Mon Sep 17 00:00:00 2001 From: "nirdosh.yadav" <> Date: Mon, 31 Aug 2026 09:38:22 +0530 Subject: [PATCH 02/10] 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 880df193f7c005f59b7b4cb8bdb225fbc6db041a Mon Sep 17 00:00:00 2001 From: "nirdosh.yadav" <> Date: Mon, 31 Aug 2026 10:18:52 +0530 Subject: [PATCH 03/10] 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 f527ca2714f5cac08f01eafe142b6f158f73295e Mon Sep 17 00:00:00 2001 From: "nirdosh.yadav" <> Date: Tue, 1 Sep 2026 06:11:37 +0530 Subject: [PATCH 04/10] 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 5c586315e95d7e28a2cdc2a5648ad22d16e79413 Mon Sep 17 00:00:00 2001 From: "nirdosh.yadav" <> Date: Tue, 1 Sep 2026 06:14:30 +0530 Subject: [PATCH 05/10] 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 a4087c44fceed2af7cb0c7d8f0fd3309b3e867be Mon Sep 17 00:00:00 2001 From: "nirdosh.yadav" <> Date: Tue, 1 Sep 2026 06:45:53 +0530 Subject: [PATCH 06/10] 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 8a60a66f975f4bd92c32643189bf325dbaeed923 Mon Sep 17 00:00:00 2001 From: "nirdosh.yadav" <> Date: Wed, 2 Sep 2026 12:11:04 +0530 Subject: [PATCH 07/10] 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 9ab296e6a37441ab4273699f5712ed8fac1844d6 Mon Sep 17 00:00:00 2001 From: "nirdosh.yadav" <> Date: Thu, 17 Sep 2026 03:39:07 +0530 Subject: [PATCH 08/10] 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); + } + } } From a88467141b061939d87467574561cfecf4d1bef3 Mon Sep 17 00:00:00 2001 From: "nirdosh.yadav" <> Date: Tue, 6 Oct 2026 15:26:42 +0530 Subject: [PATCH 09/10] HBASE-30335 Use above-watermark seqid for planted recovered.edits in TestMergeTableRegionsProcedure With the OPEN-time seed, the master's flushed watermark equals openSeqNum, so HBASE-30352 dropped the planted seqid-1 recovered.edits file as stale and the merge succeeded instead of exercising the CHECK_CLOSED_REGIONS rollback path. Name the file Long.MAX_VALUE so it always needs replay. --- .../master/assignment/TestMergeTableRegionsProcedure.java | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) 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"); From 7b5b6a7c0ee8f693c7768882fc155ce62edce65c Mon Sep 17 00:00:00 2001 From: "nirdosh.yadav" <> Date: Mon, 28 Sep 2026 14:02:27 +0530 Subject: [PATCH 10/10] HBASE-30433 Seed the master's flushedSequenceIdByRegion watermark with the region's final flushed seqid on region CLOSE On graceful region CLOSE the RegionServer now reports the region's durable flushed seqid (HRegion.getMaxFlushedSeqId()) on the CLOSED transition, and the master lifts its flushedSequenceIdByRegion watermark from it via the existing monotonic reportRegionOpen(merge(Math::max)) seed. This narrows the graceful-close window where a subsequent WAL split of a crashed source RS could otherwise treat already-durable edits as unflushed and write orphaned recovered.edits. No proto/RPC change: the RegionStateTransition message already carries an optional openSeqNum field and the CLOSED transition is already routed through the same master-side handler. HRegionServer now also sets that field for CLOSED (not just OPENED), and AssignmentManager seeds the watermark on CLOSED when seqId >= 0. Adds TestGetLastFlushedSequenceId#testFlushedSequenceIdSeededOnRegionClose. --- .../master/assignment/AssignmentManager.java | 13 ++++- .../hbase/regionserver/HRegionServer.java | 5 +- .../handler/CloseRegionHandler.java | 7 ++- .../handler/UnassignRegionHandler.java | 10 +++- .../master/TestGetLastFlushedSequenceId.java | 56 +++++++++++++++++++ 5 files changed, 83 insertions(+), 8 deletions(-) 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 0fe327fb83da..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: 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/TestGetLastFlushedSequenceId.java b/hbase-server/src/test/java/org/apache/hadoop/hbase/master/TestGetLastFlushedSequenceId.java index f2b76a07065e..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 @@ -143,4 +143,60 @@ public void testFlushedSequenceIdSeededOnRegionOpen() throws IOException, Interr 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(); + } + } }