Skip to content
Open
Original file line number Diff line number Diff line change
Expand Up @@ -289,35 +289,50 @@ ServerName regionServerStartup(RegionServerStartupRequest request, int versionNu
private void updateLastFlushedSequenceIds(ServerName sn, ServerMetrics hsl) {
for (Entry<byte[], RegionMetrics> 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<byte[], Long> storeFlushedSequenceId =
computeIfAbsent(storeFlushedSequenceIdsByRegion, encodedRegionName,
() -> new ConcurrentSkipListMap<>(Bytes.BYTES_COMPARATOR));
for (Entry<byte[], Long> 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;
});
}
}
}
Expand Down Expand Up @@ -1092,6 +1107,24 @@ public void removeRegion(final RegionInfo regionInfo) {
flushedSequenceIdByRegion.remove(encodedName);
}

/**
* Called on region OPEN to seed {@link #flushedSequenceIdByRegion} with the region's
* {@code openSeqNum}. Without this, the entry stays absent until the hosting server's next
* heartbeat, so {@link #getLastFlushedSequenceId} returns {@link HConstants#NO_SEQNUM} and
* WALSplitter conservatively treats already-durable edits as unflushed - producing orphaned
* recovered.edits when the source server crashes soon after a drain-move. Uses {@code merge} with
* {@link Math#max} so a heartbeat-supplied value (which may reflect flushes after open) is never
* regressed - and, unlike {@code putIfAbsent}, a stale-low prior value is lifted to
* {@code openSeqNum}. Safe because at OPEN a region cannot have flushed past its own
* {@code openSeqNum}. See HBASE-30335.
*/
public void reportRegionOpen(final RegionInfo regionInfo, final long openSeqNum) {
if (openSeqNum < 0) { // NO_SEQNUM == -1
return;
}
flushedSequenceIdByRegion.merge(regionInfo.getEncodedNameAsBytes(), openSeqNum, Math::max);
}

public boolean isRegionInServerManagerStates(final RegionInfo hri) {
final byte[] encodedName = hri.getEncodedNameAsBytes();
return (storeFlushedSequenceIdsByRegion.containsKey(encodedName)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -2287,6 +2296,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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -415,6 +416,18 @@ public void makeWAL(HRegionServer hrs, List<RegionInfo> 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();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -30,6 +31,7 @@
import org.apache.hadoop.hbase.TableName;
import org.apache.hadoop.hbase.client.Put;
import org.apache.hadoop.hbase.client.Table;
import org.apache.hadoop.hbase.regionserver.HRegion;
import org.apache.hadoop.hbase.regionserver.HRegionServer;
import org.apache.hadoop.hbase.regionserver.Region;
import org.apache.hadoop.hbase.testclassification.MediumTests;
Expand Down Expand Up @@ -89,10 +91,12 @@ public void test() throws IOException, InterruptedException {
Thread.sleep(2000);
RegionStoreSequenceIds ids = testUtil.getHBaseCluster().getMaster().getServerManager()
.getLastFlushedSequenceId(region.getRegionInfo().getEncodedNameAsBytes());
assertEquals(HConstants.NO_SEQNUM, ids.getLastFlushedSequenceId());
// This will be the sequenceid just before that of the earliest edit in memstore.
long storeSequenceId = ids.getStoreSequenceId(0).getSequenceId();
assertTrue(storeSequenceId > 0);
// HBASE-30335: openSeqNum is now seeded on region OPEN, so lastFlushedSequenceId is no
// longer NO_SEQNUM before the first flush.
assertNotEquals(HConstants.NO_SEQNUM, ids.getLastFlushedSequenceId());
testUtil.getAdmin().flush(tableName);
Thread.sleep(2000);
ids = testUtil.getHBaseCluster().getMaster().getServerManager()
Expand All @@ -102,4 +106,97 @@ public void test() throws IOException, InterruptedException {
assertEquals(ids.getLastFlushedSequenceId(), ids.getStoreSequenceId(0).getSequenceId());
table.close();
}

/**
* HBASE-30335: after a region is opened - and before any user write or flush - the master's
* flushedSequenceIdByRegion must already contain the region's openSeqNum. Otherwise a subsequent
* WAL split (e.g. the hosting RS crashes before its first flush heartbeat) would treat
* already-durable edits as unflushed and produce orphaned recovered.edits.
*/
@Test
public void testFlushedSequenceIdSeededOnRegionOpen() throws IOException, InterruptedException {
TableName freshTable = TableName.valueOf(getClass().getSimpleName(), "openseed");
testUtil.getAdmin()
.createNamespace(NamespaceDescriptor.create(freshTable.getNamespaceAsString()).build());
Table table = testUtil.createTable(freshTable, families);
try {
SingleProcessHBaseCluster cluster = testUtil.getMiniHBaseCluster();
HRegion region = null;
for (JVMClusterUtil.RegionServerThread rst : cluster.getRegionServerThreads()) {
for (HRegion r : rst.getRegionServer().getRegions(freshTable)) {
region = r;
break;
}
if (region != null) {
break;
}
}
assertNotNull(region);
long openSeqNum = region.getOpenSeqNum();
RegionStoreSequenceIds ids = testUtil.getHBaseCluster().getMaster().getServerManager()
.getLastFlushedSequenceId(region.getRegionInfo().getEncodedNameAsBytes());
assertNotEquals(HConstants.NO_SEQNUM, ids.getLastFlushedSequenceId(),
"flushedSequenceIdByRegion should be seeded on region OPEN (HBASE-30335)");
assertEquals(openSeqNum, ids.getLastFlushedSequenceId(),
"seeded value must equal the region's openSeqNum");
} finally {
table.close();
}
}

/**
* HBASE-30335 follow-up: on graceful region CLOSE the RS reports the region's final durable
* flushed seqid, and the master must lift its flushedSequenceIdByRegion watermark to that value.
* This complements the OPEN-time seed and covers the drain-move / disable window: if the source
* RS crashes after the close, a subsequent WAL split must see the already-durable edits as
* flushed and not resurrect them as orphaned recovered.edits. Writes and flushes so the region's
* maxFlushedSeqId advances well past its openSeqNum, then disables the table (a graceful close
* that keeps the region - unlike delete/split/merge, disable does not call
* ServerManager#removeRegion) and asserts the watermark reflects the post-close flushed seqid.
*/
@Test
public void testFlushedSequenceIdSeededOnRegionClose() throws IOException, InterruptedException {
TableName freshTable = TableName.valueOf(getClass().getSimpleName(), "closeseed");
testUtil.getAdmin()
.createNamespace(NamespaceDescriptor.create(freshTable.getNamespaceAsString()).build());
Table table = testUtil.createTable(freshTable, families);
try {
SingleProcessHBaseCluster cluster = testUtil.getMiniHBaseCluster();
HRegion region = null;
for (JVMClusterUtil.RegionServerThread rst : cluster.getRegionServerThreads()) {
for (HRegion r : rst.getRegionServer().getRegions(freshTable)) {
region = r;
break;
}
if (region != null) {
break;
}
}
assertNotNull(region);
long openSeqNum = region.getOpenSeqNum();
// Write and flush a few times so the region's durable flushed seqid advances past openSeqNum;
// this makes the CLOSE-time seed distinguishable from the OPEN-time seed.
for (int i = 0; i < 3; i++) {
table.put(new Put(Bytes.toBytes("k" + i)).addColumn(family, Bytes.toBytes("q"),
Bytes.toBytes("v" + i)));
testUtil.getAdmin().flush(freshTable);
}
long flushedSeqId = region.getMaxFlushedSeqId();
assertTrue(flushedSeqId > openSeqNum,
"test setup: flushed seqid " + flushedSeqId + " must exceed openSeqNum " + openSeqNum);
byte[] encodedName = region.getRegionInfo().getEncodedNameAsBytes();
// Gracefully close the region (disable keeps the region entry - it is not removed like a
// delete/split/merge would), driving a CLOSED transition that carries the flushed seqid.
testUtil.getAdmin().disableTable(freshTable);
RegionStoreSequenceIds ids = testUtil.getHBaseCluster().getMaster().getServerManager()
.getLastFlushedSequenceId(encodedName);
assertNotEquals(HConstants.NO_SEQNUM, ids.getLastFlushedSequenceId(),
"flushedSequenceIdByRegion should be seeded on region CLOSE");
assertTrue(ids.getLastFlushedSequenceId() >= flushedSeqId,
"CLOSE seed " + ids.getLastFlushedSequenceId() + " must be >= the region's flushed seqid "
+ flushedSeqId);
} finally {
table.close();
}
}
}
Loading