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
Original file line number Diff line number Diff line change
Expand Up @@ -662,8 +662,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 @@ -36,13 +36,21 @@
import static org.apache.hadoop.hbase.replication.ReplicationUtils.OFFSET_UPDATE_INTERVAL_MS_KEY;
import static org.apache.hadoop.hbase.replication.ReplicationUtils.OFFSET_UPDATE_SIZE_THRESHOLD_KEY;

import com.google.errorprone.annotations.RestrictedApi;
import java.io.IOException;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.UUID;
import java.util.stream.Collectors;
import org.apache.hadoop.fs.FileStatus;
import org.apache.hadoop.fs.FileSystem;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.hbase.HConstants;
import org.apache.hadoop.hbase.ServerName;
import org.apache.hadoop.hbase.TableName;
import org.apache.hadoop.hbase.backup.BackupCopyJob;
import org.apache.hadoop.hbase.backup.BackupInfo;
Expand All @@ -60,7 +68,9 @@
import org.apache.hadoop.hbase.client.TableDescriptorBuilder;
import org.apache.hadoop.hbase.replication.ReplicationException;
import org.apache.hadoop.hbase.replication.ReplicationPeerConfig;
import org.apache.hadoop.hbase.util.CommonFSUtils;
import org.apache.hadoop.hbase.util.EnvironmentEdgeManager;
import org.apache.hadoop.hbase.wal.AbstractFSWALProvider;
import org.apache.yetus.audience.InterfaceAudience;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
Expand Down Expand Up @@ -208,7 +218,12 @@ private void handleContinuousBackup(Admin admin) throws IOException {
}

private void handleNonContinuousBackup(Admin admin) throws IOException {
performLogRoll();
Map<String, Long> previousLogRollsByHost = backupManager.readRegionServerLastLogRollResult();
Map<String, Long> latestLogRollsByHost = performLogRoll();
Path walRootDir = CommonFSUtils.getWALRootDir(conf);
newTimestamps = computeLogBoundaries(walRootDir.getFileSystem(conf), walRootDir,
BackupUtils.getRolledHosts(previousLogRollsByHost, latestLogRollsByHost), admin);

performBackupSnapshots(admin);
backupManager.addIncrementalBackupTableSet(backupInfo.getTables());

Expand All @@ -219,15 +234,15 @@ private void handleNonContinuousBackup(Admin admin) throws IOException {
updateBackupMetadata();
}

private void performLogRoll() throws IOException {
private Map<String, Long> performLogRoll() throws IOException {
// We roll log here before we do the snapshot. It is possible there is duplicate data
// in the log that is already in the snapshot. But if we do it after the snapshot, we
// could have data loss.
// A better approach is to do the roll log on each RS in the same global procedure as
// the snapshot.
LOG.info("Execute roll log procedure for full backup ...");
BackupUtils.logRoll(conn, backupInfo.getBackupRootDir(), conf);
newTimestamps = backupManager.readRegionServerLastLogRollResult();
return backupManager.readRegionServerLastLogRollResult();
}

private void performBackupSnapshots(Admin admin) throws IOException {
Expand All @@ -248,6 +263,75 @@ private void performSnapshots(Admin admin) throws IOException {
}
}

/**
* Computes the per-host log boundaries stored by this full backup, using
* {@link BackupUtils#computeLogBoundaries}. A host that took part in this backup's log roll gets
* its roll result: WALs created after the roll are included in the next incremental backup, even
* though some of their edits may already be in the snapshot, which is safe because deletes are
* replayed along with the puts. For any other host, its archived WALs and the WALs of dead
* servers are covered by the snapshot, so the next incremental backup does not replay them. The
* WALs of a live server that did not take part in the roll (for example one that started during
* it) are pending, because they can still receive edits that are not in the snapshot. The live
* servers are read only after the WAL directories are listed: a region server creates its WAL
* directory only after the master registers it, so every live server whose directory was listed
* is found. If the live servers do not include every host that took part in the roll, the
* master's server list is incomplete (for example right after a master failover), and the backup
* fails instead of treating live servers as dead. It must run right after the log roll, before
* the snapshot is taken, so that a host starting later gets no boundary and has all of its WALs
* included in the next incremental backup.
*/
@RestrictedApi(
explanation = "Package-private for test visibility only. Do not use outside tests.",
link = "",
allowedOnPath = "(.*/src/test/.*|.*/org/apache/hadoop/hbase/backup/impl/FullTableBackupClient.java)")
static Map<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())) {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sadly, I can't think of a much better way to accomplish consistency here. The HMaster may be initializing when we query for live RS, and unfortunately that means we may get a partial list of RS back.

This sanity check is a little brittle, but I believe that it's robust enough to do what we need. Given that we just rolled WAL files, if we don't see any of those hosts in the live RS list we got back from the HMaster, we can likely assume the HMaster didn't return a comprehensive list, and we should abort.

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);
}

private void updateBackupMetadata() throws IOException {
// The table list in backupInfo is good for both full backup and incremental backup.
// For incremental backup, it contains the incremental backup table set.
Expand Down
Comment thread
hgromer marked this conversation as resolved.
Comment thread
hgromer marked this conversation as resolved.
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.Collections;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.FileStatus;
import org.apache.hadoop.fs.FileSystem;
Expand Down Expand Up @@ -54,12 +56,14 @@ public IncrementalBackupManager(Connection conn, Configuration conf) throws IOEx
/**
* Obtain the list of logs that need to be copied out for this incremental backup. The list is set
* in BackupInfo.
* @return The new HashMap of RS log time stamps after the log roll for this incremental backup.
* @return The new map of RS log time stamps for this incremental backup, as computed by
* {@link BackupUtils#computeLogBoundaries}: the included logs are covered, the logs held
* back for a later backup (for example, logs still being split) are pending, and a host
* that still has logs keeps its previous boundary if this backup gives it no new one.
* @throws IOException exception
*/
public Map<String, Long> getIncrBackupLogFileMap() throws IOException {
List<String> logList;
Map<String, Long> newTimestamps;
Map<String, Long> previousTimestampMins =
BackupUtils.getRSLogTimestampMins(readLogTimestampMap());

Expand All @@ -69,18 +73,28 @@ 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();
LOG.info("Execute roll log procedure for incremental backup ...");
BackupUtils.logRoll(conn, backupInfo.getBackupRootDir(), conf);

newTimestamps = readRegionServerLastLogRollResult();
Map<String, Long> rolledHosts =
BackupUtils.getRolledHosts(previousLogRollByHost, readRegionServerLastLogRollResult());

logList = getLogFilesForNewBackup(previousTimestampMins, newTimestamps, conf);
logList = excludeProcV2WALs(logList);
LogFileSelection selection = getLogFilesForNewBackup(previousTimestampMins, rolledHosts, conf);
logList = excludeProcV2WALs(selection.included());
backupInfo.setIncrBackupFileList(logList);

Map<String, Long> newTimestamps = BackupUtils.computeLogBoundaries(rolledHosts,
previousTimestampMins, selection.hostsWithLogs(), logList, selection.heldBack());
LOG.debug("Log boundaries for incremental backup {}: {}", backupInfo.getBackupId(),
newTimestamps);
return newTimestamps;
}

private record LogFileSelection(List<String> included, List<String> heldBack,
Set<String> hostsWithLogs) {
}

private List<String> excludeProcV2WALs(List<String> logList) {
List<String> list = new ArrayList<>();
for (int i = 0; i < logList.size(); i++) {
Expand All @@ -97,15 +111,17 @@ 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
* @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,

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

If a server crashes as we're doing log rolls + backups, its WALs will sit in the /WALs directory until that split is finished. We need to have a concept of WALs that were "held back", which are set as LogFileSelection#pending.

We will keep the boundary at before the oldest held back WAL, ensuring that the log cleaner doesn't delete it. These WALs will be backed up in the next incremental.

Map<String, Long> newestTimestamps, Configuration conf) throws IOException {
LOG.debug("In getLogFilesForNewBackup()\n" + "olderTimestamps: " + olderTimestamps
+ "\n newestTimestamps: " + newestTimestamps);
Expand All @@ -119,6 +135,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 @@ -143,6 +160,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 @@ -198,6 +216,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 @@ -217,18 +236,14 @@ private List<String> getLogFilesForNewBackup(Map<String, Long> olderTimestamps,
resultLogFiles.add(currentLogFile);
}

// It is possible that a host in .oldlogs is an obsolete region server
// so newestTimestamps.get(host) here can be null.
// Even if these logs belong to a obsolete region server, we still need
// to include they to avoid loss of edits for backup.
Long newTimestamp = newestTimestamps.get(host);

@hgromer hgromer Sep 23, 2026 •

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I know there was debate whether this was necessary or not here, and want to emphasize that this is a very important change. Without this change, we're subject to data loss anytime a new RS joins the cluster.

If we add just this code snippet back, our TestBackupOfflineRS tests will fail, showing this is in fact problematic. Importantly, adding a file to newestLog means it is not included in the backup. The roll time moves forward, and then the log cleaner deletes these logs without them ever being included in the backup.

There was some discussion about modifying this check due to hypothetical scenarios. I think it's a lot safer to simply remove this block which can cause data loss and provides no value. Simply removing this check keeps things safe. And I don't think the scenario here warrants potential data loss. Even if this scenario did play out, the only consequence is that we include a WAL that we don't technically need in a backup. This is a non-issue, given that our restore process handles this without any issues.

Finally, there was dicussion around trying to ensure "cross-table consistency" and data being out of sync which I disagree with as well. This backup is a point-in-time snapshot. In a distributed system, each RS is going to have it's own point in time. This is similarly true for snapshots, where one region's snapshot procedure can happen minutes prior to another one.

Regardless, I think preventing data loss scenarios and simplifying the code should be the priority here.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I still argue to not delete this check, but to change it to if (newTimestamp != null && currentLogTS > newTimestamp) instead. I think the removal (as suggested in the PR) can lead to data loss. This adjusted version fixes the original data loss issue and does not break the newly added tests.

Two things have to occur on a region server during an incremental backup to cause data loss:

  1. Two quick wall rolls. The backup's own roll creates a new WAL, C. Before the backup finishes listing oldWALs, the region server has to roll twice more: C → T, then T → U. Only then is T closed and eligible for archiving.
  2. T is archived while C isn't. cleanOldLogs archives any closed WAL whose regions have all flushed past its edits. If C holds edits for a region that T doesn't touch, and that region hasn't flushed yet, T gets archived first.

Of these 2 conditions, I expect (2) to occur commonly. For (1): The window runs from the backup's roll to the end of the oldWALs listing, and that's normally a few seconds. Two rolls might fit in it when one of these holds:

  • Heavy ingest. By default a WAL rolls when it reaches 0.5 × (2 × HDFS block size), which is about 128 MB. A region server ingesting around 50–100 MB/s rolls every 1–3 seconds. The hourly time-based roll (hbase.regionserver.logroll.period) is far too slow to matter.
  • A slow listing. getLogFilesForNewBackup runs one listStatus per region server directory, then a recursive listing of oldWALs. BackupLogCleaner keeps every WAL that a backup root still needs, so with backups enabled oldWALs can hold tens or hundreds of thousands of files. A large cluster or a slow NameNode can stretch the window to tens of seconds.

A rare case, but one that can occur on any incremental backup on any region server. For a large cluster with daily backups, it is bound to happen eventually. The result is a WAL that never gets included in a backup, leading to data loss when restoring.

I agree that simpler is better, and I wished lots of the backup code was simpler, but in this case the complexity is worth it.

Example test case (breaks on the current PR, succeeds on the current PR + suggested adjustment):

// paste me in TestIncrementalBackupManager 

  /**
   * WALs can be archived out of order. This test sets up the WAL directories as an incremental
   * backup would find them if, after its log roll, a newer WAL got archived while an older one was
   * still in the WALs directory. The older WAL must still end up in a backup.
   */
  @Test
  public void testOutOfOrderArchivedWALDoesNotSkipOlderWAL() throws Exception {
    List<TableName> tables = List.of(table1);
    HRegionServer rs = TEST_UTIL.getMiniHBaseCluster().getRegionServer(0);
    ServerName serverName = rs.getServerName();
    Path walRootDir = CommonFSUtils.getWALRootDir(conf1);
    FileSystem fs = walRootDir.getFileSystem(conf1);

    try (Connection conn = ConnectionFactory.createConnection(conf1);
      BackupAdminImpl backupAdmin = new BackupAdminImpl(conn)) {
      String fullBackupId =
        backupAdmin.backupTables(createBackupRequest(BackupType.FULL, tables, BACKUP_ROOT_DIR));
      assertTrue(checkSucceeded(fullBackupId));

      // Both test WALs are newer than any WAL that exists at this point, so they are newer than the
      // WAL reported by the log roll of the next incremental backup.
      long olderWALTs = EnvironmentEdgeManager.currentTime() + 1;
      long newerWALTs = olderWALTs + 1;
      Path walDir =
        new Path(walRootDir, AbstractFSWALProvider.getWALDirectoryName(serverName.toString()));
      Path olderWAL =
        new Path(walDir, serverName.toString() + BackupUtils.LOGNAME_SEPARATOR + olderWALTs);
      Path archiveDir = new Path(walRootDir,
        AbstractFSWALProvider.getWALArchiveDirectoryName(conf1, serverName.toString()));
      Path newerArchivedWAL =
        new Path(archiveDir, "wal" + BackupUtils.LOGNAME_SEPARATOR + newerWALTs);
      fs.create(olderWAL).close();
      fs.mkdirs(archiveDir);
      fs.create(newerArchivedWAL).close();

      try {
        List<String> firstBackupFiles;
        try (IncrementalBackupManager manager = new IncrementalBackupManager(conn, conf1)) {
          BackupInfo backupInfo = manager.createBackupInfo("backup_incr_1",
            BackupType.INCREMENTAL, tables, BACKUP_ROOT_DIR, -1, -1, false, false);
          Map<String, Long> boundaries = manager.getIncrBackupLogFileMap();
          manager.writeRegionServerLogTimestamp(backupInfo.getTables(), boundaries);
          firstBackupFiles = backupInfo.getIncrBackupFileList();
        }

        // Roll once more, so the log roll of the next backup reports a WAL that is newer than both
        // test WALs.
        TEST_UTIL.waitFor(30_000, () -> EnvironmentEdgeManager.currentTime() > newerWALTs);
        rs.getWalRoller().requestRollAll();
        rs.getWalRoller().waitUntilWalRollFinished();

        List<String> secondBackupFiles;
        try (IncrementalBackupManager manager = new IncrementalBackupManager(conn, conf1)) {
          BackupInfo backupInfo = manager.createBackupInfo("backup_incr_2",
            BackupType.INCREMENTAL, tables, BACKUP_ROOT_DIR, -1, -1, false, false);
          manager.getIncrBackupLogFileMap();
          secondBackupFiles = backupInfo.getIncrBackupFileList();
        }

        assertTrue(
          firstBackupFiles.contains(olderWAL.toString())
            || secondBackupFiles.contains(olderWAL.toString()),
          "WAL " + olderWAL + " was not included in any backup. First backup: " + firstBackupFiles
            + ", second backup: " + secondBackupFiles);
      } finally {
        fs.delete(olderWAL, false);
        fs.delete(newerArchivedWAL, false);
      }
    }
  }

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Hmm, okay I think I see your point, I didn't realize that WAL archival was random. This comment makes sense. In that case, I think your proposed change makes sense, while still preventing the data loss scenario that I'm trying to address

if (newTimestamp == null || currentLogTS > newTimestamp) {
if (newTimestamp != null && currentLogTS > newTimestamp) {
newestLogs.add(currentLogFile);
}
}
// remove newest log per host because they are still in use
resultLogFiles.removeAll(newestLogs);
return resultLogFiles;
return new LogFileSelection(resultLogFiles, newestLogs, hostsWithLogs);
}

static class NewestLogFilter implements PathFilter {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,13 +31,15 @@
import java.text.SimpleDateFormat;
import java.time.ZoneOffset;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
import java.util.Comparator;
import java.util.Date;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Map.Entry;
import java.util.Set;
import java.util.TimeZone;
import java.util.TreeSet;
import java.util.function.Predicate;
Expand Down Expand Up @@ -206,6 +208,84 @@ public static String parseHostNameFromLogFile(Path p) {
}
}

/**
* Returns the log roll result of every region server that took part in the log roll that happened
* between reading {@code previousLogRolls} and {@code latestLogRolls}. A region server that did
* not take part (for example because it is offline) keeps its old roll result, which is never
* removed, so it is recognized by its roll result not changing.
* @param previousLogRolls roll results by host, read before the log roll
* @param latestLogRolls roll results by host, read after the log roll
* @return roll results by host, for the hosts that took part in the log roll
*/
public static Map<String, Long> getRolledHosts(Map<String, Long> previousLogRolls,
Map<String, Long> 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<String, Long> computeLogBoundaries(Map<String, Long> rolledHosts,
Collection<String> coveredLogs, Collection<String> 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.
* <ul>
* <li>A host that took part in the backup's log roll gets its roll result.</li>
* <li>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.</li>
* <li>A host without covered or pending logs keeps its previous boundary while it still has logs,
* and otherwise gets no boundary.</li>
* </ul>
* 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<String, Long> computeLogBoundaries(Map<String, Long> rolledHosts,
Map<String, Long> previousBoundaries, Set<String> hostsWithLogs, Collection<String> coveredLogs,
Collection<String> pendingLogs) throws IOException {
Map<String, Long> 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<String, Long> boundaries = new HashMap<>(rolledHosts);
boundaries.putAll(otherHosts);
for (Entry<String, Long> 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
Expand Down
Loading
Loading