diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/LocalOrphanFilesClean.java b/paimon-core/src/main/java/org/apache/paimon/operation/LocalOrphanFilesClean.java index a630d8543a4a..016d08fb9c6b 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/LocalOrphanFilesClean.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/LocalOrphanFilesClean.java @@ -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; @@ -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; @@ -108,12 +110,20 @@ public CleanOrphanFilesResult clean() } candidateDeletes = new HashSet<>(candidates.keySet()); + AtomicBoolean missingManifest = new AtomicBoolean(false); + // find used files Set 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() @@ -157,35 +167,58 @@ private void cleanEmptyDataDirectory(List deleteFiles) { } private void collectWithoutDataFile( - String branch, Consumer usedFileConsumer, Consumer manifestConsumer) + String branch, + Consumer usedFileConsumer, + Consumer manifestConsumer, + Consumer liveManifestConsumer, + AtomicBoolean missingManifest) throws IOException { + Set liveSnapshots = safelyGetLiveSnapshots(branch); + Set snapshots = snapshotsIncludingTagsAndChangelogs(branch, liveSnapshots); randomlyOnlyExecute( executor, snapshot -> { try { + boolean live = liveSnapshots.contains(snapshot); + Consumer 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 getUsedFiles(String branch) { + private Set getUsedFiles(String branch, AtomicBoolean missingManifest) { Set usedFiles = ConcurrentHashMap.newKeySet(); ManifestFile manifestFile = table.switchToBranch(branch).store().manifestFileFactory().create(); try { Set manifests = ConcurrentHashMap.newKeySet(); - collectWithoutDataFile(branch, usedFiles::add, manifests::add); + Set 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.emptyList()) + Collections.emptyList(), + fnfFallback) .stream() .map(ManifestEntry::file) .forEach( diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/OrphanFilesClean.java b/paimon-core/src/main/java/org/apache/paimon/operation/OrphanFilesClean.java index 4245460225ae..48e60897a1e5 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/OrphanFilesClean.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/OrphanFilesClean.java @@ -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; @@ -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; @@ -243,22 +243,32 @@ protected boolean isManagedBlobPack(Path path) { return path.getName().endsWith(ManagedBlobReferenceFile.MANAGED_BLOB_SUFFIX); } - protected Set safelyGetAllSnapshots(String branch) throws IOException { + protected Set safelyGetLiveSnapshots(String branch) throws IOException { FileStoreTable branchTable = table.switchToBranch(branch); - SnapshotManager snapshotManager = branchTable.snapshotManager(); - ChangelogManager changelogManager = branchTable.changelogManager(); - TagManager tagManager = branchTable.tagManager(); - Set readSnapshots = new HashSet<>(snapshotManager.safelyGetAllSnapshots()); - readSnapshots.addAll(tagManager.taggedSnapshots()); - readSnapshots.addAll(changelogManager.safelyGetAllChangelogs()); - return readSnapshots; + return new HashSet<>( + branchTable.snapshotManager().safelyGetAllSnapshotsWithConsistentLatest()); + } + + protected Set snapshotsIncludingTagsAndChangelogs( + String branch, Set liveSnapshots) throws IOException { + FileStoreTable branchTable = table.switchToBranch(branch); + Set snapshots = new HashSet<>(liveSnapshots); + snapshots.addAll(branchTable.tagManager().taggedSnapshots()); + snapshots.addAll(branchTable.changelogManager().safelyGetAllChangelogs()); + return snapshots; + } + + protected Set safelyGetAllSnapshots(String branch) throws IOException { + Set liveSnapshots = safelyGetLiveSnapshots(branch); + return snapshotsIncludingTagsAndChangelogs(branch, liveSnapshots); } protected void collectWithoutDataFile( String branch, Snapshot snapshot, Consumer usedFileConsumer, - Consumer manifestConsumer) + Consumer manifestConsumer, + @Nullable AtomicBoolean missingManifest) throws IOException { Consumer> usedFileWithFlagConsumer = fileAndFlag -> { @@ -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> usedFileWithFlagConsumer) + Consumer> usedFileWithFlagConsumer, + @Nullable AtomicBoolean missingManifest) throws IOException { FileStoreTable branchTable = table.switchToBranch(branch); ManifestList manifestList = branchTable.store().manifestListFactory().create(); IndexFileHandler indexFileHandler = branchTable.store().newIndexFileHandler(); List manifestFileMetas = new ArrayList<>(); - // changelog manifest + // collect changelog, delta and base manifest lists + List 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 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) { @@ -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 indexEntries = + retryReadingFiles( () -> indexFileHandler.readManifestWithIOException(indexManifest), - Collections.emptyList()) - .stream() - .map(IndexManifestEntry::indexFile) - .map(IndexFileMeta::fileName) - .forEach(name -> usedFileWithFlagConsumer.accept(Pair.of(name, false))); + Collections.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 @@ -444,12 +457,24 @@ protected List tryBestListingDirs(Path dir) { */ protected static T retryReadingFiles(SupplierWithIOException reader, T defaultValue) throws IOException { + return retryReadingFiles(reader, defaultValue, null); + } + + protected static T retryReadingFiles( + SupplierWithIOException 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; diff --git a/paimon-core/src/main/java/org/apache/paimon/utils/SnapshotManager.java b/paimon-core/src/main/java/org/apache/paimon/utils/SnapshotManager.java index af21181214b3..37bc175fd38a 100644 --- a/paimon-core/src/main/java/org/apache/paimon/utils/SnapshotManager.java +++ b/paimon-core/src/main/java/org/apache/paimon/utils/SnapshotManager.java @@ -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; @@ -580,6 +581,23 @@ public List safelyGetAllSnapshots() throws IOException { return snapshots; } + public List safelyGetAllSnapshotsWithConsistentLatest() throws IOException { + Long latestBefore = latestSnapshotIdFromFileSystem(); + List 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 pathConsumer, List paths) throws IOException { ExecutorService executor = diff --git a/paimon-core/src/test/java/org/apache/paimon/operation/LocalOrphanFilesCleanTest.java b/paimon-core/src/test/java/org/apache/paimon/operation/LocalOrphanFilesCleanTest.java index cd6480d7e74c..c45e4d35df7b 100644 --- a/paimon-core/src/test/java/org/apache/paimon/operation/LocalOrphanFilesCleanTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/operation/LocalOrphanFilesCleanTest.java @@ -83,6 +83,7 @@ import java.util.function.Predicate; import java.util.stream.Collectors; +import static org.apache.paimon.catalog.Identifier.DEFAULT_MAIN_BRANCH; import static org.apache.paimon.utils.BranchManager.branchPath; import static org.apache.paimon.utils.FileStorePathFactory.BUCKET_PATH_PREFIX; import static org.assertj.core.api.Assertions.assertThat; @@ -212,8 +213,13 @@ public void normallyRemoving(Path dataPath) throws Throwable { table.copy(expireOptions.toMap()).newCommit("").expireSnapshots(); // randomly delete tags - List deleteTags = Collections.emptyList(); - deleteTags = randomlyPick(allTags); + String branchBaseTag = allTags.get(0); + List deletableTags = + allTags.stream() + .filter(tag -> !tag.equals(branchBaseTag)) + .collect(Collectors.toList()); + List deleteTags = + deletableTags.isEmpty() ? Collections.emptyList() : randomlyPick(deletableTags); for (String tagName : deleteTags) { table.deleteTag(tagName); } @@ -315,8 +321,13 @@ public void testNormallyRemovingMixedWithExternalPath() throws Throwable { table.copy(expireOptions.toMap()).newCommit("").expireSnapshots(); // randomly delete tags - List deleteTags = Collections.emptyList(); - deleteTags = randomlyPick(allTags); + String branchBaseTag = allTags.get(0); + List deletableTags = + allTags.stream() + .filter(tag -> !tag.equals(branchBaseTag)) + .collect(Collectors.toList()); + List deleteTags = + deletableTags.isEmpty() ? Collections.emptyList() : randomlyPick(deletableTags); for (String tagName : deleteTags) { table.deleteTag(tagName); } @@ -550,6 +561,10 @@ public void testAbnormallyRemoving() throws Exception { commit(generateData()); } + List dataFilesBefore = new ArrayList<>(); + collectDataFiles(tablePath, dataFilesBefore); + assertThat(dataFilesBefore).isNotEmpty(); + // randomly delete a manifest file of snapshot 1 SnapshotManager snapshotManager = table.snapshotManager(); Snapshot snapshot1 = snapshotManager.snapshot(1); @@ -567,7 +582,12 @@ public void testAbnormallyRemoving() throws Exception { LocalOrphanFilesClean orphanFilesClean = new LocalOrphanFilesClean( table, System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(2)); - assertThat(orphanFilesClean.clean().getDeletedFilesPath().size()).isGreaterThan(0); + CleanOrphanFilesResult result = orphanFilesClean.clean(); + + assertThat(result.getDeletedFileCount()).isEqualTo(0); + for (Path dataFile : dataFilesBefore) { + assertThat(fileIO.exists(dataFile)).isTrue(); + } } @Test @@ -704,6 +724,127 @@ void testDirectoryInSnapshotDirNotTreatedAsCandidate() throws Exception { assertThat(fileIO.exists(unknownDir)).isTrue(); } + @Test + void testAbortWhenManifestListMissingDuringConcurrentExpiration() throws Exception { + commit(Collections.singletonList(new TestPojo(1, 0, "a", "v1"))); + commit(Collections.singletonList(new TestPojo(2, 0, "b", "v2"))); + commit(Collections.singletonList(new TestPojo(3, 1, "c", "v3"))); + + List dataFilesBefore = new ArrayList<>(); + collectDataFiles(tablePath, dataFilesBefore); + assertThat(dataFilesBefore).isNotEmpty(); + + Snapshot latest = table.snapshotManager().latestSnapshot(); + Path missingManifestList = new Path(manifestDir, latest.baseManifestList()); + assertThat(fileIO.exists(missingManifestList)).isTrue(); + fileIO.deleteQuietly(missingManifestList); + + LocalOrphanFilesClean orphanFilesClean = + new LocalOrphanFilesClean( + table, System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(2)); + CleanOrphanFilesResult result = orphanFilesClean.clean(); + + assertThat(result.getDeletedFileCount()).isEqualTo(0); + for (Path dataFile : dataFilesBefore) { + assertThat(fileIO.exists(dataFile)).isTrue(); + } + } + + @Test + void testAbortWhenBranchManifestListMissingDuringConcurrentExpiration() throws Exception { + commit(Collections.singletonList(new TestPojo(1, 0, "a", "v1"))); + + String branchName = "branch1"; + table.createBranch(branchName); + FileStoreTable branchTable = table.switchToBranch(branchName); + String branchCommitUser = UUID.randomUUID().toString(); + try (TableWriteImpl branchWrite = branchTable.newWrite(branchCommitUser); + TableCommitImpl branchCommit = branchTable.newCommit(branchCommitUser)) { + branchWrite.write(new TestPojo(2, 0, "a", "v2").toRow(RowKind.INSERT)); + branchCommit.commit(0, branchWrite.prepareCommit(true, 0)); + branchWrite.write(new TestPojo(3, 0, "b", "v3").toRow(RowKind.INSERT)); + branchCommit.commit(1, branchWrite.prepareCommit(true, 1)); + branchWrite.write(new TestPojo(4, 1, "c", "v4").toRow(RowKind.INSERT)); + branchCommit.commit(2, branchWrite.prepareCommit(true, 2)); + } + + List dataFilesBefore = new ArrayList<>(); + collectDataFiles(tablePath, dataFilesBefore); + assertThat(dataFilesBefore).isNotEmpty(); + + Snapshot branchLatest = branchTable.snapshotManager().latestSnapshot(); + Path missingManifestList = new Path(manifestDir, branchLatest.baseManifestList()); + assertThat(fileIO.exists(missingManifestList)).isTrue(); + fileIO.deleteQuietly(missingManifestList); + + LocalOrphanFilesClean orphanFilesClean = + new LocalOrphanFilesClean( + table, System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(2)); + CleanOrphanFilesResult result = orphanFilesClean.clean(); + + assertThat(result.getDeletedFileCount()).isEqualTo(0); + for (Path dataFile : dataFilesBefore) { + assertThat(fileIO.exists(dataFile)).isTrue(); + } + } + + @Test + void testReuseCapturedLiveSnapshots() throws Exception { + commit(Collections.singletonList(new TestPojo(1, 0, "a", "v1"))); + Snapshot mainSnapshot = table.snapshotManager().latestSnapshot(); + + String branchName = "branch1"; + table.createBranch(branchName); + FileStoreTable branchTable = table.switchToBranch(branchName); + String branchCommitUser = UUID.randomUUID().toString(); + try (TableWriteImpl branchWrite = branchTable.newWrite(branchCommitUser); + TableCommitImpl branchCommit = branchTable.newCommit(branchCommitUser)) { + branchWrite.write(new TestPojo(2, 0, "b", "v2").toRow(RowKind.INSERT)); + branchCommit.commit(0, branchWrite.prepareCommit(true, 0)); + } + + List dataFilesBefore = new ArrayList<>(); + collectDataFiles(tablePath, dataFilesBefore); + assertThat(dataFilesBefore).isNotEmpty(); + + LocalOrphanFilesClean orphanFilesClean = + new LocalOrphanFilesClean( + table, System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(2)) { + @Override + protected Set snapshotsIncludingTagsAndChangelogs( + String branch, Set liveSnapshots) throws IOException { + if (DEFAULT_MAIN_BRANCH.equals(branch)) { + fileIO.deleteQuietly( + table.snapshotManager().snapshotPath(mainSnapshot.id())); + fileIO.deleteQuietly( + new Path(manifestDir, mainSnapshot.baseManifestList())); + } + return super.snapshotsIncludingTagsAndChangelogs(branch, liveSnapshots); + } + }; + + CleanOrphanFilesResult result = orphanFilesClean.clean(); + + assertThat(result.getDeletedFileCount()).isEqualTo(0); + for (Path dataFile : dataFilesBefore) { + assertThat(fileIO.exists(dataFile)).isTrue(); + } + } + + private void collectDataFiles(Path dir, List result) throws IOException { + FileStatus[] statuses = fileIO.listStatus(dir); + if (statuses == null) { + return; + } + for (FileStatus status : statuses) { + if (status.isDir()) { + collectDataFiles(status.getPath(), result); + } else if (status.getPath().getParent().getName().startsWith(BUCKET_PATH_PREFIX)) { + result.add(status.getPath()); + } + } + } + private void writeData( SnapshotManager snapshotManager, List> committedData, diff --git a/paimon-core/src/test/java/org/apache/paimon/utils/SnapshotManagerTest.java b/paimon-core/src/test/java/org/apache/paimon/utils/SnapshotManagerTest.java index 02b21ba90637..74bd578a573d 100644 --- a/paimon-core/src/test/java/org/apache/paimon/utils/SnapshotManagerTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/utils/SnapshotManagerTest.java @@ -40,8 +40,10 @@ import java.util.List; import java.util.Set; import java.util.concurrent.ThreadLocalRandom; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; import java.util.stream.Collectors; +import java.util.stream.Stream; import static org.apache.paimon.SnapshotTest.newChangelogManager; import static org.apache.paimon.SnapshotTest.newSnapshotManager; @@ -272,6 +274,49 @@ public void testLaterOrEqualWatermark(boolean isRaceCondition) throws IOExceptio assertThat(snapshotManager.laterOrEqualWatermark(millis + 999)).isNull(); } + @Test + public void testDetectIncompleteSnapshotEnumeration() throws IOException { + FileIO localFileIO = LocalFileIO.create(); + long millis = 1684726826L; + AtomicBoolean triggerRace = new AtomicBoolean(true); + SnapshotManager snapshotManager = + new SnapshotManager( + localFileIO, + new Path(tempDir.toString()), + DEFAULT_MAIN_BRANCH, + null, + null) { + @Override + public Stream snapshotIdStream() throws IOException { + return super.snapshotIdStream() + .peek( + snapshotId -> { + if (!triggerRace.compareAndSet(true, false)) { + return; + } + Snapshot nextSnapshot = + createSnapshotWithMillis( + snapshotId + 1, millis + 1000); + try { + localFileIO.tryToWriteAtomic( + snapshotPath(nextSnapshot.id()), + nextSnapshot.toJson()); + localFileIO.delete(snapshotPath(snapshotId), false); + } catch (IOException e) { + throw new RuntimeException(e); + } + }); + } + }; + Snapshot snapshot = createSnapshotWithMillis(0, millis); + localFileIO.tryToWriteAtomic( + snapshotManager.snapshotPath(snapshot.id()), snapshot.toJson()); + + assertThatThrownBy(snapshotManager::safelyGetAllSnapshotsWithConsistentLatest) + .isInstanceOf(IOException.class) + .hasMessageContaining("latest snapshot changed from 0 to 1"); + } + public static Snapshot createSnapshotWithMillis(long id, long millis) { return new Snapshot( id, diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/orphan/FlinkOrphanFilesClean.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/orphan/FlinkOrphanFilesClean.java index 3ce2bf82f8ae..8fb7cafa3a0d 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/orphan/FlinkOrphanFilesClean.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/orphan/FlinkOrphanFilesClean.java @@ -37,6 +37,7 @@ import org.apache.flink.api.common.functions.ReduceFunction; import org.apache.flink.api.common.typeinfo.TypeInformation; import org.apache.flink.api.java.tuple.Tuple2; +import org.apache.flink.api.java.tuple.Tuple3; import org.apache.flink.configuration.Configuration; import org.apache.flink.configuration.CoreOptions; import org.apache.flink.configuration.ExecutionOptions; @@ -63,6 +64,7 @@ import java.util.List; import java.util.Map; import java.util.Set; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicLong; import java.util.function.Consumer; @@ -75,6 +77,9 @@ public class FlinkOrphanFilesClean extends OrphanFilesClean { protected static final Logger LOG = LoggerFactory.getLogger(FlinkOrphanFilesClean.class); + private static final String MISSING_MANIFEST_SENTINEL = + "__PAIMON_ORPHAN_CLEAN_MISSING_MANIFEST__"; + @Nullable protected final Integer parallelism; public FlinkOrphanFilesClean( @@ -143,48 +148,70 @@ public void processElement( .name("branch-snapshot-deletion-result"); // branch and manifest file - final OutputTag> manifestOutputTag = - new OutputTag>("manifest-output") {}; + final OutputTag> manifestOutputTag = + new OutputTag>("manifest-output") {}; SingleOutputStreamOperator usedManifestFiles = env.fromCollection(branches) .name("branch-source") .process( - new ProcessFunction>() { + new ProcessFunction>() { @Override public void processElement( String branch, - ProcessFunction>.Context + ProcessFunction> + .Context ctx, - Collector> out) + Collector> out) throws Exception { - for (Snapshot snapshot : safelyGetAllSnapshots(branch)) { - out.collect(new Tuple2<>(branch, snapshot.toJson())); + Set liveSnapshots = + safelyGetLiveSnapshots(branch); + for (Snapshot snapshot : + snapshotsIncludingTagsAndChangelogs( + branch, liveSnapshots)) { + out.collect( + new Tuple3<>( + branch, + snapshot.toJson(), + liveSnapshots.contains(snapshot))); } } }) .name("collect-snapshots") .rebalance() .process( - new ProcessFunction, String>() { + new ProcessFunction, String>() { @Override public void processElement( - Tuple2 branchAndSnapshot, - ProcessFunction, String>.Context + Tuple3 branchAndSnapshot, + ProcessFunction, String> + .Context ctx, Collector out) throws Exception { String branch = branchAndSnapshot.f0; Snapshot snapshot = Snapshot.fromJson(branchAndSnapshot.f1); + boolean isLiveSnapshot = branchAndSnapshot.f2; Consumer manifestConsumer = - manifest -> { - Tuple2 tuple2 = - new Tuple2<>(branch, manifest); - ctx.output(manifestOutputTag, tuple2); - }; + manifest -> + ctx.output( + manifestOutputTag, + new Tuple3<>( + branch, + manifest, + isLiveSnapshot)); + AtomicBoolean missingManifest = + isLiveSnapshot ? new AtomicBoolean(false) : null; collectWithoutDataFile( - branch, snapshot, out::collect, manifestConsumer); + branch, + snapshot, + out::collect, + manifestConsumer, + missingManifest); + if (missingManifest != null && missingManifest.get()) { + out.collect(MISSING_MANIFEST_SENTINEL); + } } }) .name("collect-manifests"); @@ -192,39 +219,48 @@ public void processElement( DataStream usedFiles = usedManifestFiles .getSideOutput(manifestOutputTag) - .keyBy(tuple2 -> tuple2.f0 + ":" + tuple2.f1) + .keyBy(tuple3 -> tuple3.f0 + ":" + tuple3.f1) .transform( "collect-used-files", STRING_TYPE_INFO, - new BoundedOneInputOperator, String>() { + new BoundedOneInputOperator< + Tuple3, String>() { - private final Set> manifests = - new HashSet<>(); + private final Map> + manifests = new HashMap<>(); @Override public void processElement( - StreamRecord> element) { - manifests.add(element.getValue()); + StreamRecord> element) { + Tuple3 value = element.getValue(); + manifests.merge( + value.f0 + ":" + value.f1, + value, + (a, b) -> new Tuple3<>(a.f0, a.f1, a.f2 || b.f2)); } @Override public void endInput() throws IOException { Map branchManifests = new HashMap<>(); - for (Tuple2 tuple2 : manifests) { + for (Tuple3 tuple : + manifests.values()) { ManifestFile manifestFile = branchManifests.computeIfAbsent( - tuple2.f0, + tuple.f0, key -> table.switchToBranch(key) .store() .manifestFileFactory() .create()); + AtomicBoolean manifestMissing = + tuple.f2 ? new AtomicBoolean(false) : null; retryReadingFiles( () -> manifestFile .readWithIOException( - tuple2.f1), - Collections.emptyList()) + tuple.f1), + Collections.emptyList(), + manifestMissing) .forEach( f -> { List files = @@ -237,6 +273,11 @@ public void endInput() throws IOException { new StreamRecord<>( file))); }); + if (manifestMissing != null && manifestMissing.get()) { + output.collect( + new StreamRecord<>( + MISSING_MANIFEST_SENTINEL)); + } } } }); @@ -354,11 +395,74 @@ public void endInput() throws IOException { .setParallelism(1) .setMaxParallelism(1); - DataStream deleted = + DataStream missingManifestSignals = + usedFiles + .filter(MISSING_MANIFEST_SENTINEL::equals) + .name("missing-manifest-signals"); + DataStream normalUsedFiles = usedFiles + .filter(file -> !MISSING_MANIFEST_SENTINEL.equals(file)) + .name("normal-used-files"); + + DataStream> candidatesAfterGlobalAbort = + missingManifestSignals + .broadcast() + .connect(candidates) + .transform( + "abort-candidates-on-missing-manifest", + candidates.getType(), + new BoundedTwoInputOperator< + String, Tuple2, Tuple2>() { + + private boolean signalEnd; + private boolean abortDeletion; + + @Override + public InputSelection nextSelection() { + return signalEnd + ? InputSelection.SECOND + : InputSelection.FIRST; + } + + @Override + public void endInput(int inputId) { + if (inputId == 1) { + checkState(!signalEnd, "Signal input already ended."); + signalEnd = true; + if (abortDeletion) { + LOG.warn( + "Detected missing manifest, aborting clean globally."); + } + } else { + checkState(signalEnd, "Signal input should end first."); + } + } + + @Override + public void processElement1(StreamRecord element) { + checkState( + MISSING_MANIFEST_SENTINEL.equals( + element.getValue()), + "Unexpected global abort signal."); + abortDeletion = true; + } + + @Override + public void processElement2( + StreamRecord> element) { + checkState(signalEnd, "Signal input should end first."); + if (!abortDeletion) { + output.collect(element); + } + } + }); + + DataStream deleted = + normalUsedFiles .keyBy(f -> f) .connect( - candidates.keyBy(pathAndSize -> new Path(pathAndSize.f0).getName())) + candidatesAfterGlobalAbort.keyBy( + pathAndSize -> new Path(pathAndSize.f0).getName())) .transform( "join-used-and-candidate-files", TypeInformation.of(CleanOrphanFilesResult.class), diff --git a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/action/RemoveOrphanFilesActionITCaseBase.java b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/action/RemoveOrphanFilesActionITCaseBase.java index e54fd5c66205..49cad290592a 100644 --- a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/action/RemoveOrphanFilesActionITCaseBase.java +++ b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/action/RemoveOrphanFilesActionITCaseBase.java @@ -19,6 +19,7 @@ package org.apache.paimon.flink.action; import org.apache.paimon.CoreOptions; +import org.apache.paimon.Snapshot; import org.apache.paimon.data.BinaryString; import org.apache.paimon.data.GenericRow; import org.apache.paimon.fs.FileIO; @@ -402,6 +403,40 @@ public void testRunWithMode(boolean isNamedArgument) throws Exception { .hasMessageContaining("Unknown mode"); } + @Test + public void testDistributedCleanAbortsGloballyWhenManifestListIsMissing() throws Exception { + FileStoreTable table = createTableAndWriteData(tableName); + FileIO fileIO = table.fileIO(); + + List orphanFiles = new ArrayList<>(); + for (int i = 0; i < 32; i++) { + Path orphanFile = getOrphanFilePath(table, "bucket-0/global-abort-orphan-" + i); + fileIO.writeFile(orphanFile, "orphan", true); + orphanFiles.add(orphanFile); + } + Thread.sleep(2000); + + Snapshot latestSnapshot = table.snapshotManager().latestSnapshot(); + Path manifestList = + table.store().pathFactory().toManifestListPath(latestSnapshot.baseManifestList()); + assertThat(fileIO.exists(manifestList)).isTrue(); + fileIO.deleteQuietly(manifestList); + + String olderThan = + DateTimeUtils.formatLocalDateTime( + DateTimeUtils.toLocalDateTime(System.currentTimeMillis()), 3); + String procedure = + String.format( + "CALL sys.remove_orphan_files('%s.%s', '%s', false, 5, 'distributed')", + database, tableName, olderThan); + ImmutableList result = ImmutableList.copyOf(executeSQL(procedure)); + + assertThat(result).containsOnly(Row.of("0")); + for (Path orphanFile : orphanFiles) { + assertThat(fileIO.exists(orphanFile)).isTrue(); + } + } + @Test public void testEmptyPartitionDirectories() throws Exception { FileStoreTable table = createPartitionedTableWithData(); diff --git a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/procedure/SparkOrphanFilesClean.scala b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/procedure/SparkOrphanFilesClean.scala index 428ac6e09763..f57371ca4208 100644 --- a/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/procedure/SparkOrphanFilesClean.scala +++ b/paimon-spark/paimon-spark-common/src/main/scala/org/apache/paimon/spark/procedure/SparkOrphanFilesClean.scala @@ -34,12 +34,13 @@ import org.apache.spark.sql.catalyst.SQLConfHelper import java.util import java.util.Collections -import java.util.concurrent.atomic.AtomicLong +import java.util.concurrent.atomic.{AtomicBoolean, AtomicLong} import java.util.function.Consumer import scala.collection.JavaConverters._ import scala.collection.mutable import scala.collection.mutable.ArrayBuffer +import scala.util.control.NonFatal case class SparkOrphanFilesClean( specifiedTable: FileStoreTable, @@ -51,7 +52,19 @@ case class SparkOrphanFilesClean( with SQLConfHelper with Logging { - def doOrphanClean(): (Dataset[(Long, Long)], Dataset[BranchAndManifestFile]) = { + def doOrphanClean(): (Dataset[(Long, Long)], Seq[Dataset[_]]) = { + val cachedDatasets = new ArrayBuffer[Dataset[_]]() + try { + doOrphanClean(cachedDatasets) + } catch { + case NonFatal(t) => + cachedDatasets.foreach(_.unpersist()) + throw t + } + } + + private def doOrphanClean( + cachedDatasets: ArrayBuffer[Dataset[_]]): (Dataset[(Long, Long)], Seq[Dataset[_]]) = { import spark.implicits._ val branches = validBranches() @@ -69,29 +82,57 @@ case class SparkOrphanFilesClean( val usedManifestFiles = spark.sparkContext .parallelize(branches.asScala.toSeq, maxBranchParallelism) .mapPartitions(_.flatMap { - branch => safelyGetAllSnapshots(branch).asScala.map(snapshot => (branch, snapshot.toJson)) + branch => + val liveSnapshots = safelyGetLiveSnapshots(branch) + snapshotsIncludingTagsAndChangelogs(branch, liveSnapshots).asScala.map( + snapshot => (branch, snapshot.toJson, liveSnapshots.contains(snapshot))) }) .repartition(parallelism) .flatMap { - case (branch, snapshotJson) => + case (branch, snapshotJson, isLive) => val usedFileBuffer = new ArrayBuffer[BranchAndManifestFile]() val usedFileConsumer = new Consumer[org.apache.paimon.utils.Pair[String, java.lang.Boolean]] { override def accept(pair: utils.Pair[String, java.lang.Boolean]): Unit = { - usedFileBuffer.append(BranchAndManifestFile(branch, pair.getLeft, pair.getRight)) + usedFileBuffer.append( + BranchAndManifestFile( + branch, + pair.getLeft, + pair.getRight, + isMissing = false, + isLiveSnapshot = isLive)) } } + val missingManifest = if (isLive) new AtomicBoolean(false) else null val snapshot = Snapshot.fromJson(snapshotJson) - collectWithoutDataFileWithManifestFlag(branch, snapshot, usedFileConsumer) + collectWithoutDataFileWithManifestFlag( + branch, + snapshot, + usedFileConsumer, + missingManifest) + if (missingManifest != null && missingManifest.get()) { + usedFileBuffer.append( + BranchAndManifestFile(branch, "", isManifestFile = false, isMissing = true)) + } usedFileBuffer } .toDS() .cache() + cachedDatasets += usedManifestFiles + + if (!usedManifestFiles.filter(_.isMissing).isEmpty) { + logWarning("Detected missing manifest during used-files collection, aborting clean.") + return ( + spark.createDataset( + Seq((deletedFilesCountInLocal.get(), deletedFilesLenInBytesInLocal.get()))), + cachedDatasets.toSeq) + } - // find all data files - val dataFiles = usedManifestFiles - .filter(_.isManifestFile) - .distinct() + val dataFilesWithFlag = usedManifestFiles + .filter(f => f.isManifestFile && !f.isMissing) + .groupByKey(f => (f.branch, f.manifestName)) + .reduceGroups((a, b) => a.copy(isLiveSnapshot = a.isLiveSnapshot || b.isLiveSnapshot)) + .map(_._2) .mapPartitions { it => val branchManifests = new util.HashMap[String, ManifestFile] @@ -102,19 +143,56 @@ case class SparkOrphanFilesClean( (key: String) => specifiedTable.switchToBranch(key).store.manifestFileFactory.create) - retryReadingFiles( + val manifestMissing = + if (branchAndManifestFile.isLiveSnapshot) new AtomicBoolean(false) else null + val entries = retryReadingFiles( () => manifestFile.readWithIOException(branchAndManifestFile.manifestName), - Collections.emptyList[ManifestEntry] - ).asScala.flatMap { - manifestEntry => - manifestEntry.fileName() +: manifestEntry.file().extraFiles().asScala + Collections.emptyList[ManifestEntry], + manifestMissing + ).asScala + if (manifestMissing != null && manifestMissing.get()) { + Iterator.single(("", true)) + } else { + entries + .flatMap { + manifestEntry => + manifestEntry.fileName() +: manifestEntry.file().extraFiles().asScala + } + .map(name => (name, false)) + .iterator } } } + .cache() + cachedDatasets += dataFilesWithFlag + + if (!dataFilesWithFlag.filter(_._2).isEmpty) { + logWarning("Detected missing manifest while collecting data files, aborting clean.") + return ( + spark.createDataset( + Seq((deletedFilesCountInLocal.get(), deletedFilesLenInBytesInLocal.get()))), + cachedDatasets.toSeq) + } + + val dataFiles = dataFilesWithFlag.map { + case (name, isMissing) => + if (isMissing) { + throw new RuntimeException( + "Detected missing live manifest while collecting data files, aborting clean.") + } + name + } // union manifest and data files val usedFiles = usedManifestFiles - .map(_.manifestName) + .map { + file => + if (file.isMissing) { + throw new RuntimeException( + "Detected missing live manifest during used-files collection, aborting clean.") + } + file.manifestName + } .union(dataFiles) .toDF("used_name") @@ -181,7 +259,7 @@ case class SparkOrphanFilesClean( deleted } - (finalDeletedDataset, usedManifestFiles) + (finalDeletedDataset, cachedDatasets.toSeq) } } @@ -193,7 +271,12 @@ case class SparkOrphanFilesClean( * @param isManifestFile * If it is the manifest file */ -case class BranchAndManifestFile(branch: String, manifestName: String, isManifestFile: Boolean) +case class BranchAndManifestFile( + branch: String, + manifestName: String, + isManifestFile: Boolean, + isMissing: Boolean = false, + isLiveSnapshot: Boolean = false) object SparkOrphanFilesClean extends SQLConfHelper { def executeDatabaseOrphanFiles( @@ -227,17 +310,22 @@ object SparkOrphanFilesClean extends SQLConfHelper { if (tables.isEmpty) { return new CleanOrphanFilesResult(0, 0) } - val (deleted, waitToRelease) = tables.map { - table => - new SparkOrphanFilesClean( - table, - olderThanMillis, - parallelism, - dryRun, - spark - ).doOrphanClean() - }.unzip + val deleted = new ArrayBuffer[Dataset[(Long, Long)]]() + val waitToRelease = new ArrayBuffer[Dataset[_]]() try { + tables.foreach { + table => + val (deletedDataset, cachedDatasets) = new SparkOrphanFilesClean( + table, + olderThanMillis, + parallelism, + dryRun, + spark + ).doOrphanClean() + deleted += deletedDataset + waitToRelease ++= cachedDatasets + } + val result = deleted .reduce((l, r) => l.union(r)) .toDF("deletedFilesCount", "deletedFilesLenInBytes") diff --git a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/RemoveOrphanFilesProcedureTest.scala b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/RemoveOrphanFilesProcedureTest.scala index af82549738dd..5efbe7d3bc2b 100644 --- a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/RemoveOrphanFilesProcedureTest.scala +++ b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/RemoveOrphanFilesProcedureTest.scala @@ -24,6 +24,8 @@ import org.apache.paimon.utils.DateTimeUtils import org.apache.spark.sql.Row +import java.nio.file.{Files, Paths} +import java.nio.file.attribute.FileTime import java.util.concurrent.TimeUnit class RemoveOrphanFilesProcedureTest extends PaimonSparkTestBase { @@ -311,4 +313,55 @@ class RemoveOrphanFilesProcedureTest extends PaimonSparkTestBase { .contains(partitionValue)) } + test("Paimon procedure: abort when cached references are recomputed after expiration") { + spark.sql(s""" + |CREATE TABLE T (id INT, name STRING) + |USING PAIMON + |TBLPROPERTIES ( + | 'bucket' = '1', + | 'bucket-key' = 'id', + | 'manifest.merge-min-count' = '1') + |""".stripMargin) + + spark.sql("INSERT INTO T VALUES (1, 'a')") + spark.sql("INSERT INTO T VALUES (2, 'b')") + + val oldDataFiles = + spark.sql("SELECT file_path FROM `T$files`").collect().map(_.getString(0)).toSeq + val oldTime = System.currentTimeMillis() - TimeUnit.DAYS.toMillis(2) + oldDataFiles.foreach( + path => Files.setLastModifiedTime(Paths.get(path), FileTime.fromMillis(oldTime))) + + val (deleted, cachedDatasets) = new SparkOrphanFilesClean( + loadTable("T"), + System.currentTimeMillis() - TimeUnit.DAYS.toMillis(1), + 1, + dryRunPara = false, + spark + ).doOrphanClean() + + try { + deleted.queryExecution.executedPlan + spark.sql("INSERT INTO T VALUES (3, 'c')") + spark.sql("CALL sys.expire_snapshots(table => 'T', retain_max => 1, retain_min => 1)") + spark.sparkContext.getPersistentRDDs.values.foreach(_.unpersist(blocking = true)) + + val exception = intercept[org.apache.spark.SparkException] { + deleted.collect() + } + assert( + Iterator + .iterate[Throwable](exception)(_.getCause) + .takeWhile(_ != null) + .exists( + cause => Option(cause.getMessage).exists(_.contains("Detected missing live manifest")))) + assert(oldDataFiles.forall(path => Files.exists(Paths.get(path)))) + checkAnswer( + spark.sql("SELECT * FROM T ORDER BY id"), + Seq(Row(1, "a"), Row(2, "b"), Row(3, "c"))) + } finally { + cachedDatasets.foreach(_.unpersist()) + } + } + }