Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions hbase-backup/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -167,6 +167,11 @@
<artifactId>junit</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.mockito</groupId>
<artifactId>mockito-core</artifactId>
<scope>test</scope>
</dependency>
</dependencies>
<build>
<plugins>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -679,8 +679,9 @@ public List<BackupInfo> 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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -174,6 +177,9 @@ public void execute() throws IOException {
LogRollMasterProcedureManager.ROLLLOG_PROCEDURE_NAME, props);

Map<String, Long> 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);
Expand All @@ -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<FileStatus> 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);
Expand All @@ -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<String, Long> computeLogBoundaries(FileSystem fs, Path walRootDir,
Map<String, Long> rolledHosts, Admin admin) throws IOException {
Path logDir = new Path(walRootDir, HConstants.HREGION_LOGDIR_NAME);
Path oldLogDir = new Path(walRootDir, HConstants.HREGION_OLDLOGDIR_NAME);

Map<ServerName, List<String>> logsByServer = new HashMap<>();
for (FileStatus serverLogDir : fs.listStatus(logDir)) {
ServerName serverName =
AbstractFSWALProvider.getServerNameFromWALDirectoryName(serverLogDir.getPath());
if (serverName == null) {
continue;
}
List<String> 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<ServerName> live = new HashSet<>(admin.getRegionServers());
Set<String> 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<String> coveredLogs = new ArrayList<>();
List<String> pendingLogs = new ArrayList<>();
for (Map.Entry<ServerName, List<String>> 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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<String, Long> getIncrBackupLogFileMap() throws IOException {
Expand All @@ -82,6 +87,7 @@ public Map<String, Long> getIncrBackupLogFileMap() throws IOException {
+ "In order to create an incremental backup, at least one full backup is needed.");
}

Map<String, Long> previousLogRollByHost = readRegionServerLastLogRollResult();
if (backupInfo.getUsePreviousLogRoll()) {
LOG.info("Using previous WAL roll for backup, skipping WAL roll procedure");
} else {
Expand All @@ -94,51 +100,46 @@ public Map<String, Long> getIncrBackupLogFileMap() throws IOException {
LogRollMasterProcedureManager.ROLLLOG_PROCEDURE_NAME, props);
}
}
Map<String, Long> newTimestamps = readRegionServerLastLogRollResult();

Map<String, Long> latestLogRollByHost = readRegionServerLastLogRollResult();
for (Map.Entry<String, Long> 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<String, Long> 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<String, Long> 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<String> included;
private final List<String> heldBack;
private final Set<String> hostsWithLogs;

private LogFileSelection(List<String> included, List<String> heldBack,
Set<String> hostsWithLogs) {
this.included = included;
this.heldBack = heldBack;
this.hostsWithLogs = hostsWithLogs;
}

private List<String> getIncluded() {
return included;
}

private List<String> getHeldBack() {
return heldBack;
}

private Set<String> getHostsWithLogs() {
return hostsWithLogs;
}
}

private List<String> excludeProcV2WALs(List<String> logList) {
List<String> list = new ArrayList<>();
for (int i = 0; i < logList.size(); i++) {
Expand All @@ -155,16 +156,18 @@ private List<String> excludeProcV2WALs(List<String> 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<String> getLogFilesForNewBackup(Map<String, Long> olderTimestamps,
private LogFileSelection getLogFilesForNewBackup(Map<String, Long> olderTimestamps,
Map<String, Long> newestTimestamps, Configuration conf, String savedStartCode)
throws IOException {
LOG.debug("In getLogFilesForNewBackup()\n" + "olderTimestamps: " + olderTimestamps
Expand All @@ -178,6 +181,7 @@ private List<String> getLogFilesForNewBackup(Map<String, Long> olderTimestamps,

List<String> resultLogFiles = new ArrayList<>();
List<String> newestLogs = new ArrayList<>();
Set<String> hostsWithLogs = new HashSet<>();

/*
* The old region servers and timestamps info we kept in backup system table may be out of sync
Expand All @@ -202,6 +206,7 @@ private List<String> getLogFilesForNewBackup(Map<String, Long> 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
Expand Down Expand Up @@ -243,10 +248,10 @@ private List<String> getLogFilesForNewBackup(Map<String, Long> olderTimestamps,
}

// Include the .oldlogs files too.
FileStatus[] oldlogs = fs.listStatus(oldLogDir);
for (FileStatus oldlog : oldlogs) {
p = oldlog.getPath();
currentLogFile = p.toString();
List<String> 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);
Expand All @@ -257,6 +262,7 @@ private List<String> getLogFilesForNewBackup(Map<String, Long> olderTimestamps,
if (host == null) {
continue;
}
hostsWithLogs.add(host);
currentLogTS = BackupUtils.getCreationTime(p);
oldTimeStamp = olderTimestamps.get(host);
/*
Expand All @@ -275,10 +281,15 @@ private List<String> getLogFilesForNewBackup(Map<String, Long> 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 {
Expand Down
Loading