From 335fe85d2e9b9529f85883029802015e74e4075b Mon Sep 17 00:00:00 2001 From: Hernan Gelaf-Romer Date: Wed, 23 Sep 2026 12:45:57 -0400 Subject: [PATCH] Log filtering in IncrementalBackupManager can lead to data loss --- .../hbase/backup/impl/BackupSystemTable.java | 5 +- .../backup/impl/FullTableBackupClient.java | 90 +++- .../backup/impl/IncrementalBackupManager.java | 45 +- .../hadoop/hbase/backup/util/BackupUtils.java | 80 ++++ .../hbase/backup/TestBackupOfflineRS.java | 439 ++++++++++++++++++ .../hadoop/hbase/backup/TestBackupUtils.java | 80 ++++ .../backup/TestIncrementalBackupManager.java | 280 +++++++++++ ...estFullTableBackupClientLogBoundaries.java | 148 ++++++ 8 files changed, 1147 insertions(+), 20 deletions(-) create mode 100644 hbase-backup/src/test/java/org/apache/hadoop/hbase/backup/TestBackupOfflineRS.java create mode 100644 hbase-backup/src/test/java/org/apache/hadoop/hbase/backup/impl/TestFullTableBackupClientLogBoundaries.java diff --git a/hbase-backup/src/main/java/org/apache/hadoop/hbase/backup/impl/BackupSystemTable.java b/hbase-backup/src/main/java/org/apache/hadoop/hbase/backup/impl/BackupSystemTable.java index 73ed588b265c..b6d0d95c5942 100644 --- a/hbase-backup/src/main/java/org/apache/hadoop/hbase/backup/impl/BackupSystemTable.java +++ b/hbase-backup/src/main/java/org/apache/hadoop/hbase/backup/impl/BackupSystemTable.java @@ -662,8 +662,9 @@ public List getBackupHistory(Order order, int n, BackupInfo.Filter.. /** * Write the current timestamps for each regionserver to backup system table after a successful - * full or incremental backup. The saved timestamp is of the last log file that was backed up - * already. + * full or incremental backup. For a region server that took part in the backup's log roll, the + * saved timestamp is the result of that roll. For a region server that did not (offline or + * decommissioned), it is the creation time of the newest of its log files that was backed up. * @param tables tables * @param newTimestamps timestamps * @param backupRoot root directory path to backup diff --git a/hbase-backup/src/main/java/org/apache/hadoop/hbase/backup/impl/FullTableBackupClient.java b/hbase-backup/src/main/java/org/apache/hadoop/hbase/backup/impl/FullTableBackupClient.java index dae0be0be378..14fba9f994e1 100644 --- a/hbase-backup/src/main/java/org/apache/hadoop/hbase/backup/impl/FullTableBackupClient.java +++ b/hbase-backup/src/main/java/org/apache/hadoop/hbase/backup/impl/FullTableBackupClient.java @@ -36,13 +36,21 @@ import static org.apache.hadoop.hbase.replication.ReplicationUtils.OFFSET_UPDATE_INTERVAL_MS_KEY; import static org.apache.hadoop.hbase.replication.ReplicationUtils.OFFSET_UPDATE_SIZE_THRESHOLD_KEY; +import com.google.errorprone.annotations.RestrictedApi; import java.io.IOException; import java.util.ArrayList; import java.util.HashMap; +import java.util.HashSet; import java.util.List; import java.util.Map; +import java.util.Set; import java.util.UUID; import java.util.stream.Collectors; +import org.apache.hadoop.fs.FileStatus; +import org.apache.hadoop.fs.FileSystem; +import org.apache.hadoop.fs.Path; +import org.apache.hadoop.hbase.HConstants; +import org.apache.hadoop.hbase.ServerName; import org.apache.hadoop.hbase.TableName; import org.apache.hadoop.hbase.backup.BackupCopyJob; import org.apache.hadoop.hbase.backup.BackupInfo; @@ -60,7 +68,9 @@ import org.apache.hadoop.hbase.client.TableDescriptorBuilder; import org.apache.hadoop.hbase.replication.ReplicationException; import org.apache.hadoop.hbase.replication.ReplicationPeerConfig; +import org.apache.hadoop.hbase.util.CommonFSUtils; import org.apache.hadoop.hbase.util.EnvironmentEdgeManager; +import org.apache.hadoop.hbase.wal.AbstractFSWALProvider; import org.apache.yetus.audience.InterfaceAudience; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -208,7 +218,12 @@ private void handleContinuousBackup(Admin admin) throws IOException { } private void handleNonContinuousBackup(Admin admin) throws IOException { - performLogRoll(); + Map previousLogRollsByHost = backupManager.readRegionServerLastLogRollResult(); + Map latestLogRollsByHost = performLogRoll(); + Path walRootDir = CommonFSUtils.getWALRootDir(conf); + newTimestamps = computeLogBoundaries(walRootDir.getFileSystem(conf), walRootDir, + BackupUtils.getRolledHosts(previousLogRollsByHost, latestLogRollsByHost), admin); + performBackupSnapshots(admin); backupManager.addIncrementalBackupTableSet(backupInfo.getTables()); @@ -219,7 +234,7 @@ private void handleNonContinuousBackup(Admin admin) throws IOException { updateBackupMetadata(); } - private void performLogRoll() throws IOException { + private Map performLogRoll() throws IOException { // We roll log here before we do the snapshot. It is possible there is duplicate data // in the log that is already in the snapshot. But if we do it after the snapshot, we // could have data loss. @@ -227,7 +242,7 @@ private void performLogRoll() throws IOException { // the snapshot. LOG.info("Execute roll log procedure for full backup ..."); BackupUtils.logRoll(conn, backupInfo.getBackupRootDir(), conf); - newTimestamps = backupManager.readRegionServerLastLogRollResult(); + return backupManager.readRegionServerLastLogRollResult(); } private void performBackupSnapshots(Admin admin) throws IOException { @@ -248,6 +263,75 @@ private void performSnapshots(Admin admin) throws IOException { } } + /** + * Computes the per-host log boundaries stored by this full backup, using + * {@link BackupUtils#computeLogBoundaries}. A host that took part in this backup's log roll gets + * its roll result: WALs created after the roll are included in the next incremental backup, even + * though some of their edits may already be in the snapshot, which is safe because deletes are + * replayed along with the puts. For any other host, its archived WALs and the WALs of dead + * servers are covered by the snapshot, so the next incremental backup does not replay them. The + * WALs of a live server that did not take part in the roll (for example one that started during + * it) are pending, because they can still receive edits that are not in the snapshot. The live + * servers are read only after the WAL directories are listed: a region server creates its WAL + * directory only after the master registers it, so every live server whose directory was listed + * is found. If the live servers do not include every host that took part in the roll, the + * master's server list is incomplete (for example right after a master failover), and the backup + * fails instead of treating live servers as dead. It must run right after the log roll, before + * the snapshot is taken, so that a host starting later gets no boundary and has all of its WALs + * included in the next incremental backup. + */ + @RestrictedApi( + explanation = "Package-private for test visibility only. Do not use outside tests.", + link = "", + allowedOnPath = "(.*/src/test/.*|.*/org/apache/hadoop/hbase/backup/impl/FullTableBackupClient.java)") + static Map computeLogBoundaries(FileSystem fs, Path walRootDir, + Map rolledHosts, Admin admin) throws IOException { + Path logDir = new Path(walRootDir, HConstants.HREGION_LOGDIR_NAME); + Path oldLogDir = new Path(walRootDir, HConstants.HREGION_OLDLOGDIR_NAME); + + Map> logsByServer = new HashMap<>(); + for (FileStatus serverLogDir : fs.listStatus(logDir)) { + ServerName serverName = + AbstractFSWALProvider.getServerNameFromWALDirectoryName(serverLogDir.getPath()); + if (serverName == null) { + continue; + } + List logs = logsByServer.computeIfAbsent(serverName, k -> new ArrayList<>()); + for (FileStatus log : fs.listStatus(serverLogDir.getPath())) { + if (!AbstractFSWALProvider.isMetaFile(log.getPath())) { + logs.add(log.getPath().toString()); + } + } + } + + Set live = new HashSet<>(admin.getRegionServers()); + Set liveAddresses = + live.stream().map(sn -> sn.getAddress().toString()).collect(Collectors.toSet()); + if (!liveAddresses.containsAll(rolledHosts.keySet())) { + throw new IOException("Live region servers " + live + " do not include every host that took" + + " part in this backup's log roll " + rolledHosts.keySet() + ". The master's server list" + + " may be incomplete, for example during a master failover, so the backup is failed to be" + + " retried."); + } + + List coveredLogs = new ArrayList<>(); + List pendingLogs = new ArrayList<>(); + for (Map.Entry> entry : logsByServer.entrySet()) { + if (live.contains(entry.getKey())) { + pendingLogs.addAll(entry.getValue()); + } else { + coveredLogs.addAll(entry.getValue()); + } + } + + if (fs.exists(oldLogDir)) { + BackupUtils.getFiles(fs, oldLogDir, coveredLogs, + path -> !AbstractFSWALProvider.isMetaFile(path)); + } + + return BackupUtils.computeLogBoundaries(rolledHosts, coveredLogs, pendingLogs); + } + private void updateBackupMetadata() throws IOException { // The table list in backupInfo is good for both full backup and incremental backup. // For incremental backup, it contains the incremental backup table set. diff --git a/hbase-backup/src/main/java/org/apache/hadoop/hbase/backup/impl/IncrementalBackupManager.java b/hbase-backup/src/main/java/org/apache/hadoop/hbase/backup/impl/IncrementalBackupManager.java index 18be4c4f94ab..d50666667f69 100644 --- a/hbase-backup/src/main/java/org/apache/hadoop/hbase/backup/impl/IncrementalBackupManager.java +++ b/hbase-backup/src/main/java/org/apache/hadoop/hbase/backup/impl/IncrementalBackupManager.java @@ -20,8 +20,10 @@ import java.io.IOException; import java.util.ArrayList; import java.util.Collections; +import java.util.HashSet; import java.util.List; import java.util.Map; +import java.util.Set; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.FileStatus; import org.apache.hadoop.fs.FileSystem; @@ -54,12 +56,14 @@ public IncrementalBackupManager(Connection conn, Configuration conf) throws IOEx /** * Obtain the list of logs that need to be copied out for this incremental backup. The list is set * in BackupInfo. - * @return The new HashMap of RS log time stamps after the log roll for this incremental backup. + * @return The new map of RS log time stamps for this incremental backup, as computed by + * {@link BackupUtils#computeLogBoundaries}: the included logs are covered, the logs held + * back for a later backup (for example, logs still being split) are pending, and a host + * that still has logs keeps its previous boundary if this backup gives it no new one. * @throws IOException exception */ public Map getIncrBackupLogFileMap() throws IOException { List logList; - Map newTimestamps; Map previousTimestampMins = BackupUtils.getRSLogTimestampMins(readLogTimestampMap()); @@ -69,18 +73,28 @@ public Map getIncrBackupLogFileMap() throws IOException { + "In order to create an incremental backup, at least one full backup is needed."); } + Map previousLogRollByHost = readRegionServerLastLogRollResult(); LOG.info("Execute roll log procedure for incremental backup ..."); BackupUtils.logRoll(conn, backupInfo.getBackupRootDir(), conf); - newTimestamps = readRegionServerLastLogRollResult(); + Map rolledHosts = + BackupUtils.getRolledHosts(previousLogRollByHost, readRegionServerLastLogRollResult()); - logList = getLogFilesForNewBackup(previousTimestampMins, newTimestamps, conf); - logList = excludeProcV2WALs(logList); + LogFileSelection selection = getLogFilesForNewBackup(previousTimestampMins, rolledHosts, conf); + logList = excludeProcV2WALs(selection.included()); backupInfo.setIncrBackupFileList(logList); + Map newTimestamps = BackupUtils.computeLogBoundaries(rolledHosts, + previousTimestampMins, selection.hostsWithLogs(), logList, selection.heldBack()); + LOG.debug("Log boundaries for incremental backup {}: {}", backupInfo.getBackupId(), + newTimestamps); return newTimestamps; } + private record LogFileSelection(List included, List heldBack, + Set hostsWithLogs) { + } + private List excludeProcV2WALs(List logList) { List list = new ArrayList<>(); for (int i = 0; i < logList.size(); i++) { @@ -97,15 +111,17 @@ private List excludeProcV2WALs(List logList) { } /** - * For each region server: get all log files newer than the last timestamps but not newer than the - * newest timestamps. + * Gather all log files that either: 1) are newer than the older timestamps, but not newer than + * the newest timestamps, or 2) are archived logs whose host name does not occur in the newest + * timestamps. * @param olderTimestamps the timestamp for each region server of the last backup. * @param newestTimestamps the timestamp for each region server that the backup should lead to. * @param conf the Hadoop and Hbase configuration - * @return a list of log files to be backed up + * @return the log files to be backed up, the log files held back for a later backup, and the + * hosts that have any log files, including ones not backed up * @throws IOException exception */ - private List getLogFilesForNewBackup(Map olderTimestamps, + private LogFileSelection getLogFilesForNewBackup(Map olderTimestamps, Map newestTimestamps, Configuration conf) throws IOException { LOG.debug("In getLogFilesForNewBackup()\n" + "olderTimestamps: " + olderTimestamps + "\n newestTimestamps: " + newestTimestamps); @@ -119,6 +135,7 @@ private List getLogFilesForNewBackup(Map olderTimestamps, List resultLogFiles = new ArrayList<>(); List newestLogs = new ArrayList<>(); + Set hostsWithLogs = new HashSet<>(); /* * The old region servers and timestamps info we kept in backup system table may be out of sync @@ -143,6 +160,7 @@ private List getLogFilesForNewBackup(Map olderTimestamps, if (host == null) { continue; } + hostsWithLogs.add(host); FileStatus[] logs; oldTimeStamp = olderTimestamps.get(host); // It is possible that there is no old timestamp in backup system table for this host if @@ -198,6 +216,7 @@ private List getLogFilesForNewBackup(Map olderTimestamps, if (host == null) { continue; } + hostsWithLogs.add(host); currentLogTS = BackupUtils.getCreationTime(p); oldTimeStamp = olderTimestamps.get(host); /* @@ -217,18 +236,14 @@ private List getLogFilesForNewBackup(Map olderTimestamps, resultLogFiles.add(currentLogFile); } - // It is possible that a host in .oldlogs is an obsolete region server - // so newestTimestamps.get(host) here can be null. - // Even if these logs belong to a obsolete region server, we still need - // to include they to avoid loss of edits for backup. Long newTimestamp = newestTimestamps.get(host); - if (newTimestamp == null || currentLogTS > newTimestamp) { + if (newTimestamp != null && currentLogTS > newTimestamp) { newestLogs.add(currentLogFile); } } // remove newest log per host because they are still in use resultLogFiles.removeAll(newestLogs); - return resultLogFiles; + return new LogFileSelection(resultLogFiles, newestLogs, hostsWithLogs); } static class NewestLogFilter implements PathFilter { diff --git a/hbase-backup/src/main/java/org/apache/hadoop/hbase/backup/util/BackupUtils.java b/hbase-backup/src/main/java/org/apache/hadoop/hbase/backup/util/BackupUtils.java index 4154e5f06daa..4bfccd9e8ca6 100644 --- a/hbase-backup/src/main/java/org/apache/hadoop/hbase/backup/util/BackupUtils.java +++ b/hbase-backup/src/main/java/org/apache/hadoop/hbase/backup/util/BackupUtils.java @@ -31,6 +31,7 @@ import java.text.SimpleDateFormat; import java.time.ZoneOffset; import java.util.ArrayList; +import java.util.Collection; import java.util.Collections; import java.util.Comparator; import java.util.Date; @@ -38,6 +39,7 @@ import java.util.List; import java.util.Map; import java.util.Map.Entry; +import java.util.Set; import java.util.TimeZone; import java.util.TreeSet; import java.util.function.Predicate; @@ -206,6 +208,84 @@ public static String parseHostNameFromLogFile(Path p) { } } + /** + * Returns the log roll result of every region server that took part in the log roll that happened + * between reading {@code previousLogRolls} and {@code latestLogRolls}. A region server that did + * not take part (for example because it is offline) keeps its old roll result, which is never + * removed, so it is recognized by its roll result not changing. + * @param previousLogRolls roll results by host, read before the log roll + * @param latestLogRolls roll results by host, read after the log roll + * @return roll results by host, for the hosts that took part in the log roll + */ + public static Map getRolledHosts(Map previousLogRolls, + Map latestLogRolls) { + return latestLogRolls.entrySet().stream() + .filter(entry -> !entry.getValue().equals(previousLogRolls.get(entry.getKey()))) + .collect(Collectors.toMap(Entry::getKey, Entry::getValue)); + } + + /** + * Computes the per-host log boundaries stored by a backup that has no previous boundaries to + * carry forward, such as a full backup. See + * {@link #computeLogBoundaries(Map, Map, Set, Collection, Collection)}. + */ + public static Map computeLogBoundaries(Map rolledHosts, + Collection coveredLogs, Collection pendingLogs) throws IOException { + return computeLogBoundaries(rolledHosts, Collections.emptyMap(), Collections.emptySet(), + coveredLogs, pendingLogs); + } + + /** + * Computes the per-host log boundaries stored by a backup. Everything up to a host's boundary is + * covered by backups, so later incremental backups only include that host's logs that are newer. + *
    + *
  • A host that took part in the backup's log roll gets its roll result.
  • + *
  • Any other host gets the creation time of its newest covered log, capped to just below its + * oldest pending log, so that pending logs are included in a later backup.
  • + *
  • A host without covered or pending logs keeps its previous boundary while it still has logs, + * and otherwise gets no boundary.
  • + *
+ * Capping is needed because a host's logs can still be pending while some of its newer logs are + * already covered, for example when a dead server's logs are split and archived out of order. A + * host is also kept while any of its logs are pending, even if none are covered, because without + * a boundary the next backup would treat it as unknown and skip its older logs. For the same + * reason a host keeps its previous boundary while it still has logs: without it, the next backup + * would include its logs that are already covered, which can bring back deleted data. + * @param rolledHosts roll results of the hosts that took part in the log roll + * @param previousBoundaries boundaries stored by the previous backup + * @param hostsWithLogs hosts that still have logs, whether or not this backup includes them + * @param coveredLogs logs whose edits are covered by this backup + * @param pendingLogs logs whose edits might not be covered by this backup + * @return boundaries by host + */ + public static Map computeLogBoundaries(Map rolledHosts, + Map previousBoundaries, Set hostsWithLogs, Collection coveredLogs, + Collection pendingLogs) throws IOException { + Map otherHosts = new HashMap<>(); + for (String log : coveredLogs) { + Path path = new Path(log); + String host = parseHostNameFromLogFile(path); + if (host != null && !rolledHosts.containsKey(host)) { + otherHosts.merge(host, getCreationTime(path), Math::max); + } + } + for (String log : pendingLogs) { + Path path = new Path(log); + String host = parseHostNameFromLogFile(path); + if (host != null && !rolledHosts.containsKey(host)) { + otherHosts.merge(host, getCreationTime(path) - 1, Math::min); + } + } + Map boundaries = new HashMap<>(rolledHosts); + boundaries.putAll(otherHosts); + for (Entry previous : previousBoundaries.entrySet()) { + if (hostsWithLogs.contains(previous.getKey())) { + boundaries.putIfAbsent(previous.getKey(), previous.getValue()); + } + } + return boundaries; + } + /** * Returns WAL file name * @param walFileName WAL file name diff --git a/hbase-backup/src/test/java/org/apache/hadoop/hbase/backup/TestBackupOfflineRS.java b/hbase-backup/src/test/java/org/apache/hadoop/hbase/backup/TestBackupOfflineRS.java new file mode 100644 index 000000000000..ecf3647309b3 --- /dev/null +++ b/hbase-backup/src/test/java/org/apache/hadoop/hbase/backup/TestBackupOfflineRS.java @@ -0,0 +1,439 @@ +/* + * 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.backup; + +import static org.junit.jupiter.api.Assertions.assertAll; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.io.IOException; +import java.util.List; +import java.util.Map; +import org.apache.hadoop.hbase.HBaseTestingUtil; +import org.apache.hadoop.hbase.ServerName; +import org.apache.hadoop.hbase.SingleProcessHBaseCluster; +import org.apache.hadoop.hbase.TableName; +import org.apache.hadoop.hbase.backup.impl.BackupAdminImpl; +import org.apache.hadoop.hbase.backup.impl.BackupSystemTable; +import org.apache.hadoop.hbase.backup.impl.FullTableBackupClient; +import org.apache.hadoop.hbase.backup.impl.TableBackupClient; +import org.apache.hadoop.hbase.backup.util.BackupUtils; +import org.apache.hadoop.hbase.client.Admin; +import org.apache.hadoop.hbase.client.Connection; +import org.apache.hadoop.hbase.client.ConnectionFactory; +import org.apache.hadoop.hbase.client.Delete; +import org.apache.hadoop.hbase.client.Get; +import org.apache.hadoop.hbase.client.Put; +import org.apache.hadoop.hbase.client.RegionInfo; +import org.apache.hadoop.hbase.client.ResultScanner; +import org.apache.hadoop.hbase.client.Scan; +import org.apache.hadoop.hbase.client.Table; +import org.apache.hadoop.hbase.master.HMaster; +import org.apache.hadoop.hbase.master.procedure.ServerCrashProcedure; +import org.apache.hadoop.hbase.regionserver.HRegion; +import org.apache.hadoop.hbase.regionserver.HRegionServer; +import org.apache.hadoop.hbase.testclassification.LargeTests; +import org.apache.hadoop.hbase.util.Bytes; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.Tag; +import org.junit.jupiter.api.Test; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import org.apache.hbase.thirdparty.com.google.common.collect.Lists; + +/** + * Tests that WAL files from offline/inactive RegionServers are handled correctly during backup. + * Specifically verifies that WALs from an offline RS are: + *
    + *
  1. Backed up once in the first backup after the RS goes offline
  2. + *
  3. NOT re-backed up in subsequent backups
  4. + *
+ */ +@Tag(LargeTests.TAG) +public class TestBackupOfflineRS extends TestBackupBase { + + private static final Logger LOG = LoggerFactory.getLogger(TestBackupOfflineRS.class); + + @BeforeAll + public static void setUp() throws Exception { + TEST_UTIL = new HBaseTestingUtil(); + conf1 = TEST_UTIL.getConfiguration(); + conf1.setInt("hbase.regionserver.info.port", -1); + conf1.setInt(HMaster.HBASE_MASTER_CLEANER_INTERVAL, Integer.MAX_VALUE); + autoRestoreOnFailure = true; + useSecondCluster = false; + setUpHelper(); + TEST_UTIL.getMiniHBaseCluster().startRegionServer(); + TEST_UTIL.waitTableAvailable(table1); + } + + private static Runnable afterSnapshotHook; + + /** Full backup client that runs {@link #afterSnapshotHook} once the table snapshots exist. */ + public static class FullTableBackupClientWithHook extends FullTableBackupClient { + @Override + protected void snapshotCopy(BackupInfo backupInfo) throws IOException { + afterSnapshotHook.run(); + super.snapshotCopy(backupInfo); + } + } + + private static void stopRegionServerAndWait(ServerName serverName) throws Exception { + SingleProcessHBaseCluster cluster = TEST_UTIL.getMiniHBaseCluster(); + HMaster master = cluster.getMaster(); + cluster.stopRegionServer(serverName); + cluster.waitForRegionServerToStop(serverName, 60_000); + TEST_UTIL.waitFor(60_000, + () -> master.getProcedures().stream().filter(ServerCrashProcedure.class::isInstance) + .map(ServerCrashProcedure.class::cast) + .anyMatch(scp -> scp.getServerName().equals(serverName) && scp.isFinished())); + TEST_UTIL.waitUntilNoRegionsInTransition(60_000); + } + + private static void restore(String backupId, TableName sourceTable, TableName restoredTable) + throws Exception { + try (Connection conn = ConnectionFactory.createConnection(conf1); + BackupAdminImpl backupAdmin = new BackupAdminImpl(conn)) { + backupAdmin.restore(BackupUtils.createRestoreRequest(BACKUP_ROOT_DIR, backupId, false, + new TableName[] { sourceTable }, new TableName[] { restoredTable }, true)); + } + } + + private static void restoreAndAssertAllRowsPresent(String backupId, String restoredTableName, + String message) throws Exception { + TableName restoredTable = TableName.valueOf(restoredTableName); + restore(backupId, table1, restoredTable); + assertEquals(TEST_UTIL.countRows(table1), TEST_UTIL.countRows(restoredTable), message); + } + + private static Long boundary(BackupSystemTable sysTable, String host) throws IOException { + return sysTable.readLogTimestampMap(BACKUP_ROOT_DIR).get(table1).get(host); + } + + private static void moveRegionAndWait(TableName table, HRegionServer destination) + throws Exception { + try (Admin admin = TEST_UTIL.getConnection().getAdmin()) { + RegionInfo region = admin.getRegions(table).get(0); + admin.move(region.getEncodedNameAsBytes(), destination.getServerName()); + } + TEST_UTIL.waitFor(60_000, () -> !destination.getRegions(table).isEmpty()); + TEST_UTIL.waitUntilAllRegionsAssigned(table); + } + + /** + * Tests that when a full backup is taken while an RS is offline (with WALs in oldlogs), the + * offline host's timestamps are recorded so subsequent incremental backups don't reinclude those + * WALs. + */ + @Test + public void testBackupWithOfflineRS() throws Exception { + LOG.info("Starting testBackupWithOfflineRS"); + + SingleProcessHBaseCluster cluster = TEST_UTIL.getMiniHBaseCluster(); + List tables = Lists.newArrayList(table1); + + HRegionServer rsBeforeStop = cluster.startRegionServerAndWait(10000).getRegionServer(); + moveRegionAndWait(table1, rsBeforeStop); + + LOG.info("Inserting data to generate WAL entries"); + try (Connection conn = ConnectionFactory.createConnection(conf1)) { + insertIntoTable(conn, table1, famName, 2, 100); + } + + String offlineHost = + rsBeforeStop.getServerName().getHostname() + ":" + rsBeforeStop.getServerName().getPort(); + LOG.info("Stopping RS: {}", offlineHost); + + stopRegionServerAndWait(rsBeforeStop.getServerName()); + + LOG.info("Taking full backup (with offline RS WALs in oldlogs)"); + String fullBackupId = fullTableBackup(tables); + assertTrue(checkSucceeded(fullBackupId), "Full backup should succeed"); + + try (BackupSystemTable sysTable = new BackupSystemTable(TEST_UTIL.getConnection())) { + Map> timestamps = sysTable.readLogTimestampMap(BACKUP_ROOT_DIR); + Map rsTimestamps = timestamps.get(table1); + LOG.info("RS timestamps after full backup: {}", rsTimestamps); + + Long tsAfterFullBackup = rsTimestamps.get(offlineHost); + assertNotNull(tsAfterFullBackup, + "Offline host should have timestamp recorded in trslm after full backup"); + + LOG.info("Taking incremental backup (should NOT include offline RS WALs)"); + String incrBackupId = incrementalTableBackup(tables); + assertTrue(checkSucceeded(incrBackupId), "Incremental backup should succeed"); + + assertEquals(tsAfterFullBackup, boundary(sysTable, offlineHost), + "Incremental backup moved the boundary of the offline host"); + } + } + + /** + * Tests that WALs written to an RS after a full backup are correctly included in the subsequent + * incremental backup, even if that RS has gone offline before the incremental runs. + */ + @Test + public void testRSGoesOfflineAfterFullBackupBeforeIncremental() throws Exception { + SingleProcessHBaseCluster cluster = TEST_UTIL.getMiniHBaseCluster(); + List tables = Lists.newArrayList(table1); + + HRegionServer rsBeforeStop = cluster.startRegionServerAndWait(10000).getRegionServer(); + + try (Connection conn = ConnectionFactory.createConnection(conf1)) { + insertIntoTable(conn, table1, famName, 3, 50); + } + + String fullBackupId = fullTableBackup(tables); + assertTrue(checkSucceeded(fullBackupId), "Full backup should succeed"); + + moveRegionAndWait(table1, rsBeforeStop); + try (Connection conn = ConnectionFactory.createConnection(conf1)) { + insertIntoTable(conn, table1, famName, 4, 50); + } + + rsBeforeStop.getWalRoller().requestRollAll(); + rsBeforeStop.getWalRoller().waitUntilWalRollFinished(); + String offlineHost = + rsBeforeStop.getServerName().getHostname() + ":" + rsBeforeStop.getServerName().getPort(); + + stopRegionServerAndWait(rsBeforeStop.getServerName()); + + String incrBackupId = incrementalTableBackup(tables); + assertTrue(checkSucceeded(incrBackupId), "Incremental backup should succeed"); + + restoreAndAssertAllRowsPresent(incrBackupId, "table1_rs_offline_after_full_backup", + "Restored table should contain the rows written to the RS that went offline"); + + try (BackupSystemTable sysTable = new BackupSystemTable(TEST_UTIL.getConnection())) { + Long boundaryAfterIncr1 = boundary(sysTable, offlineHost); + assertNotNull(boundaryAfterIncr1, + "Offline RS should have a timestamp boundary after the incremental backed up its WALs"); + Long staleLogRoll = + sysTable.readRegionServerLastLogRollResult(BACKUP_ROOT_DIR).get(offlineHost); + + String incrBackupId2 = incrementalTableBackup(tables); + assertTrue(checkSucceeded(incrBackupId2), "Second incremental backup should succeed"); + Long boundaryAfterIncr2 = boundary(sysTable, offlineHost); + + String incrBackupId3 = incrementalTableBackup(tables); + assertTrue(checkSucceeded(incrBackupId3), "Third incremental backup should succeed"); + Long boundaryAfterIncr3 = boundary(sysTable, offlineHost); + + String incrBackupId4 = incrementalTableBackup(tables); + assertTrue(checkSucceeded(incrBackupId4), "Fourth incremental backup should succeed"); + Long boundaryAfterIncr4 = boundary(sysTable, offlineHost); + + String message = "Incremental backup moved the boundary of the offline RS, which would back " + + "up its WALs again (last log roll = " + staleLogRoll + ")"; + assertAll(() -> assertEquals(boundaryAfterIncr1, boundaryAfterIncr2, message), + () -> assertEquals(boundaryAfterIncr1, boundaryAfterIncr3, message), + () -> assertEquals(boundaryAfterIncr1, boundaryAfterIncr4, message)); + } + } + + /** + * Tests that a brand-new RS that comes online and goes offline before any backup correctly has + * its WALs covered by the full backup. + */ + @Test + public void testTransientRSBeforeFullBackup() throws Exception { + SingleProcessHBaseCluster cluster = TEST_UTIL.getMiniHBaseCluster(); + List tables = Lists.newArrayList(table1); + + HRegionServer transientRS = cluster.startRegionServerAndWait(10000).getRegionServer(); + moveRegionAndWait(table1, transientRS); + try (Connection conn = ConnectionFactory.createConnection(conf1)) { + insertIntoTable(conn, table1, famName, 5, 50); + } + transientRS.getWalRoller().requestRollAll(); + transientRS.getWalRoller().waitUntilWalRollFinished(); + String transientHost = + transientRS.getServerName().getHostname() + ":" + transientRS.getServerName().getPort(); + + stopRegionServerAndWait(transientRS.getServerName()); + + String fullBackupId = fullTableBackup(tables); + assertTrue(checkSucceeded(fullBackupId), "Full backup should succeed"); + + try (BackupSystemTable sysTable = new BackupSystemTable(TEST_UTIL.getConnection())) { + Long boundaryAfterFullBackup = boundary(sysTable, transientHost); + assertNotNull(boundaryAfterFullBackup, + "Transient RS that went offline before full backup should have its WAL boundary recorded"); + + String incrBackupId = incrementalTableBackup(tables); + assertTrue(checkSucceeded(incrBackupId), + "Incremental backup after transient RS should succeed"); + + assertEquals(boundaryAfterFullBackup, boundary(sysTable, transientHost), + "Incremental backup moved the boundary of the transient RS"); + } + } + + /** + * Tests that WALs from an RS that comes online and goes offline between a full backup and an + * incremental backup are correctly included in the incremental backup and not re-included in + * subsequent incremental backups. + */ + @Test + public void testTransientRSAfterFullBackupBeforeIncremental() throws Exception { + SingleProcessHBaseCluster cluster = TEST_UTIL.getMiniHBaseCluster(); + List tables = Lists.newArrayList(table1); + + try (Connection conn = ConnectionFactory.createConnection(conf1)) { + insertIntoTable(conn, table1, famName, 6, 50); + } + + String fullBackupId = fullTableBackup(tables); + assertTrue(checkSucceeded(fullBackupId), "Full backup should succeed"); + + HRegionServer transientRS = cluster.startRegionServerAndWait(10000).getRegionServer(); + moveRegionAndWait(table1, transientRS); + try (Connection conn = ConnectionFactory.createConnection(conf1)) { + insertIntoTable(conn, table1, famName, 7, 50); + } + transientRS.getWalRoller().requestRollAll(); + transientRS.getWalRoller().waitUntilWalRollFinished(); + String transientHost = + transientRS.getServerName().getHostname() + ":" + transientRS.getServerName().getPort(); + + stopRegionServerAndWait(transientRS.getServerName()); + + String incrBackupId = incrementalTableBackup(tables); + assertTrue(checkSucceeded(incrBackupId), "Incremental backup should succeed"); + + restoreAndAssertAllRowsPresent(incrBackupId, "table1_transient_rs_after_full_backup", + "Restored table should contain the rows written to the transient RS"); + + try (BackupSystemTable sysTable = new BackupSystemTable(TEST_UTIL.getConnection())) { + Long boundaryAfterIncr1 = boundary(sysTable, transientHost); + assertNotNull(boundaryAfterIncr1, + "Transient RS should have a timestamp boundary after the incremental backed up its WALs"); + + String incrBackupId2 = incrementalTableBackup(tables); + assertTrue(checkSucceeded(incrBackupId2), "Second incremental backup should succeed"); + + assertEquals(boundaryAfterIncr1, boundary(sysTable, transientHost), + "Second incremental backup moved the boundary of the transient RS"); + } + } + + /** + * Tests that edits written during a full backup, after its log roll and table snapshot, to an RS + * that started after that log roll, are included in the next incremental backup. + */ + @Test + public void testRSStartedDuringFullBackup() throws Exception { + SingleProcessHBaseCluster cluster = TEST_UTIL.getMiniHBaseCluster(); + List tables = Lists.newArrayList(table1); + + afterSnapshotHook = () -> { + try { + HRegionServer newRS = cluster.startRegionServerAndWait(10000).getRegionServer(); + moveRegionAndWait(table1, newRS); + try (Connection conn = ConnectionFactory.createConnection(conf1)) { + insertIntoTable(conn, table1, famName, 8, 50).close(); + } + stopRegionServerAndWait(newRS.getServerName()); + } catch (Exception e) { + throw new RuntimeException(e); + } + }; + conf1.set(TableBackupClient.BACKUP_CLIENT_IMPL_CLASS, + FullTableBackupClientWithHook.class.getName()); + String fullBackupId; + try { + fullBackupId = fullTableBackup(tables); + } finally { + conf1.unset(TableBackupClient.BACKUP_CLIENT_IMPL_CLASS); + } + assertTrue(checkSucceeded(fullBackupId), "Full backup should succeed"); + + String incrBackupId = incrementalTableBackup(tables); + assertTrue(checkSucceeded(incrBackupId), "Incremental backup should succeed"); + + restoreAndAssertAllRowsPresent(incrBackupId, "table1_rs_started_during_full_backup", + "Restored table should contain all rows, including those written to the new RS"); + } + + /** + * Tests that a row deleted before a full backup is not brought back by the WALs of a region + * server that went offline between two backups. The row and its delete marker are compacted away + * before the full backup, so only the offline RS's old WAL, which still holds the put, could + * bring it back. + */ + @Test + public void testDeletedRowIsNotResurrectedByOfflineRSWALs() throws Exception { + SingleProcessHBaseCluster cluster = TEST_UTIL.getMiniHBaseCluster(); + List tables = Lists.newArrayList(table1); + byte[] deletedRow = Bytes.toBytes("row-deleted"); + + HRegionServer offlineRS = cluster.startRegionServerAndWait(10000).getRegionServer(); + HRegionServer liveRS = cluster.startRegionServerAndWait(10000).getRegionServer(); + + try (Table table = TEST_UTIL.getConnection().getTable(table1)) { + moveRegionAndWait(table1, offlineRS); + + String firstFullBackupId = fullTableBackup(tables); + assertTrue(checkSucceeded(firstFullBackupId), "First full backup should succeed"); + + for (int i = 0; i < 10; i++) { + table.put(new Put(Bytes.toBytes("row-kept-" + i)).addColumn(famName, qualName, + Bytes.toBytes("value"))); + } + table.put(new Put(deletedRow).addColumn(famName, qualName, Bytes.toBytes("value"))); + for (HRegion region : offlineRS.getRegions(table1)) { + region.flush(true); + } + + moveRegionAndWait(table1, liveRS); + + table.delete(new Delete(deletedRow)); + for (HRegion region : liveRS.getRegions(table1)) { + region.flush(true); + region.compact(true); + } + + assertTrue(table.get(new Get(deletedRow)).isEmpty(), + "Deleted row should be gone from the source table"); + Scan rawScan = new Scan().withStartRow(deletedRow).withStopRow(deletedRow, true).setRaw(true); + try (ResultScanner scanner = table.getScanner(rawScan)) { + assertNull(scanner.next(), + "Major compaction should have purged the deleted row and its delete marker"); + } + + stopRegionServerAndWait(offlineRS.getServerName()); + + String secondFullBackupId = fullTableBackup(tables); + assertTrue(checkSucceeded(secondFullBackupId), "Second full backup should succeed"); + String incrBackupId = incrementalTableBackup(tables); + assertTrue(checkSucceeded(incrBackupId), "Incremental backup should succeed"); + + TableName restoredTable = TableName.valueOf("table1_deleted_row_not_resurrected"); + restore(incrBackupId, table1, restoredTable); + try (Table restored = TEST_UTIL.getConnection().getTable(restoredTable)) { + assertTrue(restored.get(new Get(deletedRow)).isEmpty(), + "Deleted row was brought back by the WALs of the offline RS"); + } + assertEquals(TEST_UTIL.countRows(table1), TEST_UTIL.countRows(restoredTable), + "Restored table should match the source table"); + } + } +} diff --git a/hbase-backup/src/test/java/org/apache/hadoop/hbase/backup/TestBackupUtils.java b/hbase-backup/src/test/java/org/apache/hadoop/hbase/backup/TestBackupUtils.java index 0fc4393a0a61..dc009b82a558 100644 --- a/hbase-backup/src/test/java/org/apache/hadoop/hbase/backup/TestBackupUtils.java +++ b/hbase-backup/src/test/java/org/apache/hadoop/hbase/backup/TestBackupUtils.java @@ -28,6 +28,8 @@ import java.time.ZoneId; import java.util.ArrayList; import java.util.List; +import java.util.Map; +import java.util.Set; import java.util.TimeZone; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.FileSystem; @@ -184,6 +186,84 @@ public void testGetValidWalDirOneMsBeforeMidnightUTC() throws IOException { testGetValidWalDirs(startAndEndTime, startAndEndTime, walDir, walDateDirs, 1, walDateDirs); } + @Test + public void testGetRolledHostsKeepsOnlyHostsWhoseRollResultChanged() { + Map previousLogRolls = + Map.of("rolled:16020", 100L, "offline:16020", 200L, "removed:16020", 300L); + Map latestLogRolls = + Map.of("rolled:16020", 150L, "offline:16020", 200L, "new:16020", 400L); + + assertEquals(Map.of("rolled:16020", 150L, "new:16020", 400L), + BackupUtils.getRolledHosts(previousLogRolls, latestLogRolls)); + } + + @Test + public void testComputeLogBoundariesUsesRollResultForRolledHosts() throws IOException { + Map rolledHosts = Map.of("rolled:16020", 100L); + List coveredLogs = List.of("/hbase/oldWALs/rolled%2C16020%2C1.500"); + List pendingLogs = List.of("/hbase/WALs/rolled,16020,1/rolled%2C16020%2C1.50"); + + assertEquals(Map.of("rolled:16020", 100L), + BackupUtils.computeLogBoundaries(rolledHosts, coveredLogs, pendingLogs)); + } + + @Test + public void testComputeLogBoundariesUsesNewestCoveredLogForOtherHosts() throws IOException { + List coveredLogs = + List.of("/hbase/oldWALs/offline%2C16020%2C1.200", "/hbase/oldWALs/offline%2C16020%2C1.300", + "/hbase/WALs/offline,16020,1/offline%2C16020%2C1.250"); + + assertEquals(Map.of("offline:16020", 300L), + BackupUtils.computeLogBoundaries(Map.of(), coveredLogs, List.of())); + } + + @Test + public void testComputeLogBoundariesCapsBelowOldestPendingLog() throws IOException { + List coveredLogs = List.of("/hbase/oldWALs/splitting%2C16020%2C1.400"); + List pendingLogs = + List.of("/hbase/WALs/splitting,16020,1-splitting/splitting%2C16020%2C1.350", + "/hbase/WALs/splitting,16020,1-splitting/splitting%2C16020%2C1.370"); + + assertEquals(Map.of("splitting:16020", 349L), + BackupUtils.computeLogBoundaries(Map.of(), coveredLogs, pendingLogs)); + } + + @Test + public void testComputeLogBoundariesKeepsHostWithOnlyPendingLogs() throws IOException { + List pendingLogs = List.of("/hbase/WALs/joined,16020,1/joined%2C16020%2C1.500"); + + assertEquals(Map.of("joined:16020", 499L), + BackupUtils.computeLogBoundaries(Map.of(), List.of(), pendingLogs)); + } + + @Test + public void testComputeLogBoundariesSkipsUnparseableLogs() throws IOException { + List coveredLogs = List.of("/hbase/oldWALs/not-a-wal"); + + assertEquals(Map.of("rolled:16020", 100L), + BackupUtils.computeLogBoundaries(Map.of("rolled:16020", 100L), coveredLogs, List.of())); + } + + @Test + public void testComputeLogBoundariesKeepsPreviousBoundaryWhileHostHasLogs() throws IOException { + Map previousBoundaries = Map.of("stuck:16020", 1200L, "gone:16020", 800L); + + assertEquals(Map.of("stuck:16020", 1200L), BackupUtils.computeLogBoundaries(Map.of(), + previousBoundaries, Set.of("stuck:16020"), List.of(), List.of())); + } + + @Test + public void testComputeLogBoundariesPrefersNewBoundaryOverPreviousOne() throws IOException { + Map previousBoundaries = + Map.of("rolled:16020", 100L, "offline:16020", 200L, "pending:16020", 300L); + List coveredLogs = List.of("/hbase/oldWALs/offline%2C16020%2C1.250"); + List pendingLogs = List.of("/hbase/WALs/pending,16020,1/pending%2C16020%2C1.400"); + + assertEquals(Map.of("rolled:16020", 150L, "offline:16020", 250L, "pending:16020", 399L), + BackupUtils.computeLogBoundaries(Map.of("rolled:16020", 150L), previousBoundaries, + Set.of("rolled:16020", "offline:16020", "pending:16020"), coveredLogs, pendingLogs)); + } + protected void testGetValidWalDirs(long startTime, long endTime, Path walDir, List availableWalDateDirs, int numExpectedValidWalDirs, List expectedValidWalDirs) throws IOException { diff --git a/hbase-backup/src/test/java/org/apache/hadoop/hbase/backup/TestIncrementalBackupManager.java b/hbase-backup/src/test/java/org/apache/hadoop/hbase/backup/TestIncrementalBackupManager.java index b8334d3d0e0a..fc9b45debd7c 100644 --- a/hbase-backup/src/test/java/org/apache/hadoop/hbase/backup/TestIncrementalBackupManager.java +++ b/hbase-backup/src/test/java/org/apache/hadoop/hbase/backup/TestIncrementalBackupManager.java @@ -17,9 +17,12 @@ */ package org.apache.hadoop.hbase.backup; +import static org.junit.jupiter.api.Assertions.assertAll; +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 java.util.ArrayList; import java.util.List; import java.util.Map; import org.apache.hadoop.conf.Configuration; @@ -33,9 +36,11 @@ import org.apache.hadoop.hbase.backup.util.BackupUtils; import org.apache.hadoop.hbase.client.Connection; import org.apache.hadoop.hbase.client.ConnectionFactory; +import org.apache.hadoop.hbase.regionserver.HRegionServer; import org.apache.hadoop.hbase.testclassification.LargeTests; import org.apache.hadoop.hbase.util.CommonFSUtils; import org.apache.hadoop.hbase.util.EnvironmentEdgeManager; +import org.apache.hadoop.hbase.util.JVMClusterUtil; import org.apache.hadoop.hbase.wal.AbstractFSWALProvider; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.Tag; @@ -101,4 +106,279 @@ private void testCollectWALFiles(boolean separateOldLogDir) throws Exception { } } } + + /** + * WALs can be archived out of order, so a region server that took part in the log roll can have + * an archived WAL newer than its roll result while an older WAL is still in the WALs directory. + * The newer archived WAL must be deferred to a later backup instead of being backed up twice, and + * the older WAL must still end up in a backup. + */ + @Test + public void testOutOfOrderArchivedWALDoesNotSkipOlderWAL() throws Exception { + List tables = List.of(table1); + HRegionServer rs = TEST_UTIL.getMiniHBaseCluster().getRegionServer(0); + ServerName serverName = rs.getServerName(); + Path walRootDir = CommonFSUtils.getWALRootDir(conf1); + FileSystem fs = walRootDir.getFileSystem(conf1); + + try (Connection conn = ConnectionFactory.createConnection(conf1); + BackupAdminImpl backupAdmin = new BackupAdminImpl(conn)) { + String fullBackupId = + backupAdmin.backupTables(createBackupRequest(BackupType.FULL, tables, BACKUP_ROOT_DIR)); + assertTrue(checkSucceeded(fullBackupId)); + + long olderWALTs = EnvironmentEdgeManager.currentTime() + 1; + long newerWALTs = olderWALTs + 1; + Path walDir = + new Path(walRootDir, AbstractFSWALProvider.getWALDirectoryName(serverName.toString())); + Path olderWAL = + new Path(walDir, serverName.toString() + BackupUtils.LOGNAME_SEPARATOR + olderWALTs); + Path archiveDir = new Path(walRootDir, + AbstractFSWALProvider.getWALArchiveDirectoryName(conf1, serverName.toString())); + Path newerArchivedWAL = + new Path(archiveDir, "wal" + BackupUtils.LOGNAME_SEPARATOR + newerWALTs); + fs.create(olderWAL).close(); + fs.mkdirs(archiveDir); + fs.create(newerArchivedWAL).close(); + + try { + List firstBackupFiles; + try (IncrementalBackupManager manager = new IncrementalBackupManager(conn, conf1)) { + BackupInfo backupInfo = manager.createBackupInfo("backup_incr_1", BackupType.INCREMENTAL, + tables, BACKUP_ROOT_DIR, -1, -1, false, false); + Map boundaries = manager.getIncrBackupLogFileMap(); + manager.writeRegionServerLogTimestamp(backupInfo.getTables(), boundaries); + firstBackupFiles = backupInfo.getIncrBackupFileList(); + } + assertFalse(firstBackupFiles.contains(newerArchivedWAL.toString()), + "Archived WAL newer than the roll result should be deferred to a later backup: " + + firstBackupFiles); + + TEST_UTIL.waitFor(30_000, () -> EnvironmentEdgeManager.currentTime() > newerWALTs); + rs.getWalRoller().requestRollAll(); + rs.getWalRoller().waitUntilWalRollFinished(); + + List secondBackupFiles; + try (IncrementalBackupManager manager = new IncrementalBackupManager(conn, conf1)) { + BackupInfo backupInfo = manager.createBackupInfo("backup_incr_2", BackupType.INCREMENTAL, + tables, BACKUP_ROOT_DIR, -1, -1, false, false); + manager.getIncrBackupLogFileMap(); + secondBackupFiles = backupInfo.getIncrBackupFileList(); + } + + assertTrue( + firstBackupFiles.contains(olderWAL.toString()) + || secondBackupFiles.contains(olderWAL.toString()), + "WAL " + olderWAL + " was not included in any backup. First backup: " + firstBackupFiles + + ", second backup: " + secondBackupFiles); + } finally { + fs.delete(olderWAL, false); + fs.delete(newerArchivedWAL, false); + } + } + } + + /** + * A dead region server's WALs stay in its -splitting directory until WAL splitting finishes. An + * incremental backup that runs during the split must not lose them once they are archived. + */ + @Test + public void testWALOfDeadServerStillSplittingIsBackedUpAfterArchiving() throws Exception { + List tables = List.of(table1); + ServerName deadServer = ServerName.valueOf("deadhost", 16020, 1001L); + Path walRootDir = CommonFSUtils.getWALRootDir(conf1); + FileSystem fs = walRootDir.getFileSystem(conf1); + Path splittingDir = splittingDir(walRootDir, deadServer); + Path archiveDir = new Path(walRootDir, + AbstractFSWALProvider.getWALArchiveDirectoryName(conf1, deadServer.toString())); + + try (Connection conn = ConnectionFactory.createConnection(conf1); + BackupAdminImpl backupAdmin = new BackupAdminImpl(conn)) { + String fullBackupId = + backupAdmin.backupTables(createBackupRequest(BackupType.FULL, tables, BACKUP_ROOT_DIR)); + assertTrue(checkSucceeded(fullBackupId)); + + long walTs = EnvironmentEdgeManager.currentTime() + 1; + Path splittingWAL = new Path(splittingDir, walName(deadServer, walTs)); + Path archivedWAL = archivedWAL(archiveDir, walTs); + fs.create(splittingWAL).close(); + + try { + // Make the live servers' roll results, and so the next boundaries, newer than the WAL. + TEST_UTIL.waitFor(30_000, () -> EnvironmentEdgeManager.currentTime() > walTs); + rollAllLiveRegionServers(); + + List firstBackupFiles = runIncrementalBackup(conn, tables, "backup_split_1"); + + fs.mkdirs(archiveDir); + assertTrue(fs.rename(splittingWAL, archivedWAL)); + fs.delete(splittingDir, true); + + List secondBackupFiles = runIncrementalBackup(conn, tables, "backup_split_2"); + + assertTrue( + firstBackupFiles.contains(splittingWAL.toString()) + || secondBackupFiles.contains(archivedWAL.toString()), + "WAL of the dead server was not included in any backup. First backup: " + firstBackupFiles + + ", second backup: " + secondBackupFiles); + } finally { + fs.delete(splittingDir, true); + fs.delete(archivedWAL, false); + } + } + } + + /** + * WAL splitting archives a dead region server's WALs in parallel, so a newer WAL can be archived + * while older ones are still in the -splitting directory. Backing up the newer WAL must not move + * the boundary past the older ones. + */ + @Test + public void testOlderWALsOfDeadServerAreNotSkippedWhenNewerWALIsArchivedFirst() throws Exception { + List tables = List.of(table1); + ServerName deadServer = ServerName.valueOf("deadhost", 16020, 2002L); + Path walRootDir = CommonFSUtils.getWALRootDir(conf1); + FileSystem fs = walRootDir.getFileSystem(conf1); + Path splittingDir = splittingDir(walRootDir, deadServer); + Path archiveDir = new Path(walRootDir, + AbstractFSWALProvider.getWALArchiveDirectoryName(conf1, deadServer.toString())); + + try (Connection conn = ConnectionFactory.createConnection(conf1); + BackupAdminImpl backupAdmin = new BackupAdminImpl(conn)) { + String fullBackupId = + backupAdmin.backupTables(createBackupRequest(BackupType.FULL, tables, BACKUP_ROOT_DIR)); + assertTrue(checkSucceeded(fullBackupId)); + + long oldestTs = EnvironmentEdgeManager.currentTime() + 1; + long middleTs = oldestTs + 1; + long newestTs = oldestTs + 2; + Path oldestSplittingWAL = new Path(splittingDir, walName(deadServer, oldestTs)); + Path middleSplittingWAL = new Path(splittingDir, walName(deadServer, middleTs)); + Path oldestArchivedWAL = archivedWAL(archiveDir, oldestTs); + Path middleArchivedWAL = archivedWAL(archiveDir, middleTs); + Path newestArchivedWAL = archivedWAL(archiveDir, newestTs); + fs.create(oldestSplittingWAL).close(); + fs.create(middleSplittingWAL).close(); + fs.mkdirs(archiveDir); + fs.create(newestArchivedWAL).close(); + + try { + TEST_UTIL.waitFor(30_000, () -> EnvironmentEdgeManager.currentTime() > newestTs); + + List firstBackupFiles = runIncrementalBackup(conn, tables, "backup_order_1"); + + assertTrue(fs.rename(oldestSplittingWAL, oldestArchivedWAL)); + assertTrue(fs.rename(middleSplittingWAL, middleArchivedWAL)); + fs.delete(splittingDir, true); + + List secondBackupFiles = runIncrementalBackup(conn, tables, "backup_order_2"); + + assertAll( + () -> assertTrue( + firstBackupFiles.contains(oldestSplittingWAL.toString()) + || secondBackupFiles.contains(oldestArchivedWAL.toString()), + "Oldest WAL of the dead server was not included in any backup. First backup: " + + firstBackupFiles + ", second backup: " + secondBackupFiles), + () -> assertTrue( + firstBackupFiles.contains(middleSplittingWAL.toString()) + || secondBackupFiles.contains(middleArchivedWAL.toString()), + "Middle WAL of the dead server was not included in any backup. First backup: " + + firstBackupFiles + ", second backup: " + secondBackupFiles)); + } finally { + fs.delete(splittingDir, true); + fs.delete(oldestArchivedWAL, false); + fs.delete(middleArchivedWAL, false); + fs.delete(newestArchivedWAL, false); + } + } + } + + /** + * A dead region server can keep an old WAL, already covered by its boundary, in its -splitting + * directory across several incremental backups. Its boundary must survive those backups, or a + * later backup would include that WAL again and could bring back deleted data. + */ + @Test + public void testDeadServerKeepsBoundaryWhileOldWALIsStillSplitting() throws Exception { + List tables = List.of(table1); + ServerName deadServer = ServerName.valueOf("deadhost", 16020, 3003L); + Path walRootDir = CommonFSUtils.getWALRootDir(conf1); + FileSystem fs = walRootDir.getFileSystem(conf1); + Path splittingDir = splittingDir(walRootDir, deadServer); + Path archiveDir = new Path(walRootDir, + AbstractFSWALProvider.getWALArchiveDirectoryName(conf1, deadServer.toString())); + + try (Connection conn = ConnectionFactory.createConnection(conf1); + BackupAdminImpl backupAdmin = new BackupAdminImpl(conn)) { + String fullBackupId = + backupAdmin.backupTables(createBackupRequest(BackupType.FULL, tables, BACKUP_ROOT_DIR)); + assertTrue(checkSucceeded(fullBackupId)); + + long stuckTs = EnvironmentEdgeManager.currentTime() + 1; + long archivedTs = stuckTs + 1; + Path archivedWAL = archivedWAL(archiveDir, archivedTs); + Path stuckSplittingWAL = new Path(splittingDir, walName(deadServer, stuckTs)); + Path stuckArchivedWAL = archivedWAL(archiveDir, stuckTs); + fs.mkdirs(archiveDir); + fs.create(archivedWAL).close(); + + try { + TEST_UTIL.waitFor(30_000, () -> EnvironmentEdgeManager.currentTime() > archivedTs); + List firstBackupFiles = runIncrementalBackup(conn, tables, "backup_stuck_1"); + assertTrue(firstBackupFiles.contains(archivedWAL.toString()), + "Archived WAL of the dead server should be in the first backup: " + firstBackupFiles); + + fs.create(stuckSplittingWAL).close(); + List laterBackupFiles = new ArrayList<>(); + laterBackupFiles.addAll(runIncrementalBackup(conn, tables, "backup_stuck_2")); + laterBackupFiles.addAll(runIncrementalBackup(conn, tables, "backup_stuck_3")); + + assertTrue(fs.rename(stuckSplittingWAL, stuckArchivedWAL)); + fs.delete(splittingDir, true); + laterBackupFiles.addAll(runIncrementalBackup(conn, tables, "backup_stuck_4")); + + assertFalse( + laterBackupFiles.contains(stuckSplittingWAL.toString()) + || laterBackupFiles.contains(stuckArchivedWAL.toString()), + "WAL already covered by the dead server's boundary was backed up again: " + + laterBackupFiles); + } finally { + fs.delete(splittingDir, true); + fs.delete(archivedWAL, false); + fs.delete(stuckArchivedWAL, false); + } + } + } + + private static Path splittingDir(Path walRootDir, ServerName serverName) { + return new Path(walRootDir, AbstractFSWALProvider.getWALDirectoryName(serverName.toString()) + + AbstractFSWALProvider.SPLITTING_EXT); + } + + private static Path archivedWAL(Path archiveDir, long ts) { + return new Path(archiveDir, "wal" + BackupUtils.LOGNAME_SEPARATOR + ts); + } + + private static String walName(ServerName serverName, long ts) { + return serverName.toString() + BackupUtils.LOGNAME_SEPARATOR + ts; + } + + private static void rollAllLiveRegionServers() throws Exception { + for (JVMClusterUtil.RegionServerThread rst : TEST_UTIL.getMiniHBaseCluster() + .getLiveRegionServerThreads()) { + rst.getRegionServer().getWalRoller().requestRollAll(); + rst.getRegionServer().getWalRoller().waitUntilWalRollFinished(); + } + } + + private static List runIncrementalBackup(Connection conn, List tables, + String backupId) throws Exception { + try (IncrementalBackupManager manager = new IncrementalBackupManager(conn, conf1)) { + BackupInfo backupInfo = manager.createBackupInfo(backupId, BackupType.INCREMENTAL, tables, + BACKUP_ROOT_DIR, -1, -1, false, false); + Map boundaries = manager.getIncrBackupLogFileMap(); + manager.writeRegionServerLogTimestamp(backupInfo.getTables(), boundaries); + return backupInfo.getIncrBackupFileList(); + } + } } diff --git a/hbase-backup/src/test/java/org/apache/hadoop/hbase/backup/impl/TestFullTableBackupClientLogBoundaries.java b/hbase-backup/src/test/java/org/apache/hadoop/hbase/backup/impl/TestFullTableBackupClientLogBoundaries.java new file mode 100644 index 000000000000..91aaa76c6425 --- /dev/null +++ b/hbase-backup/src/test/java/org/apache/hadoop/hbase/backup/impl/TestFullTableBackupClientLogBoundaries.java @@ -0,0 +1,148 @@ +/* + * 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.backup.impl; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +import java.io.IOException; +import java.util.List; +import java.util.Map; +import java.util.Set; +import org.apache.hadoop.fs.FileSystem; +import org.apache.hadoop.fs.Path; +import org.apache.hadoop.hbase.HBaseTestingUtil; +import org.apache.hadoop.hbase.HConstants; +import org.apache.hadoop.hbase.ServerName; +import org.apache.hadoop.hbase.client.Admin; +import org.apache.hadoop.hbase.testclassification.SmallTests; +import org.apache.hadoop.hbase.wal.AbstractFSWALProvider; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Tag; +import org.junit.jupiter.api.Test; + +@Tag(SmallTests.TAG) +public class TestFullTableBackupClientLogBoundaries { + + private static final HBaseTestingUtil TEST_UTIL = new HBaseTestingUtil(); + + private FileSystem fs; + private Path walRootDir; + + @BeforeEach + public void setUp() throws IOException { + fs = TEST_UTIL.getTestFileSystem(); + walRootDir = TEST_UTIL.getDataTestDirOnTestFS("walRoot"); + } + + @AfterEach + public void tearDown() throws IOException { + fs.delete(walRootDir, true); + } + + @Test + public void testLiveServerWithoutRollResultIsCappedBelowItsOldestWAL() throws IOException { + ServerName joined = ServerName.valueOf("joined", 16020, 1L); + createWAL(walDir(joined), joined, 500); + createWAL(walDir(joined), joined, 600); + Path metaWAL = + new Path(walDir(joined), walName(joined, 450) + AbstractFSWALProvider.META_WAL_PROVIDER_ID); + fs.create(metaWAL).close(); + + assertEquals(Map.of("joined:16020", 499L), computeLogBoundaries(Map.of(), Set.of(joined))); + } + + @Test + public void testDeadServerGetsItsNewestWAL() throws IOException { + ServerName dead = ServerName.valueOf("dead", 16020, 2L); + createWAL(new Path(walDir(dead).toString() + AbstractFSWALProvider.SPLITTING_EXT), dead, 300); + createWAL(new Path(walRootDir, HConstants.HREGION_OLDLOGDIR_NAME), dead, 350); + + assertEquals(Map.of("dead:16020", 350L), computeLogBoundaries(Map.of(), Set.of())); + } + + @Test + public void testRestartedServerCoversOldInstanceOnly() throws IOException { + ServerName oldInstance = ServerName.valueOf("restarted", 16020, 3L); + ServerName newInstance = ServerName.valueOf("restarted", 16020, 4L); + createWAL(walDir(oldInstance), oldInstance, 700); + createWAL(walDir(newInstance), newInstance, 800); + + assertEquals(Map.of("restarted:16020", 700L), + computeLogBoundaries(Map.of(), Set.of(newInstance))); + } + + @Test + public void testRolledServerKeepsItsRollResult() throws IOException { + ServerName rolled = ServerName.valueOf("rolled", 16020, 5L); + createWAL(walDir(rolled), rolled, 900); + createWAL(new Path(walRootDir, HConstants.HREGION_OLDLOGDIR_NAME), rolled, 950); + + assertEquals(Map.of("rolled:16020", 850L), + computeLogBoundaries(Map.of("rolled:16020", 850L), Set.of(rolled))); + } + + @Test + public void testServerRegisteringAfterTheListingGetsNoBoundary() throws IOException { + ServerName joined = ServerName.valueOf("joined", 16020, 6L); + ServerName late = ServerName.valueOf("late", 16020, 7L); + createWAL(walDir(joined), joined, 500); + Admin admin = mock(Admin.class); + when(admin.getRegionServers()).thenAnswer(invocation -> { + createWAL(walDir(late), late, 600); + return List.of(joined); + }); + + assertEquals(Map.of("joined:16020", 499L), + FullTableBackupClient.computeLogBoundaries(fs, walRootDir, Map.of(), admin)); + } + + @Test + public void testIncompleteLiveServerListFailsTheBackup() throws IOException { + ServerName rolled = ServerName.valueOf("rolled", 16020, 8L); + ServerName joined = ServerName.valueOf("joined", 16020, 9L); + createWAL(walDir(rolled), rolled, 900); + createWAL(walDir(joined), joined, 500); + + assertThrows(IOException.class, + () -> computeLogBoundaries(Map.of("rolled:16020", 850L), Set.of(joined))); + } + + private Map computeLogBoundaries(Map rolledHosts, + Set liveServers) throws IOException { + Admin admin = mock(Admin.class); + when(admin.getRegionServers()).thenReturn(liveServers); + return FullTableBackupClient.computeLogBoundaries(fs, walRootDir, rolledHosts, admin); + } + + private Path walDir(ServerName serverName) { + return new Path(walRootDir, AbstractFSWALProvider.getWALDirectoryName(serverName.toString())); + } + + private void createWAL(Path dir, ServerName serverName, long ts) throws IOException { + fs.mkdirs(dir); + fs.create(new Path(dir, walName(serverName, ts))).close(); + } + + private static String walName(ServerName serverName, long ts) { + return serverName.toString().replace(",", "%2C") + "." + ts; + } +}