Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
package org.apache.paimon.operation;

import org.apache.paimon.CoreOptions;
import org.apache.paimon.Snapshot;
import org.apache.paimon.catalog.Catalog;
import org.apache.paimon.catalog.Identifier;
import org.apache.paimon.fs.FileStatus;
Expand Down Expand Up @@ -47,6 +48,7 @@
import java.util.concurrent.Future;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.Consumer;
import java.util.function.Function;
Expand Down Expand Up @@ -108,12 +110,20 @@ public CleanOrphanFilesResult clean()
}
candidateDeletes = new HashSet<>(candidates.keySet());

AtomicBoolean missingManifest = new AtomicBoolean(false);

// find used files
Set<String> usedFiles =
branches.stream()
.flatMap(branch -> getUsedFiles(branch).stream())
.flatMap(branch -> getUsedFiles(branch, missingManifest).stream())
.collect(Collectors.toSet());

if (missingManifest.get()) {
LOG.warn("Detected missing manifest during used-files collection, aborting clean.");
return new CleanOrphanFilesResult(
deleteFiles.size(), deletedFilesLenInBytes.get(), deleteFiles);
}

// delete unused files
candidateDeletes.removeAll(usedFiles);
candidateDeletes.stream()
Expand Down Expand Up @@ -157,35 +167,58 @@ private void cleanEmptyDataDirectory(List<Path> deleteFiles) {
}

private void collectWithoutDataFile(
String branch, Consumer<String> usedFileConsumer, Consumer<String> manifestConsumer)
String branch,
Consumer<String> usedFileConsumer,
Consumer<String> manifestConsumer,
Consumer<String> liveManifestConsumer,
AtomicBoolean missingManifest)
throws IOException {
Set<Snapshot> liveSnapshots = safelyGetLiveSnapshots(branch);
Set<Snapshot> snapshots = snapshotsIncludingTagsAndChangelogs(branch, liveSnapshots);
randomlyOnlyExecute(
executor,
snapshot -> {
try {
boolean live = liveSnapshots.contains(snapshot);
Consumer<String> perSnapshotManifestConsumer =
live
? manifest -> {
manifestConsumer.accept(manifest);
liveManifestConsumer.accept(manifest);
}
: manifestConsumer;
collectWithoutDataFile(
branch, snapshot, usedFileConsumer, manifestConsumer);
branch,
snapshot,
usedFileConsumer,
perSnapshotManifestConsumer,
live ? missingManifest : null);
} catch (IOException e) {
throw new RuntimeException(e);
}
},
safelyGetAllSnapshots(branch));
snapshots);
}

private Set<String> getUsedFiles(String branch) {
private Set<String> getUsedFiles(String branch, AtomicBoolean missingManifest) {
Set<String> usedFiles = ConcurrentHashMap.newKeySet();
ManifestFile manifestFile =
table.switchToBranch(branch).store().manifestFileFactory().create();
try {
Set<String> manifests = ConcurrentHashMap.newKeySet();
collectWithoutDataFile(branch, usedFiles::add, manifests::add);
Set<String> liveManifests = ConcurrentHashMap.newKeySet();
collectWithoutDataFile(
branch, usedFiles::add, manifests::add, liveManifests::add, missingManifest);
randomlyOnlyExecute(
executor,
manifestName -> {
try {
AtomicBoolean fnfFallback =
liveManifests.contains(manifestName) ? missingManifest : null;
retryReadingFiles(
() -> manifestFile.readWithIOException(manifestName),
Collections.<ManifestEntry>emptyList())
Collections.<ManifestEntry>emptyList(),
fnfFallback)
.stream()
.map(ManifestEntry::file)
.forEach(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,6 @@
import org.apache.paimon.utils.Preconditions;
import org.apache.paimon.utils.SnapshotManager;
import org.apache.paimon.utils.SupplierWithIOException;
import org.apache.paimon.utils.TagManager;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
Expand All @@ -56,6 +55,7 @@
import java.util.Set;
import java.util.TimeZone;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.function.Consumer;
import java.util.function.Predicate;
import java.util.stream.Collectors;
Expand Down Expand Up @@ -243,22 +243,32 @@ protected boolean isManagedBlobPack(Path path) {
return path.getName().endsWith(ManagedBlobReferenceFile.MANAGED_BLOB_SUFFIX);
}

protected Set<Snapshot> safelyGetAllSnapshots(String branch) throws IOException {
protected Set<Snapshot> safelyGetLiveSnapshots(String branch) throws IOException {
FileStoreTable branchTable = table.switchToBranch(branch);
SnapshotManager snapshotManager = branchTable.snapshotManager();
ChangelogManager changelogManager = branchTable.changelogManager();
TagManager tagManager = branchTable.tagManager();
Set<Snapshot> readSnapshots = new HashSet<>(snapshotManager.safelyGetAllSnapshots());
readSnapshots.addAll(tagManager.taggedSnapshots());
readSnapshots.addAll(changelogManager.safelyGetAllChangelogs());
return readSnapshots;
return new HashSet<>(
branchTable.snapshotManager().safelyGetAllSnapshotsWithConsistentLatest());
}

protected Set<Snapshot> snapshotsIncludingTagsAndChangelogs(
String branch, Set<Snapshot> liveSnapshots) throws IOException {
FileStoreTable branchTable = table.switchToBranch(branch);
Set<Snapshot> snapshots = new HashSet<>(liveSnapshots);
snapshots.addAll(branchTable.tagManager().taggedSnapshots());
snapshots.addAll(branchTable.changelogManager().safelyGetAllChangelogs());
return snapshots;
}

protected Set<Snapshot> safelyGetAllSnapshots(String branch) throws IOException {
Set<Snapshot> liveSnapshots = safelyGetLiveSnapshots(branch);
return snapshotsIncludingTagsAndChangelogs(branch, liveSnapshots);
}

protected void collectWithoutDataFile(
String branch,
Snapshot snapshot,
Consumer<String> usedFileConsumer,
Consumer<String> manifestConsumer)
Consumer<String> manifestConsumer,
@Nullable AtomicBoolean missingManifest)
throws IOException {
Consumer<Pair<String, Boolean>> usedFileWithFlagConsumer =
fileAndFlag -> {
Expand All @@ -267,44 +277,43 @@ protected void collectWithoutDataFile(
}
usedFileConsumer.accept(fileAndFlag.getLeft());
};
collectWithoutDataFileWithManifestFlag(branch, snapshot, usedFileWithFlagConsumer);
collectWithoutDataFileWithManifestFlag(
branch, snapshot, usedFileWithFlagConsumer, missingManifest);
}

protected void collectWithoutDataFileWithManifestFlag(
String branch,
Snapshot snapshot,
Consumer<Pair<String, Boolean>> usedFileWithFlagConsumer)
Consumer<Pair<String, Boolean>> usedFileWithFlagConsumer,
@Nullable AtomicBoolean missingManifest)
throws IOException {
FileStoreTable branchTable = table.switchToBranch(branch);
ManifestList manifestList = branchTable.store().manifestListFactory().create();
IndexFileHandler indexFileHandler = branchTable.store().newIndexFileHandler();
List<ManifestFileMeta> manifestFileMetas = new ArrayList<>();
// changelog manifest
// collect changelog, delta and base manifest lists
List<String> manifestListNames = new ArrayList<>();
if (snapshot.changelogManifestList() != null) {
usedFileWithFlagConsumer.accept(Pair.of(snapshot.changelogManifestList(), false));
manifestFileMetas.addAll(
retryReadingFiles(
() ->
manifestList.readWithIOException(
snapshot.changelogManifestList()),
emptyList()));
manifestListNames.add(snapshot.changelogManifestList());
}

// delta manifest
if (snapshot.deltaManifestList() != null) {
usedFileWithFlagConsumer.accept(Pair.of(snapshot.deltaManifestList(), false));
manifestFileMetas.addAll(
retryReadingFiles(
() -> manifestList.readWithIOException(snapshot.deltaManifestList()),
emptyList()));
manifestListNames.add(snapshot.deltaManifestList());
}
manifestListNames.add(snapshot.baseManifestList());

// base manifest
usedFileWithFlagConsumer.accept(Pair.of(snapshot.baseManifestList(), false));
manifestFileMetas.addAll(
retryReadingFiles(
() -> manifestList.readWithIOException(snapshot.baseManifestList()),
emptyList()));
for (String manifestListName : manifestListNames) {
List<ManifestFileMeta> metas =
retryReadingFiles(
() -> manifestList.readWithIOException(manifestListName),
emptyList(),
missingManifest);
// A missing manifest list means we cannot determine all used files, so abort.
if (missingManifest != null && missingManifest.get()) {
return;
}
usedFileWithFlagConsumer.accept(Pair.of(manifestListName, false));
manifestFileMetas.addAll(metas);
}

// collect manifests
for (ManifestFileMeta manifest : manifestFileMetas) {
Expand All @@ -313,15 +322,19 @@ protected void collectWithoutDataFileWithManifestFlag(

// index files
String indexManifest = snapshot.indexManifest();
if (indexManifest != null && indexFileHandler.existsManifest(indexManifest)) {
usedFileWithFlagConsumer.accept(Pair.of(indexManifest, false));
retryReadingFiles(
if (indexManifest != null) {
List<IndexManifestEntry> indexEntries =
retryReadingFiles(
() -> indexFileHandler.readManifestWithIOException(indexManifest),
Collections.<IndexManifestEntry>emptyList())
.stream()
.map(IndexManifestEntry::indexFile)
.map(IndexFileMeta::fileName)
.forEach(name -> usedFileWithFlagConsumer.accept(Pair.of(name, false)));
Collections.<IndexManifestEntry>emptyList(),
missingManifest);
if (missingManifest == null || !missingManifest.get()) {
usedFileWithFlagConsumer.accept(Pair.of(indexManifest, false));
indexEntries.stream()
.map(IndexManifestEntry::indexFile)
.map(IndexFileMeta::fileName)
.forEach(name -> usedFileWithFlagConsumer.accept(Pair.of(name, false)));
}
}

// statistic file
Expand Down Expand Up @@ -444,12 +457,24 @@ protected List<FileStatus> tryBestListingDirs(Path dir) {
*/
protected static <T> T retryReadingFiles(SupplierWithIOException<T> reader, T defaultValue)
throws IOException {
return retryReadingFiles(reader, defaultValue, null);
}

protected static <T> T retryReadingFiles(
SupplierWithIOException<T> reader, T defaultValue, @Nullable AtomicBoolean fnfFallback)
throws IOException {
int retryNumber = 0;
IOException caught = null;
while (retryNumber++ < READ_FILE_RETRY_NUM) {
try {
return reader.get();
} catch (FileNotFoundException e) {
if (fnfFallback != null) {
fnfFallback.set(true);
LOG.warn(
"File not found while collecting used files, aborting orphan files clean.",
e);
}
return defaultValue;
} catch (IOException e) {
caught = e;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@
import java.util.HashSet;
import java.util.Iterator;
import java.util.List;
import java.util.Objects;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.ExecutorService;
Expand Down Expand Up @@ -580,6 +581,23 @@ public List<Snapshot> safelyGetAllSnapshots() throws IOException {
return snapshots;
}

public List<Snapshot> safelyGetAllSnapshotsWithConsistentLatest() throws IOException {
Long latestBefore = latestSnapshotIdFromFileSystem();
List<Snapshot> snapshots = safelyGetAllSnapshots();
Long latestAfter = latestSnapshotIdFromFileSystem();

boolean latestIncluded =
latestAfter == null
|| snapshots.stream().anyMatch(snapshot -> snapshot.id() == latestAfter);
if (!Objects.equals(latestBefore, latestAfter) || !latestIncluded) {
throw new IOException(
String.format(
"Incomplete snapshot enumeration: latest snapshot changed from %s to %s or was not included.",
latestBefore, latestAfter));
}
return snapshots;
}

private static void collectSnapshots(Consumer<Path> pathConsumer, List<Path> paths)
throws IOException {
ExecutorService executor =
Expand Down
Loading
Loading