From fef5ed23ea8758e62ca8e42d6f698c650c5b1dd5 Mon Sep 17 00:00:00 2001 From: yehe Date: Mon, 20 Jul 2026 17:35:13 +0800 Subject: [PATCH 01/11] fix-orphan-clean-manifest-missing --- .../operation/LocalOrphanFilesClean.java | 63 +++++++++- .../paimon/operation/OrphanFilesClean.java | 112 ++++++++++++++---- .../flink/orphan/FlinkOrphanFilesClean.java | 43 ++++++- .../procedure/SparkOrphanFilesClean.scala | 73 ++++++++++-- 4 files changed, 245 insertions(+), 46 deletions(-) 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..df693d15e232 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 @@ -47,6 +47,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 +109,43 @@ public CleanOrphanFilesResult clean() } candidateDeletes = new HashSet<>(candidates.keySet()); + AtomicLong missingManifestListCount = new AtomicLong(0); + AtomicLong missingManifestCount = new AtomicLong(0); + // find used files Set usedFiles = branches.stream() - .flatMap(branch -> getUsedFiles(branch).stream()) + .flatMap( + branch -> + getUsedFiles( + branch, + missingManifestListCount, + missingManifestCount) + .stream()) .collect(Collectors.toSet()); + if (usedFiles.isEmpty()) { + LOG.warn( + "Collected used files is empty while {} candidates are present, " + + "aborting orphan files clean to prevent data loss.", + candidateDeletes.size()); + return new CleanOrphanFilesResult( + deleteFiles.size(), deletedFilesLenInBytes.get(), deleteFiles); + } + + if (missingManifestListCount.get() > 0 || missingManifestCount.get() > 0) { + LOG.warn( + "Detected {} missing manifest-list/index-manifest and {} missing manifest " + + "file(s) during used-files collection while {} candidates are " + + "present; this indicates a concurrent snapshot expiration. " + + "Aborting orphan files clean to prevent data loss.", + missingManifestListCount.get(), + missingManifestCount.get(), + candidateDeletes.size()); + return new CleanOrphanFilesResult( + deleteFiles.size(), deletedFilesLenInBytes.get(), deleteFiles); + } + // delete unused files candidateDeletes.removeAll(usedFiles); candidateDeletes.stream() @@ -157,14 +189,21 @@ private void cleanEmptyDataDirectory(List deleteFiles) { } private void collectWithoutDataFile( - String branch, Consumer usedFileConsumer, Consumer manifestConsumer) + String branch, + Consumer usedFileConsumer, + Consumer manifestConsumer, + Consumer missingManifestListConsumer) throws IOException { randomlyOnlyExecute( executor, snapshot -> { try { collectWithoutDataFile( - branch, snapshot, usedFileConsumer, manifestConsumer); + branch, + snapshot, + usedFileConsumer, + manifestConsumer, + missingManifestListConsumer); } catch (IOException e) { throw new RuntimeException(e); } @@ -172,20 +211,29 @@ private void collectWithoutDataFile( safelyGetAllSnapshots(branch)); } - private Set getUsedFiles(String branch) { + private Set getUsedFiles( + String branch, + AtomicLong missingManifestListCount, + AtomicLong missingManifestCount) { Set usedFiles = ConcurrentHashMap.newKeySet(); ManifestFile manifestFile = table.switchToBranch(branch).store().manifestFileFactory().create(); try { Set manifests = ConcurrentHashMap.newKeySet(); - collectWithoutDataFile(branch, usedFiles::add, manifests::add); + collectWithoutDataFile( + branch, + usedFiles::add, + manifests::add, + name -> missingManifestListCount.incrementAndGet()); randomlyOnlyExecute( executor, manifestName -> { try { + AtomicBoolean manifestMissing = new AtomicBoolean(false); retryReadingFiles( () -> manifestFile.readWithIOException(manifestName), - Collections.emptyList()) + Collections.emptyList(), + manifestMissing) .stream() .map(ManifestEntry::file) .forEach( @@ -197,6 +245,9 @@ private Set getUsedFiles(String branch) { .filter(candidateDeletes::contains) .forEach(usedFiles::add); }); + if (manifestMissing.get()) { + missingManifestCount.incrementAndGet(); + } } catch (IOException e) { throw new RuntimeException(e); } 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..6008282144fc 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 @@ -56,6 +56,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; @@ -260,6 +261,16 @@ protected void collectWithoutDataFile( Consumer usedFileConsumer, Consumer manifestConsumer) throws IOException { + collectWithoutDataFile(branch, snapshot, usedFileConsumer, manifestConsumer, name -> {}); + } + + protected void collectWithoutDataFile( + String branch, + Snapshot snapshot, + Consumer usedFileConsumer, + Consumer manifestConsumer, + Consumer missingManifestListConsumer) + throws IOException { Consumer> usedFileWithFlagConsumer = fileAndFlag -> { if (fileAndFlag.getRight()) { @@ -267,7 +278,8 @@ protected void collectWithoutDataFile( } usedFileConsumer.accept(fileAndFlag.getLeft()); }; - collectWithoutDataFileWithManifestFlag(branch, snapshot, usedFileWithFlagConsumer); + collectWithoutDataFileWithManifestFlag( + branch, snapshot, usedFileWithFlagConsumer, missingManifestListConsumer); } protected void collectWithoutDataFileWithManifestFlag( @@ -275,36 +287,47 @@ protected void collectWithoutDataFileWithManifestFlag( Snapshot snapshot, Consumer> usedFileWithFlagConsumer) throws IOException { + collectWithoutDataFileWithManifestFlag( + branch, snapshot, usedFileWithFlagConsumer, name -> {}); + } + + protected void collectWithoutDataFileWithManifestFlag( + String branch, + Snapshot snapshot, + Consumer> usedFileWithFlagConsumer, + Consumer missingManifestListConsumer) + throws IOException { FileStoreTable branchTable = table.switchToBranch(branch); ManifestList manifestList = branchTable.store().manifestListFactory().create(); IndexFileHandler indexFileHandler = branchTable.store().newIndexFileHandler(); List manifestFileMetas = new ArrayList<>(); // changelog manifest if (snapshot.changelogManifestList() != null) { - usedFileWithFlagConsumer.accept(Pair.of(snapshot.changelogManifestList(), false)); - manifestFileMetas.addAll( - retryReadingFiles( - () -> - manifestList.readWithIOException( - snapshot.changelogManifestList()), - emptyList())); + collectManifestList( + manifestList, + snapshot.changelogManifestList(), + manifestFileMetas, + usedFileWithFlagConsumer, + missingManifestListConsumer); } // delta manifest if (snapshot.deltaManifestList() != null) { - usedFileWithFlagConsumer.accept(Pair.of(snapshot.deltaManifestList(), false)); - manifestFileMetas.addAll( - retryReadingFiles( - () -> manifestList.readWithIOException(snapshot.deltaManifestList()), - emptyList())); + collectManifestList( + manifestList, + snapshot.deltaManifestList(), + manifestFileMetas, + usedFileWithFlagConsumer, + missingManifestListConsumer); } // base manifest - usedFileWithFlagConsumer.accept(Pair.of(snapshot.baseManifestList(), false)); - manifestFileMetas.addAll( - retryReadingFiles( - () -> manifestList.readWithIOException(snapshot.baseManifestList()), - emptyList())); + collectManifestList( + manifestList, + snapshot.baseManifestList(), + manifestFileMetas, + usedFileWithFlagConsumer, + missingManifestListConsumer); // collect manifests for (ManifestFileMeta manifest : manifestFileMetas) { @@ -314,14 +337,21 @@ protected void collectWithoutDataFileWithManifestFlag( // index files String indexManifest = snapshot.indexManifest(); if (indexManifest != null && indexFileHandler.existsManifest(indexManifest)) { - usedFileWithFlagConsumer.accept(Pair.of(indexManifest, false)); - retryReadingFiles( + AtomicBoolean indexManifestMissing = new AtomicBoolean(false); + 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(), + indexManifestMissing); + if (indexManifestMissing.get()) { + missingManifestListConsumer.accept(indexManifest); + } else { + usedFileWithFlagConsumer.accept(Pair.of(indexManifest, false)); + indexEntries.stream() + .map(IndexManifestEntry::indexFile) + .map(IndexFileMeta::fileName) + .forEach(name -> usedFileWithFlagConsumer.accept(Pair.of(name, false))); + } } // statistic file @@ -330,6 +360,27 @@ protected void collectWithoutDataFileWithManifestFlag( } } + private static void collectManifestList( + ManifestList manifestList, + String manifestListName, + List manifestFileMetas, + Consumer> usedFileWithFlagConsumer, + Consumer missingManifestListConsumer) + throws IOException { + AtomicBoolean missing = new AtomicBoolean(false); + List metas = + retryReadingFiles( + () -> manifestList.readWithIOException(manifestListName), + emptyList(), + missing); + if (missing.get()) { + missingManifestListConsumer.accept(manifestListName); + return; + } + usedFileWithFlagConsumer.accept(Pair.of(manifestListName, false)); + manifestFileMetas.addAll(metas); + } + /** List directories that contains data files and manifest files. */ protected List listPaimonFileDirs() { FileStorePathFactory pathFactory = table.store().pathFactory(); @@ -444,12 +495,23 @@ 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); + } return defaultValue; } catch (IOException e) { caught = e; 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..3ad187d72c31 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 @@ -63,6 +63,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 +76,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( @@ -183,8 +187,14 @@ public void processElement( new Tuple2<>(branch, manifest); ctx.output(manifestOutputTag, tuple2); }; + Consumer missingManifestListConsumer = + name -> out.collect(MISSING_MANIFEST_SENTINEL); collectWithoutDataFile( - branch, snapshot, out::collect, manifestConsumer); + branch, + snapshot, + out::collect, + manifestConsumer, + missingManifestListConsumer); } }) .name("collect-manifests"); @@ -219,12 +229,15 @@ public void endInput() throws IOException { .store() .manifestFileFactory() .create()); + AtomicBoolean manifestMissing = + new AtomicBoolean(false); retryReadingFiles( () -> manifestFile .readWithIOException( tuple2.f1), - Collections.emptyList()) + Collections.emptyList(), + manifestMissing) .forEach( f -> { List files = @@ -237,6 +250,11 @@ public void endInput() throws IOException { new StreamRecord<>( file))); }); + if (manifestMissing.get()) { + output.collect( + new StreamRecord<>( + MISSING_MANIFEST_SENTINEL)); + } } } }); @@ -369,6 +387,8 @@ public void endInput() throws IOException { private long emittedFilesCount; private long emittedFilesLen; + private boolean abortDeletion; + private final Set used = new HashSet<>(); @Override @@ -385,6 +405,15 @@ public void endInput(int inputId) { checkState(!buildEnd, "Should not build ended."); LOG.info("Finish build phase."); buildEnd = true; + if (abortDeletion) { + LOG.warn( + "Detected missing manifest-list/manifest " + + "during used-files collection; " + + "this indicates a concurrent " + + "snapshot expiration. Aborting " + + "orphan files clean to prevent " + + "data loss."); + } break; case 2: checkState(buildEnd, "Should build ended."); @@ -404,13 +433,21 @@ public void endInput(int inputId) { @Override public void processElement1(StreamRecord element) { - used.add(element.getValue()); + String value = element.getValue(); + if (MISSING_MANIFEST_SENTINEL.equals(value)) { + abortDeletion = true; + return; + } + used.add(value); } @Override public void processElement2( StreamRecord> element) { checkState(buildEnd, "Should build ended."); + if (abortDeletion) { + return; + } Tuple2 fileInfo = element.getValue(); String value = fileInfo.f0; Path path = new Path(value); 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..02bf667d4319 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,7 +34,7 @@ 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._ @@ -78,19 +78,40 @@ case class SparkOrphanFilesClean( 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)) } } + val missingConsumer = new Consumer[String] { + override def accept(name: String): Unit = { + usedFileBuffer.append( + BranchAndManifestFile(branch, name, isManifestFile = false, isMissing = true)) + } + } val snapshot = Snapshot.fromJson(snapshotJson) - collectWithoutDataFileWithManifestFlag(branch, snapshot, usedFileConsumer) + collectWithoutDataFileWithManifestFlag( + branch, + snapshot, + usedFileConsumer, + missingConsumer) usedFileBuffer } .toDS() .cache() - // find all data files - val dataFiles = usedManifestFiles - .filter(_.isManifestFile) + if (!usedManifestFiles.filter(_.isMissing).isEmpty) { + logWarning( + "Detected missing manifest-list/index-manifest during used-files collection; this " + + "indicates a concurrent snapshot expiration. Aborting orphan files clean to prevent " + + "data loss.") + return ( + spark.createDataset( + Seq((deletedFilesCountInLocal.get(), deletedFilesLenInBytesInLocal.get()))), + usedManifestFiles) + } + + val dataFilesWithFlag = usedManifestFiles + .filter(f => f.isManifestFile && !f.isMissing) .distinct() .mapPartitions { it => @@ -102,18 +123,42 @@ case class SparkOrphanFilesClean( (key: String) => specifiedTable.switchToBranch(key).store.manifestFileFactory.create) - retryReadingFiles( + val manifestMissing = new AtomicBoolean(false) + 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.get()) { + Iterator.single(("", true)) + } else { + entries + .flatMap { + manifestEntry => + manifestEntry.fileName() +: manifestEntry.file().extraFiles().asScala + } + .map(name => (name, false)) + .iterator } } } + .cache() + + if (!dataFilesWithFlag.filter(_._2).isEmpty) { + logWarning( + "Detected missing manifest file(s) while collecting data files; this indicates a " + + "concurrent snapshot expiration. Aborting orphan files clean to prevent data loss.") + return ( + spark.createDataset( + Seq((deletedFilesCountInLocal.get(), deletedFilesLenInBytesInLocal.get()))), + usedManifestFiles) + } + + val dataFiles = dataFilesWithFlag.filter(!_._2).map(_._1) // union manifest and data files val usedFiles = usedManifestFiles + .filter(!_.isMissing) .map(_.manifestName) .union(dataFiles) .toDF("used_name") @@ -193,7 +238,11 @@ 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) object SparkOrphanFilesClean extends SQLConfHelper { def executeDatabaseOrphanFiles( From 489eb6287dcc5801a2d8d9b3e20962ff015b441c Mon Sep 17 00:00:00 2001 From: yehe Date: Tue, 21 Jul 2026 18:12:13 +0800 Subject: [PATCH 02/11] [core] Simplify missing-manifest handling in orphan files clean - Unify missing-manifest signaling to a single AtomicBoolean across Local/Flink/Spark clean implementations, removing the Consumer callback and its unused file-name argument - Move FileNotFoundException diagnostic logging into retryReadingFiles so all metadata reads are covered in one place - Inline collectManifestList into a loop and short-circuit return on the first missing manifest list --- .../operation/LocalOrphanFilesClean.java | 48 ++------- .../paimon/operation/OrphanFilesClean.java | 99 +++++-------------- .../flink/orphan/FlinkOrphanFilesClean.java | 15 ++- .../procedure/SparkOrphanFilesClean.scala | 22 ++--- 4 files changed, 51 insertions(+), 133 deletions(-) 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 df693d15e232..c585d9f2cb86 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 @@ -109,39 +109,22 @@ public CleanOrphanFilesResult clean() } candidateDeletes = new HashSet<>(candidates.keySet()); - AtomicLong missingManifestListCount = new AtomicLong(0); - AtomicLong missingManifestCount = new AtomicLong(0); + AtomicBoolean missingManifest = new AtomicBoolean(false); // find used files Set usedFiles = branches.stream() - .flatMap( - branch -> - getUsedFiles( - branch, - missingManifestListCount, - missingManifestCount) - .stream()) + .flatMap(branch -> getUsedFiles(branch, missingManifest).stream()) .collect(Collectors.toSet()); if (usedFiles.isEmpty()) { - LOG.warn( - "Collected used files is empty while {} candidates are present, " - + "aborting orphan files clean to prevent data loss.", - candidateDeletes.size()); + LOG.warn("Collected used files is empty, aborting orphan files clean."); return new CleanOrphanFilesResult( deleteFiles.size(), deletedFilesLenInBytes.get(), deleteFiles); } - if (missingManifestListCount.get() > 0 || missingManifestCount.get() > 0) { - LOG.warn( - "Detected {} missing manifest-list/index-manifest and {} missing manifest " - + "file(s) during used-files collection while {} candidates are " - + "present; this indicates a concurrent snapshot expiration. " - + "Aborting orphan files clean to prevent data loss.", - missingManifestListCount.get(), - missingManifestCount.get(), - candidateDeletes.size()); + if (missingManifest.get()) { + LOG.warn("Detected missing manifest during used-files collection, aborting clean."); return new CleanOrphanFilesResult( deleteFiles.size(), deletedFilesLenInBytes.get(), deleteFiles); } @@ -192,7 +175,7 @@ private void collectWithoutDataFile( String branch, Consumer usedFileConsumer, Consumer manifestConsumer, - Consumer missingManifestListConsumer) + AtomicBoolean missingManifest) throws IOException { randomlyOnlyExecute( executor, @@ -203,7 +186,7 @@ private void collectWithoutDataFile( snapshot, usedFileConsumer, manifestConsumer, - missingManifestListConsumer); + missingManifest); } catch (IOException e) { throw new RuntimeException(e); } @@ -211,29 +194,21 @@ private void collectWithoutDataFile( safelyGetAllSnapshots(branch)); } - private Set getUsedFiles( - String branch, - AtomicLong missingManifestListCount, - AtomicLong missingManifestCount) { + 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, - name -> missingManifestListCount.incrementAndGet()); + collectWithoutDataFile(branch, usedFiles::add, manifests::add, missingManifest); randomlyOnlyExecute( executor, manifestName -> { try { - AtomicBoolean manifestMissing = new AtomicBoolean(false); retryReadingFiles( () -> manifestFile.readWithIOException(manifestName), Collections.emptyList(), - manifestMissing) + missingManifest) .stream() .map(ManifestEntry::file) .forEach( @@ -245,9 +220,6 @@ private Set getUsedFiles( .filter(candidateDeletes::contains) .forEach(usedFiles::add); }); - if (manifestMissing.get()) { - missingManifestCount.incrementAndGet(); - } } catch (IOException e) { throw new RuntimeException(e); } 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 6008282144fc..51adf8199582 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 @@ -255,21 +255,12 @@ protected Set safelyGetAllSnapshots(String branch) throws IOException return readSnapshots; } - protected void collectWithoutDataFile( - String branch, - Snapshot snapshot, - Consumer usedFileConsumer, - Consumer manifestConsumer) - throws IOException { - collectWithoutDataFile(branch, snapshot, usedFileConsumer, manifestConsumer, name -> {}); - } - protected void collectWithoutDataFile( String branch, Snapshot snapshot, Consumer usedFileConsumer, Consumer manifestConsumer, - Consumer missingManifestListConsumer) + AtomicBoolean missingManifest) throws IOException { Consumer> usedFileWithFlagConsumer = fileAndFlag -> { @@ -279,55 +270,42 @@ protected void collectWithoutDataFile( usedFileConsumer.accept(fileAndFlag.getLeft()); }; collectWithoutDataFileWithManifestFlag( - branch, snapshot, usedFileWithFlagConsumer, missingManifestListConsumer); - } - - protected void collectWithoutDataFileWithManifestFlag( - String branch, - Snapshot snapshot, - Consumer> usedFileWithFlagConsumer) - throws IOException { - collectWithoutDataFileWithManifestFlag( - branch, snapshot, usedFileWithFlagConsumer, name -> {}); + branch, snapshot, usedFileWithFlagConsumer, missingManifest); } protected void collectWithoutDataFileWithManifestFlag( String branch, Snapshot snapshot, Consumer> usedFileWithFlagConsumer, - Consumer missingManifestListConsumer) + 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) { - collectManifestList( - manifestList, - snapshot.changelogManifestList(), - manifestFileMetas, - usedFileWithFlagConsumer, - missingManifestListConsumer); + manifestListNames.add(snapshot.changelogManifestList()); } - - // delta manifest if (snapshot.deltaManifestList() != null) { - collectManifestList( - manifestList, - snapshot.deltaManifestList(), - manifestFileMetas, - usedFileWithFlagConsumer, - missingManifestListConsumer); + manifestListNames.add(snapshot.deltaManifestList()); } + manifestListNames.add(snapshot.baseManifestList()); - // base manifest - collectManifestList( - manifestList, - snapshot.baseManifestList(), - manifestFileMetas, - usedFileWithFlagConsumer, - missingManifestListConsumer); + 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.get()) { + return; + } + usedFileWithFlagConsumer.accept(Pair.of(manifestListName, false)); + manifestFileMetas.addAll(metas); + } // collect manifests for (ManifestFileMeta manifest : manifestFileMetas) { @@ -337,15 +315,12 @@ protected void collectWithoutDataFileWithManifestFlag( // index files String indexManifest = snapshot.indexManifest(); if (indexManifest != null && indexFileHandler.existsManifest(indexManifest)) { - AtomicBoolean indexManifestMissing = new AtomicBoolean(false); List indexEntries = retryReadingFiles( () -> indexFileHandler.readManifestWithIOException(indexManifest), Collections.emptyList(), - indexManifestMissing); - if (indexManifestMissing.get()) { - missingManifestListConsumer.accept(indexManifest); - } else { + missingManifest); + if (!missingManifest.get()) { usedFileWithFlagConsumer.accept(Pair.of(indexManifest, false)); indexEntries.stream() .map(IndexManifestEntry::indexFile) @@ -360,27 +335,6 @@ protected void collectWithoutDataFileWithManifestFlag( } } - private static void collectManifestList( - ManifestList manifestList, - String manifestListName, - List manifestFileMetas, - Consumer> usedFileWithFlagConsumer, - Consumer missingManifestListConsumer) - throws IOException { - AtomicBoolean missing = new AtomicBoolean(false); - List metas = - retryReadingFiles( - () -> manifestList.readWithIOException(manifestListName), - emptyList(), - missing); - if (missing.get()) { - missingManifestListConsumer.accept(manifestListName); - return; - } - usedFileWithFlagConsumer.accept(Pair.of(manifestListName, false)); - manifestFileMetas.addAll(metas); - } - /** List directories that contains data files and manifest files. */ protected List listPaimonFileDirs() { FileStorePathFactory pathFactory = table.store().pathFactory(); @@ -499,9 +453,7 @@ protected static T retryReadingFiles(SupplierWithIOException reader, T de } protected static T retryReadingFiles( - SupplierWithIOException reader, - T defaultValue, - @Nullable AtomicBoolean fnfFallback) + SupplierWithIOException reader, T defaultValue, @Nullable AtomicBoolean fnfFallback) throws IOException { int retryNumber = 0; IOException caught = null; @@ -511,6 +463,9 @@ protected static T retryReadingFiles( } 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) { 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 3ad187d72c31..408cdef12f1d 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 @@ -187,14 +187,16 @@ public void processElement( new Tuple2<>(branch, manifest); ctx.output(manifestOutputTag, tuple2); }; - Consumer missingManifestListConsumer = - name -> out.collect(MISSING_MANIFEST_SENTINEL); + AtomicBoolean missingManifest = new AtomicBoolean(false); collectWithoutDataFile( branch, snapshot, out::collect, manifestConsumer, - missingManifestListConsumer); + missingManifest); + if (missingManifest.get()) { + out.collect(MISSING_MANIFEST_SENTINEL); + } } }) .name("collect-manifests"); @@ -407,12 +409,7 @@ public void endInput(int inputId) { buildEnd = true; if (abortDeletion) { LOG.warn( - "Detected missing manifest-list/manifest " - + "during used-files collection; " - + "this indicates a concurrent " - + "snapshot expiration. Aborting " - + "orphan files clean to prevent " - + "data loss."); + "Detected missing manifest, aborting clean."); } break; case 2: 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 02bf667d4319..c4af89de722b 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 @@ -82,28 +82,24 @@ case class SparkOrphanFilesClean( BranchAndManifestFile(branch, pair.getLeft, pair.getRight, isMissing = false)) } } - val missingConsumer = new Consumer[String] { - override def accept(name: String): Unit = { - usedFileBuffer.append( - BranchAndManifestFile(branch, name, isManifestFile = false, isMissing = true)) - } - } + val missingManifest = new AtomicBoolean(false) val snapshot = Snapshot.fromJson(snapshotJson) collectWithoutDataFileWithManifestFlag( branch, snapshot, usedFileConsumer, - missingConsumer) + missingManifest) + if (missingManifest.get()) { + usedFileBuffer.append( + BranchAndManifestFile(branch, "", isManifestFile = false, isMissing = true)) + } usedFileBuffer } .toDS() .cache() if (!usedManifestFiles.filter(_.isMissing).isEmpty) { - logWarning( - "Detected missing manifest-list/index-manifest during used-files collection; this " + - "indicates a concurrent snapshot expiration. Aborting orphan files clean to prevent " + - "data loss.") + logWarning("Detected missing manifest during used-files collection, aborting clean.") return ( spark.createDataset( Seq((deletedFilesCountInLocal.get(), deletedFilesLenInBytesInLocal.get()))), @@ -145,9 +141,7 @@ case class SparkOrphanFilesClean( .cache() if (!dataFilesWithFlag.filter(_._2).isEmpty) { - logWarning( - "Detected missing manifest file(s) while collecting data files; this indicates a " + - "concurrent snapshot expiration. Aborting orphan files clean to prevent data loss.") + logWarning("Detected missing manifest while collecting data files, aborting clean.") return ( spark.createDataset( Seq((deletedFilesCountInLocal.get(), deletedFilesLenInBytesInLocal.get()))), From 0ecc4f342addf46e01ccfe3ef273bf759f0ebd0d Mon Sep 17 00:00:00 2001 From: yehe Date: Tue, 21 Jul 2026 20:01:17 +0800 Subject: [PATCH 03/11] [core] Add test reproducing orphan clean data-file deletion race on missing manifest-list --- .../operation/LocalOrphanFilesCleanTest.java | 40 +++++++++++++++++++ 1 file changed, 40 insertions(+) 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..b8ab44881523 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 @@ -704,6 +704,46 @@ 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(); + } + } + + 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, From 1c9683ecc1bf7cb544bbf9fead4c108826f44ceb Mon Sep 17 00:00:00 2001 From: yehe Date: Tue, 21 Jul 2026 23:43:32 +0800 Subject: [PATCH 04/11] [core] Abort orphan files clean when a live main-branch snapshot's manifest is missing --- .../operation/LocalOrphanFilesClean.java | 30 ++++++- .../paimon/operation/OrphanFilesClean.java | 8 +- .../operation/LocalOrphanFilesCleanTest.java | 11 ++- .../flink/orphan/FlinkOrphanFilesClean.java | 81 ++++++++++++------- .../procedure/SparkOrphanFilesClean.scala | 41 +++++++--- 5 files changed, 126 insertions(+), 45 deletions(-) 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 c585d9f2cb86..b981ee8b8671 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; @@ -53,6 +54,7 @@ import java.util.function.Function; import java.util.stream.Collectors; +import static org.apache.paimon.catalog.Identifier.DEFAULT_MAIN_BRANCH; import static org.apache.paimon.utils.FileStorePathFactory.BUCKET_PATH_PREFIX; import static org.apache.paimon.utils.Preconditions.checkArgument; import static org.apache.paimon.utils.ThreadPoolUtils.createCachedThreadPool; @@ -175,18 +177,34 @@ private void collectWithoutDataFile( String branch, Consumer usedFileConsumer, Consumer manifestConsumer, + Consumer liveManifestConsumer, AtomicBoolean missingManifest) throws IOException { + Set liveSnapshots = + DEFAULT_MAIN_BRANCH.equals(branch) + ? new HashSet<>( + table.switchToBranch(branch) + .snapshotManager() + .safelyGetAllSnapshots()) + : Collections.emptySet(); 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, - missingManifest); + perSnapshotManifestConsumer, + live ? missingManifest : null); } catch (IOException e) { throw new RuntimeException(e); } @@ -200,15 +218,19 @@ private Set getUsedFiles(String branch, AtomicBoolean missingManifest) { table.switchToBranch(branch).store().manifestFileFactory().create(); try { Set manifests = ConcurrentHashMap.newKeySet(); - collectWithoutDataFile(branch, usedFiles::add, manifests::add, missingManifest); + 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(), - missingManifest) + 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 51adf8199582..35b131877ad2 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 @@ -260,7 +260,7 @@ protected void collectWithoutDataFile( Snapshot snapshot, Consumer usedFileConsumer, Consumer manifestConsumer, - AtomicBoolean missingManifest) + @Nullable AtomicBoolean missingManifest) throws IOException { Consumer> usedFileWithFlagConsumer = fileAndFlag -> { @@ -277,7 +277,7 @@ protected void collectWithoutDataFileWithManifestFlag( String branch, Snapshot snapshot, Consumer> usedFileWithFlagConsumer, - AtomicBoolean missingManifest) + @Nullable AtomicBoolean missingManifest) throws IOException { FileStoreTable branchTable = table.switchToBranch(branch); ManifestList manifestList = branchTable.store().manifestListFactory().create(); @@ -300,7 +300,7 @@ protected void collectWithoutDataFileWithManifestFlag( emptyList(), missingManifest); // A missing manifest list means we cannot determine all used files, so abort. - if (missingManifest.get()) { + if (missingManifest != null && missingManifest.get()) { return; } usedFileWithFlagConsumer.accept(Pair.of(manifestListName, false)); @@ -320,7 +320,7 @@ protected void collectWithoutDataFileWithManifestFlag( () -> indexFileHandler.readManifestWithIOException(indexManifest), Collections.emptyList(), missingManifest); - if (!missingManifest.get()) { + if (missingManifest == null || !missingManifest.get()) { usedFileWithFlagConsumer.accept(Pair.of(indexManifest, false)); indexEntries.stream() .map(IndexManifestEntry::indexFile) 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 b8ab44881523..c6659eec59a2 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 @@ -550,6 +550,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 +571,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 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 408cdef12f1d..48a3634f9b7b 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; @@ -66,6 +67,7 @@ import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicLong; import java.util.function.Consumer; +import java.util.stream.Collectors; import static org.apache.flink.api.common.typeinfo.BasicTypeInfo.STRING_TYPE_INFO; import static org.apache.flink.util.Preconditions.checkState; @@ -147,54 +149,73 @@ 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 { + Set liveSnapshotIds = + Identifier.DEFAULT_MAIN_BRANCH.equals(branch) + ? table.switchToBranch(branch) + .snapshotManager() + .safelyGetAllSnapshots().stream() + .map(Snapshot::id) + .collect(Collectors.toSet()) + : Collections.emptySet(); for (Snapshot snapshot : safelyGetAllSnapshots(branch)) { - out.collect(new Tuple2<>(branch, snapshot.toJson())); + out.collect( + new Tuple3<>( + branch, + snapshot.toJson(), + liveSnapshotIds.contains( + snapshot.id()))); } } }) .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 liveOnMainBranch = branchAndSnapshot.f2; Consumer manifestConsumer = - manifest -> { - Tuple2 tuple2 = - new Tuple2<>(branch, manifest); - ctx.output(manifestOutputTag, tuple2); - }; - AtomicBoolean missingManifest = new AtomicBoolean(false); + manifest -> + ctx.output( + manifestOutputTag, + new Tuple3<>( + branch, + manifest, + liveOnMainBranch)); + AtomicBoolean missingManifest = + liveOnMainBranch ? new AtomicBoolean(false) : null; collectWithoutDataFile( branch, snapshot, out::collect, manifestConsumer, missingManifest); - if (missingManifest.get()) { + if (missingManifest != null && missingManifest.get()) { out.collect(MISSING_MANIFEST_SENTINEL); } } @@ -204,40 +225,46 @@ 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 = - new AtomicBoolean(false); + tuple.f2 ? new AtomicBoolean(false) : null; retryReadingFiles( () -> manifestFile .readWithIOException( - tuple2.f1), + tuple.f1), Collections.emptyList(), manifestMissing) .forEach( @@ -252,7 +279,7 @@ public void endInput() throws IOException { new StreamRecord<>( file))); }); - if (manifestMissing.get()) { + if (manifestMissing != null && manifestMissing.get()) { output.collect( new StreamRecord<>( MISSING_MANIFEST_SENTINEL)); 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 c4af89de722b..9a4d2fb12a93 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 @@ -69,27 +69,46 @@ 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 liveSnapshotIds = + if (Identifier.DEFAULT_MAIN_BRANCH.equals(branch)) { + specifiedTable + .switchToBranch(branch) + .snapshotManager() + .safelyGetAllSnapshots() + .asScala + .map(_.id()) + .toSet + } else { + Set.empty[Long] + } + safelyGetAllSnapshots(branch).asScala.map( + snapshot => (branch, snapshot.toJson, liveSnapshotIds.contains(snapshot.id()))) }) .repartition(parallelism) .flatMap { - case (branch, snapshotJson) => + case (branch, snapshotJson, liveOnMainBranch) => 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, isMissing = false)) + BranchAndManifestFile( + branch, + pair.getLeft, + pair.getRight, + isMissing = false, + isLiveSnapshot = liveOnMainBranch)) } } - val missingManifest = new AtomicBoolean(false) + val missingManifest = if (liveOnMainBranch) new AtomicBoolean(false) else null val snapshot = Snapshot.fromJson(snapshotJson) collectWithoutDataFileWithManifestFlag( branch, snapshot, usedFileConsumer, missingManifest) - if (missingManifest.get()) { + if (missingManifest != null && missingManifest.get()) { usedFileBuffer.append( BranchAndManifestFile(branch, "", isManifestFile = false, isMissing = true)) } @@ -108,7 +127,9 @@ case class SparkOrphanFilesClean( val dataFilesWithFlag = usedManifestFiles .filter(f => f.isManifestFile && !f.isMissing) - .distinct() + .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] @@ -119,13 +140,14 @@ case class SparkOrphanFilesClean( (key: String) => specifiedTable.switchToBranch(key).store.manifestFileFactory.create) - val manifestMissing = new AtomicBoolean(false) + val manifestMissing = + if (branchAndManifestFile.isLiveSnapshot) new AtomicBoolean(false) else null val entries = retryReadingFiles( () => manifestFile.readWithIOException(branchAndManifestFile.manifestName), Collections.emptyList[ManifestEntry], manifestMissing ).asScala - if (manifestMissing.get()) { + if (manifestMissing != null && manifestMissing.get()) { Iterator.single(("", true)) } else { entries @@ -236,7 +258,8 @@ case class BranchAndManifestFile( branch: String, manifestName: String, isManifestFile: Boolean, - isMissing: Boolean = false) + isMissing: Boolean = false, + isLiveSnapshot: Boolean = false) object SparkOrphanFilesClean extends SQLConfHelper { def executeDatabaseOrphanFiles( From 3d5030e8806d1eefb4086a5cfc4b379c5871891f Mon Sep 17 00:00:00 2001 From: yehe Date: Tue, 28 Jul 2026 12:28:02 +0800 Subject: [PATCH 05/11] [core] Extend missing-manifest abort in orphan clean to all active branches Previously only main-branch live snapshots adopted the missing-manifest abort behavior, while other active branches passed a null missingManifest and treated read failures as empty results. Because branches have independent snapshots and can commit/expire concurrently, this could wrongly mark data files still referenced by a branch as orphans and delete them. Now every active branch collects its live snapshots and aborts on a missing manifest (local/Flink/Spark); tags remain best-effort. Adds a regression test for concurrent branch commit/expiration. --- .../operation/LocalOrphanFilesClean.java | 9 +-- .../operation/LocalOrphanFilesCleanTest.java | 56 +++++++++++++++++-- .../flink/orphan/FlinkOrphanFilesClean.java | 17 +++--- .../procedure/SparkOrphanFilesClean.scala | 24 ++++---- 4 files changed, 71 insertions(+), 35 deletions(-) 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 b981ee8b8671..9a91d7ec7c40 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 @@ -54,7 +54,6 @@ import java.util.function.Function; import java.util.stream.Collectors; -import static org.apache.paimon.catalog.Identifier.DEFAULT_MAIN_BRANCH; import static org.apache.paimon.utils.FileStorePathFactory.BUCKET_PATH_PREFIX; import static org.apache.paimon.utils.Preconditions.checkArgument; import static org.apache.paimon.utils.ThreadPoolUtils.createCachedThreadPool; @@ -181,12 +180,8 @@ private void collectWithoutDataFile( AtomicBoolean missingManifest) throws IOException { Set liveSnapshots = - DEFAULT_MAIN_BRANCH.equals(branch) - ? new HashSet<>( - table.switchToBranch(branch) - .snapshotManager() - .safelyGetAllSnapshots()) - : Collections.emptySet(); + new HashSet<>( + table.switchToBranch(branch).snapshotManager().safelyGetAllSnapshots()); randomlyOnlyExecute( executor, snapshot -> { 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 c6659eec59a2..06ef4923c374 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 @@ -212,8 +212,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 +320,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); } @@ -739,6 +749,44 @@ void testAbortWhenManifestListMissingDuringConcurrentExpiration() throws Excepti } } + @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(); + } + } + private void collectDataFiles(Path dir, List result) throws IOException { FileStatus[] statuses = fileIO.listStatus(dir); if (statuses == null) { 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 48a3634f9b7b..45e386136652 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 @@ -166,13 +166,10 @@ public void processElement( Collector> out) throws Exception { Set liveSnapshotIds = - Identifier.DEFAULT_MAIN_BRANCH.equals(branch) - ? table.switchToBranch(branch) - .snapshotManager() - .safelyGetAllSnapshots().stream() - .map(Snapshot::id) - .collect(Collectors.toSet()) - : Collections.emptySet(); + table.switchToBranch(branch).snapshotManager() + .safelyGetAllSnapshots().stream() + .map(Snapshot::id) + .collect(Collectors.toSet()); for (Snapshot snapshot : safelyGetAllSnapshots(branch)) { out.collect( new Tuple3<>( @@ -198,7 +195,7 @@ public void processElement( throws Exception { String branch = branchAndSnapshot.f0; Snapshot snapshot = Snapshot.fromJson(branchAndSnapshot.f1); - boolean liveOnMainBranch = branchAndSnapshot.f2; + boolean isLiveSnapshot = branchAndSnapshot.f2; Consumer manifestConsumer = manifest -> ctx.output( @@ -206,9 +203,9 @@ public void processElement( new Tuple3<>( branch, manifest, - liveOnMainBranch)); + isLiveSnapshot)); AtomicBoolean missingManifest = - liveOnMainBranch ? new AtomicBoolean(false) : null; + isLiveSnapshot ? new AtomicBoolean(false) : null; collectWithoutDataFile( branch, snapshot, 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 9a4d2fb12a93..dd822db716f5 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 @@ -71,23 +71,19 @@ case class SparkOrphanFilesClean( .mapPartitions(_.flatMap { branch => val liveSnapshotIds = - if (Identifier.DEFAULT_MAIN_BRANCH.equals(branch)) { - specifiedTable - .switchToBranch(branch) - .snapshotManager() - .safelyGetAllSnapshots() - .asScala - .map(_.id()) - .toSet - } else { - Set.empty[Long] - } + specifiedTable + .switchToBranch(branch) + .snapshotManager() + .safelyGetAllSnapshots() + .asScala + .map(_.id()) + .toSet safelyGetAllSnapshots(branch).asScala.map( snapshot => (branch, snapshot.toJson, liveSnapshotIds.contains(snapshot.id()))) }) .repartition(parallelism) .flatMap { - case (branch, snapshotJson, liveOnMainBranch) => + case (branch, snapshotJson, isLive) => val usedFileBuffer = new ArrayBuffer[BranchAndManifestFile]() val usedFileConsumer = new Consumer[org.apache.paimon.utils.Pair[String, java.lang.Boolean]] { @@ -98,10 +94,10 @@ case class SparkOrphanFilesClean( pair.getLeft, pair.getRight, isMissing = false, - isLiveSnapshot = liveOnMainBranch)) + isLiveSnapshot = isLive)) } } - val missingManifest = if (liveOnMainBranch) new AtomicBoolean(false) else null + val missingManifest = if (isLive) new AtomicBoolean(false) else null val snapshot = Snapshot.fromJson(snapshotJson) collectWithoutDataFileWithManifestFlag( branch, From d142d342ede737f9ce37c9bc325c5828456ee013 Mon Sep 17 00:00:00 2001 From: yehe Date: Sat, 15 Aug 2026 16:09:54 +0800 Subject: [PATCH 06/11] [flink] Fix global abort for distributed orphan file cleanup Signed-off-by: yehe --- .../flink/orphan/FlinkOrphanFilesClean.java | 83 +++++++++++++++---- .../RemoveOrphanFilesActionITCaseBase.java | 35 ++++++++ 2 files changed, 101 insertions(+), 17 deletions(-) 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 45e386136652..8fc8e9497cfc 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 @@ -398,11 +398,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), @@ -413,8 +476,6 @@ public void endInput() throws IOException { private long emittedFilesCount; private long emittedFilesLen; - private boolean abortDeletion; - private final Set used = new HashSet<>(); @Override @@ -431,10 +492,6 @@ public void endInput(int inputId) { checkState(!buildEnd, "Should not build ended."); LOG.info("Finish build phase."); buildEnd = true; - if (abortDeletion) { - LOG.warn( - "Detected missing manifest, aborting clean."); - } break; case 2: checkState(buildEnd, "Should build ended."); @@ -454,21 +511,13 @@ public void endInput(int inputId) { @Override public void processElement1(StreamRecord element) { - String value = element.getValue(); - if (MISSING_MANIFEST_SENTINEL.equals(value)) { - abortDeletion = true; - return; - } - used.add(value); + used.add(element.getValue()); } @Override public void processElement2( StreamRecord> element) { checkState(buildEnd, "Should build ended."); - if (abortDeletion) { - return; - } Tuple2 fileInfo = element.getValue(); String value = fileInfo.f0; Path path = new Path(value); 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(); From 0a746a53775effa9c8328b668bd4866ed65ba21b Mon Sep 17 00:00:00 2001 From: yehe Date: Sun, 30 Aug 2026 21:07:45 +0800 Subject: [PATCH 07/11] reuse live snapshots --- .../operation/LocalOrphanFilesClean.java | 7 ++- .../paimon/operation/OrphanFilesClean.java | 20 +++++---- .../operation/LocalOrphanFilesCleanTest.java | 44 +++++++++++++++++++ .../flink/orphan/FlinkOrphanFilesClean.java | 15 +++---- .../procedure/SparkOrphanFilesClean.scala | 13 ++---- 5 files changed, 67 insertions(+), 32 deletions(-) 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 9a91d7ec7c40..998e2317c701 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 @@ -179,9 +179,8 @@ private void collectWithoutDataFile( Consumer liveManifestConsumer, AtomicBoolean missingManifest) throws IOException { - Set liveSnapshots = - new HashSet<>( - table.switchToBranch(branch).snapshotManager().safelyGetAllSnapshots()); + Set liveSnapshots = safelyGetLiveSnapshots(branch); + Set snapshots = snapshotsIncludingTagsAndChangelogs(branch, liveSnapshots); randomlyOnlyExecute( executor, snapshot -> { @@ -204,7 +203,7 @@ private void collectWithoutDataFile( throw new RuntimeException(e); } }, - safelyGetAllSnapshots(branch)); + snapshots); } private Set getUsedFiles(String branch, AtomicBoolean missingManifest) { 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 35b131877ad2..d8f488dc9c02 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; @@ -244,15 +243,18 @@ 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().safelyGetAllSnapshots()); + } + + 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 void collectWithoutDataFile( 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 06ef4923c374..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; @@ -787,6 +788,49 @@ void testAbortWhenBranchManifestListMissingDuringConcurrentExpiration() throws E } } + @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) { 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 8fc8e9497cfc..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 @@ -67,7 +67,6 @@ import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicLong; import java.util.function.Consumer; -import java.util.stream.Collectors; import static org.apache.flink.api.common.typeinfo.BasicTypeInfo.STRING_TYPE_INFO; import static org.apache.flink.util.Preconditions.checkState; @@ -165,18 +164,16 @@ public void processElement( ctx, Collector> out) throws Exception { - Set liveSnapshotIds = - table.switchToBranch(branch).snapshotManager() - .safelyGetAllSnapshots().stream() - .map(Snapshot::id) - .collect(Collectors.toSet()); - for (Snapshot snapshot : safelyGetAllSnapshots(branch)) { + Set liveSnapshots = + safelyGetLiveSnapshots(branch); + for (Snapshot snapshot : + snapshotsIncludingTagsAndChangelogs( + branch, liveSnapshots)) { out.collect( new Tuple3<>( branch, snapshot.toJson(), - liveSnapshotIds.contains( - snapshot.id()))); + liveSnapshots.contains(snapshot))); } } }) 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 dd822db716f5..1838ba8439fc 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 @@ -70,16 +70,9 @@ case class SparkOrphanFilesClean( .parallelize(branches.asScala.toSeq, maxBranchParallelism) .mapPartitions(_.flatMap { branch => - val liveSnapshotIds = - specifiedTable - .switchToBranch(branch) - .snapshotManager() - .safelyGetAllSnapshots() - .asScala - .map(_.id()) - .toSet - safelyGetAllSnapshots(branch).asScala.map( - snapshot => (branch, snapshot.toJson, liveSnapshotIds.contains(snapshot.id()))) + val liveSnapshots = safelyGetLiveSnapshots(branch) + snapshotsIncludingTagsAndChangelogs(branch, liveSnapshots).asScala.map( + snapshot => (branch, snapshot.toJson, liveSnapshots.contains(snapshot))) }) .repartition(parallelism) .flatMap { From 3f717a0862690221660937991282b2f1e03c7b4a Mon Sep 17 00:00:00 2001 From: yehe Date: Sun, 30 Aug 2026 22:00:39 +0800 Subject: [PATCH 08/11] fix orphan cleanup resource --- .../operation/LocalOrphanFilesClean.java | 6 --- .../procedure/SparkOrphanFilesClean.scala | 48 +++++++++++++------ 2 files changed, 34 insertions(+), 20 deletions(-) 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 998e2317c701..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 @@ -118,12 +118,6 @@ public CleanOrphanFilesResult clean() .flatMap(branch -> getUsedFiles(branch, missingManifest).stream()) .collect(Collectors.toSet()); - if (usedFiles.isEmpty()) { - LOG.warn("Collected used files is empty, aborting orphan files clean."); - return new CleanOrphanFilesResult( - deleteFiles.size(), deletedFilesLenInBytes.get(), deleteFiles); - } - if (missingManifest.get()) { LOG.warn("Detected missing manifest during used-files collection, aborting clean."); return new CleanOrphanFilesResult( 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 1838ba8439fc..e2a78ca92fce 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 @@ -40,6 +40,7 @@ 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() @@ -105,13 +118,14 @@ case class SparkOrphanFilesClean( } .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()))), - usedManifestFiles) + cachedDatasets.toSeq) } val dataFilesWithFlag = usedManifestFiles @@ -150,13 +164,14 @@ case class SparkOrphanFilesClean( } } .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()))), - usedManifestFiles) + cachedDatasets.toSeq) } val dataFiles = dataFilesWithFlag.filter(!_._2).map(_._1) @@ -231,7 +246,7 @@ case class SparkOrphanFilesClean( deleted } - (finalDeletedDataset, usedManifestFiles) + (finalDeletedDataset, cachedDatasets.toSeq) } } @@ -282,17 +297,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") From c2cb7f0535cafa04ec6bc4d1b1822abecee8813b Mon Sep 17 00:00:00 2001 From: yehe Date: Mon, 5 Oct 2026 17:03:39 +0800 Subject: [PATCH 09/11] Fix incomplete live snapshot enumeration --- .../paimon/operation/OrphanFilesClean.java | 10 ++++- .../apache/paimon/utils/SnapshotManager.java | 18 ++++++++ .../paimon/utils/SnapshotManagerTest.java | 45 +++++++++++++++++++ 3 files changed, 71 insertions(+), 2 deletions(-) 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 d8f488dc9c02..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 @@ -245,7 +245,8 @@ protected boolean isManagedBlobPack(Path path) { protected Set safelyGetLiveSnapshots(String branch) throws IOException { FileStoreTable branchTable = table.switchToBranch(branch); - return new HashSet<>(branchTable.snapshotManager().safelyGetAllSnapshots()); + return new HashSet<>( + branchTable.snapshotManager().safelyGetAllSnapshotsWithConsistentLatest()); } protected Set snapshotsIncludingTagsAndChangelogs( @@ -257,6 +258,11 @@ protected Set snapshotsIncludingTagsAndChangelogs( return snapshots; } + protected Set safelyGetAllSnapshots(String branch) throws IOException { + Set liveSnapshots = safelyGetLiveSnapshots(branch); + return snapshotsIncludingTagsAndChangelogs(branch, liveSnapshots); + } + protected void collectWithoutDataFile( String branch, Snapshot snapshot, @@ -316,7 +322,7 @@ protected void collectWithoutDataFileWithManifestFlag( // index files String indexManifest = snapshot.indexManifest(); - if (indexManifest != null && indexFileHandler.existsManifest(indexManifest)) { + if (indexManifest != null) { List indexEntries = retryReadingFiles( () -> indexFileHandler.readManifestWithIOException(indexManifest), 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/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, From 582c393998993591719b86c3d937840982f3fb52 Mon Sep 17 00:00:00 2001 From: yehe Date: Tue, 6 Oct 2026 01:10:54 +0800 Subject: [PATCH 10/11] Fix Spark orphan cleanup cache recomputation --- .../procedure/SparkOrphanFilesClean.scala | 19 +++++-- .../RemoveOrphanFilesProcedureTest.scala | 53 +++++++++++++++++++ 2 files changed, 69 insertions(+), 3 deletions(-) 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 e2a78ca92fce..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 @@ -174,12 +174,25 @@ case class SparkOrphanFilesClean( cachedDatasets.toSeq) } - val dataFiles = dataFilesWithFlag.filter(!_._2).map(_._1) + 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 - .filter(!_.isMissing) - .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") 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()) + } + } + } From 09a86e021779569d7ea79bd50f58729a26fec170 Mon Sep 17 00:00:00 2001 From: yehe Date: Tue, 6 Oct 2026 11:23:17 +0800 Subject: [PATCH 11/11] Trigger CI