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 @@ -1037,6 +1037,12 @@ CommitResult tryCommitOnce(
if (snapshot.commitUser().equals(commitUser)
&& snapshot.commitIdentifier() == identifier
&& snapshot.commitKind() == commitKind) {
LOG.warn(
"Snapshot #{} of table {} was already committed by a previous "
+ "attempt of this commit. The LATEST hint may not have been "
+ "updated.",
snapshot.id(),
tableName);
lastCommittedSnapshotId = snapshot.id();
return new SuccessCommitResult();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -502,12 +502,19 @@ public ChangelogManager changelogManager() {

@Override
public ExpireSnapshots newExpireSnapshots() {
// The LATEST hint is only written when snapshots are committed by renaming files, that is
// unless the catalog manages the snapshots of this table (see
// CatalogEnvironment#snapshotCommit).
boolean protectLatestHint =
!(catalogEnvironment.catalogLoader() != null
&& catalogEnvironment.supportsVersionManagement());
return new ExpireSnapshotsImpl(
snapshotManager(),
changelogManager(),
store().newSnapshotDeletion(),
store().newTagManager(),
store().options().scanManifestParallelism());
store().options().scanManifestParallelism(),
protectLatestHint);
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -72,16 +72,38 @@ public class ExpireSnapshotsImpl implements ExpireSnapshots {
private final Executor fileExecutor;
private final TagManager tagManager;
private final int snapshotExpireBatchSize;
/** Whether commits keep the LATEST hint up to date, see {@link #keepHintedSnapshots}. */
private final boolean protectLatestHint;

@Nullable private Long lastWarnedLatestHint;
private boolean latestHintUncheckable;

private ExpireConfig expireConfig;
private Supplier<Long> currentTimeMillis = System::currentTimeMillis;

@VisibleForTesting
public ExpireSnapshotsImpl(
SnapshotManager snapshotManager,
ChangelogManager changelogManager,
SnapshotDeletion snapshotDeletion,
TagManager tagManager,
@Nullable Integer scanManifestParallelism) {
this(
snapshotManager,
changelogManager,
snapshotDeletion,
tagManager,
scanManifestParallelism,
true);
}

public ExpireSnapshotsImpl(
SnapshotManager snapshotManager,
ChangelogManager changelogManager,
SnapshotDeletion snapshotDeletion,
TagManager tagManager,
@Nullable Integer scanManifestParallelism,
boolean protectLatestHint) {
this.snapshotManager = snapshotManager;
this.changelogManager = changelogManager;
this.consumerManager =
Expand All @@ -97,6 +119,7 @@ public ExpireSnapshotsImpl(
scanManifestParallelism == null || scanManifestParallelism <= 0
? Runtime.getRuntime().availableProcessors()
: scanManifestParallelism;
this.protectLatestHint = protectLatestHint;
}

@VisibleForTesting
Expand Down Expand Up @@ -176,8 +199,91 @@ public int expire() {
return expireUntil(earliest, maxExclusive);
}

/**
* Keeps the snapshot the LATEST hint points to and all later ones. {@code findLatest} trusts
* the hint as long as the snapshot after it does not exist, so expiring the hinted snapshot and
* the one after it would make it return a missing snapshot. If the hinted snapshot is already
* missing, {@code findLatest} only works while the snapshot after it exists, so that one and
* all later ones are kept.
*
* @return the exclusive end to expire until, or null if nothing should be expired because the
* boundary cannot be established
*/
@Nullable
private Long keepHintedSnapshots(long endExclusiveId) {
Optional<Long> latestHint;
Long kept;
try {
latestHint = snapshotManager.readLatestHintStrictly();
if (!latestHint.isPresent() || latestHint.get() >= endExclusiveId) {
latestHintUncheckable = false;
return endExclusiveId;
}
long hint = latestHint.get();
if (snapshotManager.snapshotExists(hint)) {
kept = hint;
} else if (snapshotManager.snapshotExists(hint + 1)) {
kept = hint + 1;
} else {
kept = null;
}
} catch (Exception e) {
if (!latestHintUncheckable) {
latestHintUncheckable = true;
LOG.warn(
"Cannot check the LATEST hint in {}, skip expiring snapshots "
+ "until it can be checked.",
snapshotManager.snapshotDirectory(),
e);
}
return null;
}
latestHintUncheckable = false;

Long hint = latestHint.get();
if (!hint.equals(lastWarnedLatestHint)) {
lastWarnedLatestHint = hint;
if (kept != null) {
LOG.warn(
"The LATEST hint {} in {} is behind the latest snapshot {}. "
+ "Keeping snapshot {} and later ones instead of expiring up to {}.",
hint,
snapshotManager.snapshotDirectory(),
listLatestSnapshotId(),
kept,
endExclusiveId);
} else {
LOG.warn(
"The LATEST hint {} in {} points to a snapshot that does not exist, "
+ "neither does the next one, while the latest snapshot is {}. "
+ "Skip expiring snapshots until the hint is updated.",
hint,
snapshotManager.snapshotDirectory(),
listLatestSnapshotId());
}
}
return kept;
}

private String listLatestSnapshotId() {
try {
return String.valueOf(
snapshotManager.snapshotIdStream().reduce(Math::max).orElse(null));
} catch (Exception e) {
return "unknown";
}
}

@VisibleForTesting
public int expireUntil(long earliestId, long endExclusiveId) {
if (protectLatestHint && endExclusiveId > earliestId) {
Long kept = keepHintedSnapshots(endExclusiveId);
if (kept == null) {
return 0;
}
endExclusiveId = kept;
}

try {
return innerExpireUntil(earliestId, endExclusiveId);
} catch (InterruptedException e) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
import javax.annotation.Nullable;

import java.io.IOException;
import java.util.Optional;
import java.util.concurrent.ThreadLocalRandom;
import java.util.concurrent.TimeUnit;
import java.util.function.BinaryOperator;
Expand All @@ -40,6 +41,42 @@ public class HintFileUtils {
private static final int READ_HINT_RETRY_NUM = 3;
private static final int READ_HINT_RETRY_INTERVAL = 1;

/**
* Reads a hint like {@link #readHint}, but throws if the hint file cannot be read instead of
* returning null. Content that is not a positive number is ignored by {@link #findLatest}, so
* it is returned as absent.
*
* @return the hinted id, or empty if there is no usable hint
* @throws IOException if the hint file cannot be read after retries
*/
public static Optional<Long> readHintStrictly(FileIO fileIO, String fileName, Path dir)
throws IOException {
Path path = new Path(dir, fileName);
int retryNumber = 0;
while (true) {
Optional<String> content;
try {
content = fileIO.readOverwrittenFileUtf8(path);
} catch (IOException e) {
if (++retryNumber >= READ_HINT_RETRY_NUM) {
throw e;
}
try {
TimeUnit.MILLISECONDS.sleep(READ_HINT_RETRY_INTERVAL);
} catch (InterruptedException ie) {
Thread.currentThread().interrupt();
throw new RuntimeException(ie);
}
continue;
}
try {
return content.map(Long::parseLong).filter(id -> id > 0);
} catch (NumberFormatException e) {
return Optional.empty();
}
}
}

@Nullable
public static Long findLatest(FileIO fileIO, Path dir, String prefix, Function<Long, Path> file)
throws IOException {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -832,6 +832,16 @@ public static int findPreviousOrEqualSnapshot(
return -1;
}

/**
* Reads the LATEST hint, telling an absent hint apart from one that cannot be read.
*
* @return the hinted snapshot id, or empty if there is no usable LATEST hint
* @throws IOException if the hint file cannot be read
*/
public Optional<Long> readLatestHintStrictly() throws IOException {
return HintFileUtils.readHintStrictly(fileIO, HintFileUtils.LATEST, snapshotDirectory());
}

public void deleteLatestHint() throws IOException {
HintFileUtils.deleteLatestHint(fileIO, snapshotDirectory());
}
Expand Down
Loading
Loading