From e5ec7d4392941a668a5b8af30eadc3caa1e63987 Mon Sep 17 00:00:00 2001 From: Hernan Gelaf-Romer Date: Fri, 2 Oct 2026 23:32:15 -0400 Subject: [PATCH] Log filtering in IncrementalBackupManager can lead to data loss --- hbase-backup/pom.xml | 5 + .../hbase/backup/impl/BackupSystemTable.java | 5 +- .../backup/impl/FullTableBackupClient.java | 120 ++++-- .../backup/impl/IncrementalBackupManager.java | 109 ++--- .../hadoop/hbase/backup/util/BackupUtils.java | 80 ++++ .../hbase/backup/TestBackupOfflineRS.java | 348 +++++++++++++++- .../hadoop/hbase/backup/TestBackupUtils.java | 92 +++++ .../backup/TestIncrementalBackupManager.java | 389 ++++++++++++++++++ ...estFullTableBackupClientLogBoundaries.java | 158 +++++++ 9 files changed, 1194 insertions(+), 112 deletions(-) create mode 100644 hbase-backup/src/test/java/org/apache/hadoop/hbase/backup/TestIncrementalBackupManager.java create mode 100644 hbase-backup/src/test/java/org/apache/hadoop/hbase/backup/impl/TestFullTableBackupClientLogBoundaries.java 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; + } +}