diff --git a/hbase-backup/pom.xml b/hbase-backup/pom.xml
index 7bbf59b1e96a..82ece8460281 100644
--- a/hbase-backup/pom.xml
+++ b/hbase-backup/pom.xml
@@ -167,6 +167,11 @@
junit
test
+
+ org.mockito
+ mockito-core
+ test
+
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 ccadca010562..61e24470a802 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
@@ -679,8 +679,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 a9381b1409a3..724fefb85427 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
@@ -23,18 +23,21 @@
import static org.apache.hadoop.hbase.backup.BackupRestoreConstants.DEFAULT_BACKUP_MAX_ATTEMPTS;
import static org.apache.hadoop.hbase.backup.BackupRestoreConstants.JOB_NAME_CONF_KEY;
+import com.google.errorprone.annotations.RestrictedApi;
import java.io.IOException;
import java.util.ArrayList;
-import java.util.Arrays;
import java.util.HashMap;
+import java.util.HashSet;
import java.util.List;
import java.util.Map;
+import java.util.Set;
import java.util.regex.Pattern;
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;
@@ -174,6 +177,9 @@ public void execute() throws IOException {
LogRollMasterProcedureManager.ROLLLOG_PROCEDURE_NAME, props);
Map latestLogRollsByHost = backupManager.readRegionServerLastLogRollResult();
+ Path walRootDir = CommonFSUtils.getWALRootDir(conf);
+ newTimestamps = computeLogBoundaries(walRootDir.getFileSystem(conf), walRootDir,
+ BackupUtils.getRolledHosts(previousLogRollsByHost, latestLogRollsByHost), admin);
// SNAPSHOT_TABLES:
backupInfo.setPhase(BackupPhase.SNAPSHOT);
@@ -197,49 +203,6 @@ public void execute() throws IOException {
// After this checkpoint, even if entering cancel process, will let the backup finished
backupInfo.setState(BackupState.COMPLETE);
- // Scan oldlogs for dead/decommissioned hosts and add their max WAL timestamps
- // to newTimestamps. This ensures subsequent incremental backups won't try to back up
- // WALs that are already covered by this full backup's snapshot.
- Path walRootDir = CommonFSUtils.getWALRootDir(conf);
- Path logDir = new Path(walRootDir, HConstants.HREGION_LOGDIR_NAME);
- Path oldLogDir = new Path(walRootDir, HConstants.HREGION_OLDLOGDIR_NAME);
- FileSystem fs = walRootDir.getFileSystem(conf);
-
- List allLogs = new ArrayList<>();
- for (FileStatus hostLogDir : fs.listStatus(logDir)) {
- String host = BackupUtils.parseHostNameFromLogFile(hostLogDir.getPath());
- if (host == null) {
- continue;
- }
- allLogs.addAll(Arrays.asList(fs.listStatus(hostLogDir.getPath())));
- }
- allLogs.addAll(Arrays.asList(fs.listStatus(oldLogDir)));
-
- newTimestamps = new HashMap<>();
-
- for (FileStatus log : allLogs) {
- if (AbstractFSWALProvider.isMetaFile(log.getPath())) {
- continue;
- }
- String host = BackupUtils.parseHostNameFromLogFile(log.getPath());
- if (host == null) {
- continue;
- }
- long timestamp = BackupUtils.getCreationTime(log.getPath());
- Long previousLogRoll = previousLogRollsByHost.get(host);
- Long latestLogRoll = latestLogRollsByHost.get(host);
- boolean isInactive = latestLogRoll == null || latestLogRoll.equals(previousLogRoll);
-
- if (isInactive) {
- long currentTs = newTimestamps.getOrDefault(host, 0L);
- if (timestamp > currentTs) {
- newTimestamps.put(host, timestamp);
- }
- } else {
- newTimestamps.put(host, latestLogRoll);
- }
- }
-
// The table list in backupInfo is good for both full backup and incremental backup.
// For incremental backup, it contains the incremental backup table set.
backupManager.writeRegionServerLogTimestamp(backupInfo.getTables(), newTimestamps);
@@ -264,6 +227,75 @@ public void execute() 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);
+ }
+
protected void snapshotTable(Admin admin, TableName tableName, String snapshotName)
throws IOException {
int maxAttempts = conf.getInt(BACKUP_MAX_ATTEMPTS_KEY, DEFAULT_BACKUP_MAX_ATTEMPTS);
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 ff6855e5c166..13d280d001a2 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.HashMap;
+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;
@@ -56,7 +58,10 @@ 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 {
@@ -82,6 +87,7 @@ public Map getIncrBackupLogFileMap() throws IOException {
+ "In order to create an incremental backup, at least one full backup is needed.");
}
+ Map previousLogRollByHost = readRegionServerLastLogRollResult();
if (backupInfo.getUsePreviousLogRoll()) {
LOG.info("Using previous WAL roll for backup, skipping WAL roll procedure");
} else {
@@ -94,51 +100,46 @@ public Map getIncrBackupLogFileMap() throws IOException {
LogRollMasterProcedureManager.ROLLLOG_PROCEDURE_NAME, props);
}
}
- Map newTimestamps = readRegionServerLastLogRollResult();
-
- Map latestLogRollByHost = readRegionServerLastLogRollResult();
- for (Map.Entry entry : latestLogRollByHost.entrySet()) {
- String host = entry.getKey();
- long latestLogRoll = entry.getValue();
- Long earliestTimestampToIncludeInBackup = previousTimestampMins.get(host);
-
- boolean isInactive = earliestTimestampToIncludeInBackup != null
- && earliestTimestampToIncludeInBackup > latestLogRoll;
-
- long latestTimestampToIncludeInBackup;
- if (isInactive) {
- LOG.debug("Avoided resetting latest timestamp boundary for {} from {} to {}", host,
- earliestTimestampToIncludeInBackup, latestLogRoll);
- latestTimestampToIncludeInBackup = earliestTimestampToIncludeInBackup;
- } else {
- latestTimestampToIncludeInBackup = latestLogRoll;
- }
- newTimestamps.put(host, latestTimestampToIncludeInBackup);
- }
+ Map rolledHosts =
+ BackupUtils.getRolledHosts(previousLogRollByHost, readRegionServerLastLogRollResult());
- logList = getLogFilesForNewBackup(previousTimestampMins, newTimestamps, conf, savedStartCode);
- logList = excludeProcV2WALs(logList);
+ LogFileSelection selection =
+ getLogFilesForNewBackup(previousTimestampMins, rolledHosts, conf, savedStartCode);
+ logList = excludeProcV2WALs(selection.getIncluded());
backupInfo.setIncrBackupFileList(logList);
- // Update boundaries based on WALs that will be backed up
- for (String logFile : logList) {
- Path logPath = new Path(logFile);
- String logHost = BackupUtils.parseHostFromOldLog(logPath);
- if (logHost == null) {
- logHost = BackupUtils.parseHostNameFromLogFile(logPath.getParent());
- }
- if (logHost != null) {
- long logTs = BackupUtils.getCreationTime(logPath);
- Long latestTimestampToIncludeInBackup = newTimestamps.get(logHost);
- if (latestTimestampToIncludeInBackup == null || logTs > latestTimestampToIncludeInBackup) {
- LOG.info("Updating backup boundary for inactive host {}: timestamp={}", logHost, logTs);
- newTimestamps.put(logHost, logTs);
- }
- }
- }
+ Map newTimestamps = BackupUtils.computeLogBoundaries(rolledHosts,
+ previousTimestampMins, selection.getHostsWithLogs(), logList, selection.getHeldBack());
+ LOG.debug("Log boundaries for incremental backup {}: {}", backupInfo.getBackupId(),
+ newTimestamps);
return newTimestamps;
}
+ private static final class LogFileSelection {
+ private final List included;
+ private final List heldBack;
+ private final Set hostsWithLogs;
+
+ private LogFileSelection(List included, List heldBack,
+ Set hostsWithLogs) {
+ this.included = included;
+ this.heldBack = heldBack;
+ this.hostsWithLogs = hostsWithLogs;
+ }
+
+ private List getIncluded() {
+ return included;
+ }
+
+ private List getHeldBack() {
+ return heldBack;
+ }
+
+ private Set getHostsWithLogs() {
+ return hostsWithLogs;
+ }
+ }
+
private List excludeProcV2WALs(List logList) {
List list = new ArrayList<>();
for (int i = 0; i < logList.size(); i++) {
@@ -155,16 +156,18 @@ 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
* @param savedStartCode the startcode (timestamp) of last successful backup.
- * @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, String savedStartCode)
throws IOException {
LOG.debug("In getLogFilesForNewBackup()\n" + "olderTimestamps: " + olderTimestamps
@@ -178,6 +181,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
@@ -202,6 +206,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
@@ -243,10 +248,10 @@ private List getLogFilesForNewBackup(Map olderTimestamps,
}
// Include the .oldlogs files too.
- FileStatus[] oldlogs = fs.listStatus(oldLogDir);
- for (FileStatus oldlog : oldlogs) {
- p = oldlog.getPath();
- currentLogFile = p.toString();
+ List oldlogs = BackupUtils.getFiles(fs, oldLogDir, new ArrayList<>(), path -> true);
+ for (String oldlog : oldlogs) {
+ p = new Path(oldlog);
+ currentLogFile = oldlog;
if (AbstractFSWALProvider.isMetaFile(p)) {
if (LOG.isDebugEnabled()) {
LOG.debug("Skip .meta log file: " + currentLogFile);
@@ -257,6 +262,7 @@ private List getLogFilesForNewBackup(Map olderTimestamps,
if (host == null) {
continue;
}
+ hostsWithLogs.add(host);
currentLogTS = BackupUtils.getCreationTime(p);
oldTimeStamp = olderTimestamps.get(host);
/*
@@ -275,10 +281,15 @@ private List getLogFilesForNewBackup(Map olderTimestamps,
} else if (currentLogTS > oldTimeStamp) {
resultLogFiles.add(currentLogFile);
}
+
+ Long newTimestamp = newestTimestamps.get(host);
+ 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 ebc8ad13be29..e2284cd87ce4 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
@@ -21,12 +21,14 @@
import java.io.IOException;
import java.net.URLDecoder;
import java.util.ArrayList;
+import java.util.Collection;
import java.util.Collections;
import java.util.Comparator;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Map.Entry;
+import java.util.Set;
import java.util.TreeSet;
import java.util.function.Predicate;
import java.util.stream.Collectors;
@@ -193,6 +195,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
index 10e33380f111..bc5ffe1118a0 100644
--- 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
@@ -17,21 +17,40 @@
*/
package org.apache.hadoop.hbase.backup;
-import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
+import static org.junit.Assert.assertNull;
import static org.junit.Assert.assertTrue;
+import java.io.IOException;
import java.util.List;
import java.util.Map;
import org.apache.hadoop.hbase.HBaseClassTestRule;
import org.apache.hadoop.hbase.HBaseTestingUtility;
import org.apache.hadoop.hbase.MiniHBaseCluster;
+import org.apache.hadoop.hbase.ServerName;
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.BeforeClass;
import org.junit.ClassRule;
import org.junit.Test;
@@ -63,45 +82,92 @@ public static void setUp() throws Exception {
TEST_UTIL = new HBaseTestingUtility();
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();
- // Start an additional RS so we have at least 2
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 Exception {
+ afterSnapshotHook.run();
+ super.snapshotCopy(backupInfo);
+ }
+ }
+
+ private static void stopRegionServerAndWait(ServerName serverName) throws Exception {
+ MiniHBaseCluster 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(message, TEST_UTIL.countRows(table1), TEST_UTIL.countRows(restoredTable));
+ }
+
+ 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 re-include those
+ * offline host's timestamps are recorded so subsequent incremental backups don't reinclude those
* WALs.
*/
@Test
public void testBackupWithOfflineRS() throws Exception {
- LOG.info("Starting testFullBackupWithOfflineRS");
+ LOG.info("Starting testBackupWithOfflineRS");
MiniHBaseCluster cluster = TEST_UTIL.getMiniHBaseCluster();
List tables = Lists.newArrayList(table1);
- if (cluster.getNumLiveRegionServers() < 2) {
- cluster.startRegionServer();
- Thread.sleep(2000);
- }
+ 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);
}
- int rsToStop = 0;
- HRegionServer rsBeforeStop = cluster.getRegionServer(rsToStop);
String offlineHost =
rsBeforeStop.getServerName().getHostname() + ":" + rsBeforeStop.getServerName().getPort();
LOG.info("Stopping RS: {}", offlineHost);
- cluster.stopRegionServer(rsToStop);
- // Wait for WALs to be moved to oldlogs
- Thread.sleep(5000);
+ stopRegionServerAndWait(rsBeforeStop.getServerName());
LOG.info("Taking full backup (with offline RS WALs in oldlogs)");
String fullBackupId = fullTableBackup(tables);
@@ -120,10 +186,258 @@ public void testBackupWithOfflineRS() throws Exception {
String incrBackupId = incrementalTableBackup(tables);
assertTrue("Incremental backup should succeed", checkSucceeded(incrBackupId));
- timestamps = sysTable.readLogTimestampMap(BACKUP_ROOT_DIR);
- rsTimestamps = timestamps.get(table1);
- assertFalse("Offline host should not have a boundary ",
- rsTimestamps.containsKey(offlineHost));
+ assertEquals("Incremental backup moved the boundary of the offline host", tsAfterFullBackup,
+ boundary(sysTable, offlineHost));
+ }
+ }
+
+ /**
+ * 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 {
+ MiniHBaseCluster 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("Full backup should succeed", checkSucceeded(fullBackupId));
+
+ 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("Incremental backup should succeed", checkSucceeded(incrBackupId));
+
+ 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(
+ "Offline RS should have a timestamp boundary after the incremental backed up its WALs",
+ boundaryAfterIncr1);
+ Long staleLogRoll =
+ sysTable.readRegionServerLastLogRollResult(BACKUP_ROOT_DIR).get(offlineHost);
+ String message = "Incremental backup moved the boundary of the offline RS, which would back "
+ + "up its WALs again (last log roll = " + staleLogRoll + ")";
+
+ String incrBackupId2 = incrementalTableBackup(tables);
+ assertTrue("Second incremental backup should succeed", checkSucceeded(incrBackupId2));
+ assertEquals(message, boundaryAfterIncr1, boundary(sysTable, offlineHost));
+
+ String incrBackupId3 = incrementalTableBackup(tables);
+ assertTrue("Third incremental backup should succeed", checkSucceeded(incrBackupId3));
+ assertEquals(message, boundaryAfterIncr1, boundary(sysTable, offlineHost));
+
+ String incrBackupId4 = incrementalTableBackup(tables);
+ assertTrue("Fourth incremental backup should succeed", checkSucceeded(incrBackupId4));
+ assertEquals(message, boundaryAfterIncr1, boundary(sysTable, offlineHost));
+ }
+ }
+
+ /**
+ * 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 {
+ MiniHBaseCluster 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("Full backup should succeed", checkSucceeded(fullBackupId));
+
+ try (BackupSystemTable sysTable = new BackupSystemTable(TEST_UTIL.getConnection())) {
+ Long boundaryAfterFullBackup = boundary(sysTable, transientHost);
+ assertNotNull(
+ "Transient RS that went offline before full backup should have its WAL boundary recorded",
+ boundaryAfterFullBackup);
+
+ String incrBackupId = incrementalTableBackup(tables);
+ assertTrue("Incremental backup after transient RS should succeed",
+ checkSucceeded(incrBackupId));
+
+ assertEquals("Incremental backup moved the boundary of the transient RS",
+ boundaryAfterFullBackup, boundary(sysTable, transientHost));
+ }
+ }
+
+ /**
+ * 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 {
+ MiniHBaseCluster 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("Full backup should succeed", checkSucceeded(fullBackupId));
+
+ 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("Incremental backup should succeed", checkSucceeded(incrBackupId));
+
+ 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(
+ "Transient RS should have a timestamp boundary after the incremental backed up its WALs",
+ boundaryAfterIncr1);
+
+ String incrBackupId2 = incrementalTableBackup(tables);
+ assertTrue("Second incremental backup should succeed", checkSucceeded(incrBackupId2));
+
+ assertEquals("Second incremental backup moved the boundary of the transient RS",
+ boundaryAfterIncr1, boundary(sysTable, transientHost));
+ }
+ }
+
+ /**
+ * 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 {
+ MiniHBaseCluster 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("Full backup should succeed", checkSucceeded(fullBackupId));
+
+ String incrBackupId = incrementalTableBackup(tables);
+ assertTrue("Incremental backup should succeed", checkSucceeded(incrBackupId));
+
+ 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 {
+ MiniHBaseCluster 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("First full backup should succeed", checkSucceeded(firstFullBackupId));
+
+ 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("Deleted row should be gone from the source table",
+ table.get(new Get(deletedRow)).isEmpty());
+ Scan rawScan = new Scan().withStartRow(deletedRow).withStopRow(deletedRow, true).setRaw(true);
+ try (ResultScanner scanner = table.getScanner(rawScan)) {
+ assertNull("Major compaction should have purged the deleted row and its delete marker",
+ scanner.next());
+ }
+
+ stopRegionServerAndWait(offlineRS.getServerName());
+
+ String secondFullBackupId = fullTableBackup(tables);
+ assertTrue("Second full backup should succeed", checkSucceeded(secondFullBackupId));
+ String incrBackupId = incrementalTableBackup(tables);
+ assertTrue("Incremental backup should succeed", checkSucceeded(incrBackupId));
+
+ TableName restoredTable = TableName.valueOf("table1_deleted_row_not_resurrected");
+ restore(incrBackupId, table1, restoredTable);
+ try (Table restored = TEST_UTIL.getConnection().getTable(restoredTable)) {
+ assertTrue("Deleted row was brought back by the WALs of the offline RS",
+ restored.get(new Get(deletedRow)).isEmpty());
+ }
+ assertEquals("Restored table should match the source table", TEST_UTIL.countRows(table1),
+ TEST_UTIL.countRows(restoredTable));
}
}
}
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 d4e6900fb28b..a7a85450f666 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
@@ -19,6 +19,8 @@
import java.io.IOException;
import java.security.PrivilegedAction;
+import java.util.List;
+import java.util.Map;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;
@@ -40,6 +42,10 @@
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import org.apache.hbase.thirdparty.com.google.common.collect.ImmutableList;
+import org.apache.hbase.thirdparty.com.google.common.collect.ImmutableMap;
+import org.apache.hbase.thirdparty.com.google.common.collect.ImmutableSet;
+
@Category(SmallTests.class)
public class TestBackupUtils {
@ClassRule
@@ -115,4 +121,90 @@ public void testFilesystemWalHostNameParsing() throws IOException {
}
}
+
+ @Test
+ public void testGetRolledHostsKeepsOnlyHostsWhoseRollResultChanged() {
+ Map previousLogRolls =
+ ImmutableMap.of("rolled:16020", 100L, "offline:16020", 200L, "removed:16020", 300L);
+ Map latestLogRolls =
+ ImmutableMap.of("rolled:16020", 150L, "offline:16020", 200L, "new:16020", 400L);
+
+ Assert.assertEquals(ImmutableMap.of("rolled:16020", 150L, "new:16020", 400L),
+ BackupUtils.getRolledHosts(previousLogRolls, latestLogRolls));
+ }
+
+ @Test
+ public void testComputeLogBoundariesUsesRollResultForRolledHosts() throws IOException {
+ Map rolledHosts = ImmutableMap.of("rolled:16020", 100L);
+ List coveredLogs = ImmutableList.of("/hbase/oldWALs/rolled%2C16020%2C1.500");
+ List pendingLogs = ImmutableList.of("/hbase/WALs/rolled,16020,1/rolled%2C16020%2C1.50");
+
+ Assert.assertEquals(ImmutableMap.of("rolled:16020", 100L),
+ BackupUtils.computeLogBoundaries(rolledHosts, coveredLogs, pendingLogs));
+ }
+
+ @Test
+ public void testComputeLogBoundariesUsesNewestCoveredLogForOtherHosts() throws IOException {
+ List coveredLogs = ImmutableList.of("/hbase/oldWALs/offline%2C16020%2C1.200",
+ "/hbase/oldWALs/offline%2C16020%2C1.300",
+ "/hbase/WALs/offline,16020,1/offline%2C16020%2C1.250");
+
+ Assert.assertEquals(ImmutableMap.of("offline:16020", 300L),
+ BackupUtils.computeLogBoundaries(ImmutableMap.of(), coveredLogs, ImmutableList.of()));
+ }
+
+ @Test
+ public void testComputeLogBoundariesCapsBelowOldestPendingLog() throws IOException {
+ List coveredLogs = ImmutableList.of("/hbase/oldWALs/splitting%2C16020%2C1.400");
+ List pendingLogs =
+ ImmutableList.of("/hbase/WALs/splitting,16020,1-splitting/splitting%2C16020%2C1.350",
+ "/hbase/WALs/splitting,16020,1-splitting/splitting%2C16020%2C1.370");
+
+ Assert.assertEquals(ImmutableMap.of("splitting:16020", 349L),
+ BackupUtils.computeLogBoundaries(ImmutableMap.of(), coveredLogs, pendingLogs));
+ }
+
+ @Test
+ public void testComputeLogBoundariesKeepsHostWithOnlyPendingLogs() throws IOException {
+ List pendingLogs =
+ ImmutableList.of("/hbase/WALs/joined,16020,1/joined%2C16020%2C1.500");
+
+ Assert.assertEquals(ImmutableMap.of("joined:16020", 499L),
+ BackupUtils.computeLogBoundaries(ImmutableMap.of(), ImmutableList.of(), pendingLogs));
+ }
+
+ @Test
+ public void testComputeLogBoundariesSkipsUnparseableLogs() throws IOException {
+ List coveredLogs = ImmutableList.of("/hbase/oldWALs/not-a-wal");
+
+ Map rolledHosts = ImmutableMap.of("rolled:16020", 100L);
+
+ Assert.assertEquals(rolledHosts,
+ BackupUtils.computeLogBoundaries(rolledHosts, coveredLogs, ImmutableList.of()));
+ }
+
+ @Test
+ public void testComputeLogBoundariesKeepsPreviousBoundaryWhileHostHasLogs() throws IOException {
+ Map previousBoundaries =
+ ImmutableMap.of("stuck:16020", 1200L, "gone:16020", 800L);
+
+ Assert.assertEquals(ImmutableMap.of("stuck:16020", 1200L),
+ BackupUtils.computeLogBoundaries(ImmutableMap.of(), previousBoundaries,
+ ImmutableSet.of("stuck:16020"), ImmutableList.of(), ImmutableList.of()));
+ }
+
+ @Test
+ public void testComputeLogBoundariesPrefersNewBoundaryOverPreviousOne() throws IOException {
+ Map previousBoundaries =
+ ImmutableMap.of("rolled:16020", 100L, "offline:16020", 200L, "pending:16020", 300L);
+ List coveredLogs = ImmutableList.of("/hbase/oldWALs/offline%2C16020%2C1.250");
+ List pendingLogs =
+ ImmutableList.of("/hbase/WALs/pending,16020,1/pending%2C16020%2C1.400");
+
+ Assert.assertEquals(
+ ImmutableMap.of("rolled:16020", 150L, "offline:16020", 250L, "pending:16020", 399L),
+ BackupUtils.computeLogBoundaries(ImmutableMap.of("rolled:16020", 150L), previousBoundaries,
+ ImmutableSet.of("rolled:16020", "offline:16020", "pending:16020"), coveredLogs,
+ pendingLogs));
+ }
}
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
new file mode 100644
index 000000000000..b6b865eb8a90
--- /dev/null
+++ b/hbase-backup/src/test/java/org/apache/hadoop/hbase/backup/TestIncrementalBackupManager.java
@@ -0,0 +1,389 @@
+/*
+ * 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.Assert.assertFalse;
+import static org.junit.Assert.assertNotNull;
+import static org.junit.Assert.assertTrue;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.Map;
+import org.apache.hadoop.conf.Configuration;
+import org.apache.hadoop.fs.FileSystem;
+import org.apache.hadoop.fs.Path;
+import org.apache.hadoop.hbase.HBaseClassTestRule;
+import org.apache.hadoop.hbase.HBaseTestingUtility;
+import org.apache.hadoop.hbase.ServerName;
+import org.apache.hadoop.hbase.TableName;
+import org.apache.hadoop.hbase.backup.impl.BackupAdminImpl;
+import org.apache.hadoop.hbase.backup.impl.IncrementalBackupManager;
+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.BeforeClass;
+import org.junit.ClassRule;
+import org.junit.Test;
+import org.junit.experimental.categories.Category;
+
+@Category(LargeTests.class)
+public class TestIncrementalBackupManager extends TestBackupBase {
+
+ @ClassRule
+ public static final HBaseClassTestRule CLASS_RULE =
+ HBaseClassTestRule.forClass(TestIncrementalBackupManager.class);
+
+ @BeforeClass
+ public static void setUp() throws Exception {
+ TEST_UTIL = new HBaseTestingUtility();
+ conf1 = TEST_UTIL.getConfiguration();
+ autoRestoreOnFailure = true;
+ useSecondCluster = false;
+ setUpHelper();
+ }
+
+ @Test
+ public void testCollectWALFilesFromRegionServerDirectories() throws Exception {
+ testCollectWALFiles(true);
+ }
+
+ @Test
+ public void testCollectWALFilesFromFlatOldWALDirectory() throws Exception {
+ testCollectWALFiles(false);
+ }
+
+ private void testCollectWALFiles(boolean separateOldLogDir) throws Exception {
+ Configuration testConf = new Configuration(conf1);
+ testConf.setBoolean(AbstractFSWALProvider.SEPARATE_OLDLOGDIR, separateOldLogDir);
+ List tables = Collections.singletonList(table1);
+ Path walRootDir = CommonFSUtils.getWALRootDir(conf1);
+ FileSystem fs = walRootDir.getFileSystem(conf1);
+
+ try (Connection conn = ConnectionFactory.createConnection(testConf);
+ BackupAdminImpl backupAdmin = new BackupAdminImpl(conn)) {
+ String fullBackupId = backupAdmin
+ .backupTables(createBackupRequest(BackupType.FULL, tables, BACKUP_ROOT_DIR)).getBackupId();
+ assertTrue(checkSucceeded(fullBackupId));
+
+ try (IncrementalBackupManager manager = new IncrementalBackupManager(conn, testConf)) {
+ BackupInfo backupInfo = manager.createBackupInfo("backup_test", BackupType.INCREMENTAL,
+ tables, BACKUP_ROOT_DIR, -1, -1, false);
+ Map previousTimestamps =
+ BackupUtils.getRSLogTimestampMins(manager.readLogTimestampMap());
+ ServerName serverName = TEST_UTIL.getMiniHBaseCluster().getRegionServer(0).getServerName();
+ Long previousTimestamp = previousTimestamps.get(serverName.getAddress().toString());
+ assertNotNull(previousTimestamp);
+
+ TEST_UTIL.waitFor(30_000,
+ () -> EnvironmentEdgeManager.currentTime() > previousTimestamp + 1);
+ Path archiveDir = new Path(walRootDir,
+ AbstractFSWALProvider.getWALArchiveDirectoryName(testConf, serverName.toString()));
+ Path archivedWAL = new Path(archiveDir, walName(serverName, previousTimestamp + 1));
+ fs.mkdirs(archiveDir);
+ fs.create(archivedWAL).close();
+
+ try {
+ manager.getIncrBackupLogFileMap();
+
+ assertTrue("Archived WAL should be in the backup: " + backupInfo.getIncrBackupFileList(),
+ backupInfo.getIncrBackupFileList().contains(archivedWAL.toString()));
+ } finally {
+ fs.delete(archivedWAL, false);
+ }
+ }
+ }
+ }
+
+ /**
+ * 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 = Collections.singletonList(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)).getBackupId();
+ 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, serverName.toString() + 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);
+ Map boundaries = manager.getIncrBackupLogFileMap();
+ manager.writeRegionServerLogTimestamp(backupInfo.getTables(), boundaries);
+ firstBackupFiles = backupInfo.getIncrBackupFileList();
+ }
+ assertFalse("Archived WAL newer than the roll result should be deferred to a later backup: "
+ + firstBackupFiles, firstBackupFiles.contains(newerArchivedWAL.toString()));
+
+ 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);
+ manager.getIncrBackupLogFileMap();
+ secondBackupFiles = backupInfo.getIncrBackupFileList();
+ }
+
+ assertTrue(
+ "WAL " + olderWAL + " was not included in any backup. First backup: " + firstBackupFiles
+ + ", second backup: " + secondBackupFiles,
+ firstBackupFiles.contains(olderWAL.toString())
+ || secondBackupFiles.contains(olderWAL.toString()));
+ } 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 = Collections.singletonList(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)).getBackupId();
+ assertTrue(checkSucceeded(fullBackupId));
+
+ long walTs = EnvironmentEdgeManager.currentTime() + 1;
+ Path splittingWAL = new Path(splittingDir, walName(deadServer, walTs));
+ Path archivedWAL = new Path(archiveDir, walName(deadServer, walTs));
+ fs.create(splittingWAL).close();
+
+ try {
+ 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(
+ "WAL of the dead server was not included in any backup. First backup: " + firstBackupFiles
+ + ", second backup: " + secondBackupFiles,
+ firstBackupFiles.contains(splittingWAL.toString())
+ || secondBackupFiles.contains(archivedWAL.toString()));
+ } 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 = Collections.singletonList(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)).getBackupId();
+ 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 = new Path(archiveDir, walName(deadServer, oldestTs));
+ Path middleArchivedWAL = new Path(archiveDir, walName(deadServer, middleTs));
+ Path newestArchivedWAL = new Path(archiveDir, walName(deadServer, 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");
+
+ assertTrue(
+ "Oldest WAL of the dead server was not included in any backup. First backup: "
+ + firstBackupFiles + ", second backup: " + secondBackupFiles,
+ firstBackupFiles.contains(oldestSplittingWAL.toString())
+ || secondBackupFiles.contains(oldestArchivedWAL.toString()));
+ assertTrue(
+ "Middle WAL of the dead server was not included in any backup. First backup: "
+ + firstBackupFiles + ", second backup: " + secondBackupFiles,
+ firstBackupFiles.contains(middleSplittingWAL.toString())
+ || secondBackupFiles.contains(middleArchivedWAL.toString()));
+ } 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 = Collections.singletonList(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)).getBackupId();
+ assertTrue(checkSucceeded(fullBackupId));
+
+ long stuckTs = EnvironmentEdgeManager.currentTime() + 1;
+ long archivedTs = stuckTs + 1;
+ Path archivedWAL = new Path(archiveDir, walName(deadServer, archivedTs));
+ Path stuckSplittingWAL = new Path(splittingDir, walName(deadServer, stuckTs));
+ Path stuckArchivedWAL = new Path(archiveDir, walName(deadServer, 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(
+ "Archived WAL of the dead server should be in the first backup: " + firstBackupFiles,
+ firstBackupFiles.contains(archivedWAL.toString()));
+
+ 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(
+ "WAL already covered by the dead server's boundary was backed up again: "
+ + laterBackupFiles,
+ laterBackupFiles.contains(stuckSplittingWAL.toString())
+ || laterBackupFiles.contains(stuckArchivedWAL.toString()));
+ } 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 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);
+ Map boundaries = manager.getIncrBackupLogFileMap();
+ manager.writeRegionServerLogTimestamp(backupInfo.getTables(), boundaries);
+ manager.writeBackupStartCode(
+ BackupUtils.getMinValue(BackupUtils.getRSLogTimestampMins(manager.readLogTimestampMap())));
+ 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..f04dac028f44
--- /dev/null
+++ b/hbase-backup/src/test/java/org/apache/hadoop/hbase/backup/impl/TestFullTableBackupClientLogBoundaries.java
@@ -0,0 +1,158 @@
+/*
+ * 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.Assert.assertEquals;
+import static org.junit.Assert.assertThrows;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+import java.io.IOException;
+import java.util.Collection;
+import java.util.Map;
+import org.apache.hadoop.fs.FileSystem;
+import org.apache.hadoop.fs.Path;
+import org.apache.hadoop.hbase.HBaseClassTestRule;
+import org.apache.hadoop.hbase.HBaseTestingUtility;
+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.After;
+import org.junit.Before;
+import org.junit.ClassRule;
+import org.junit.Test;
+import org.junit.experimental.categories.Category;
+
+import org.apache.hbase.thirdparty.com.google.common.collect.ImmutableList;
+import org.apache.hbase.thirdparty.com.google.common.collect.ImmutableMap;
+
+@Category(SmallTests.class)
+public class TestFullTableBackupClientLogBoundaries {
+
+ @ClassRule
+ public static final HBaseClassTestRule CLASS_RULE =
+ HBaseClassTestRule.forClass(TestFullTableBackupClientLogBoundaries.class);
+
+ private static final HBaseTestingUtility TEST_UTIL = new HBaseTestingUtility();
+
+ private FileSystem fs;
+ private Path walRootDir;
+
+ @Before
+ public void setUp() throws IOException {
+ fs = TEST_UTIL.getTestFileSystem();
+ walRootDir = TEST_UTIL.getDataTestDirOnTestFS("walRoot");
+ }
+
+ @After
+ 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(ImmutableMap.of("joined:16020", 499L),
+ computeLogBoundaries(ImmutableMap.of(), ImmutableList.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(ImmutableMap.of("dead:16020", 350L),
+ computeLogBoundaries(ImmutableMap.of(), ImmutableList.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(ImmutableMap.of("restarted:16020", 700L),
+ computeLogBoundaries(ImmutableMap.of(), ImmutableList.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(ImmutableMap.of("rolled:16020", 850L),
+ computeLogBoundaries(ImmutableMap.of("rolled:16020", 850L), ImmutableList.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 ImmutableList.of(joined);
+ });
+
+ assertEquals(ImmutableMap.of("joined:16020", 499L),
+ FullTableBackupClient.computeLogBoundaries(fs, walRootDir, ImmutableMap.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(ImmutableMap.of("rolled:16020", 850L), ImmutableList.of(joined)));
+ }
+
+ private Map computeLogBoundaries(Map rolledHosts,
+ Collection 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;
+ }
+}