diff --git a/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestAvroReader.java b/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestAvroReader.java index 95134f36a002..1c6c121eb15b 100644 --- a/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestAvroReader.java +++ b/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestAvroReader.java @@ -439,6 +439,18 @@ public RowIterator toRows(RowType projectedType) throws IOException { return toRows(projectedType, null, null, true); } + /** + * Lazily decompresses this block, applies manifest filters before decoding file metadata, + * and returns an iterator over one reusable row. + */ + public RowIterator toRows( + RowType projectedType, + @Nullable PartitionPredicate partitionFilter, + @Nullable BucketFilter bucketFilter) + throws IOException { + return toRows(projectedType, partitionFilter, bucketFilter, true); + } + private RowIterator toRows( RowType projectedType, @Nullable PartitionPredicate partitionFilter, diff --git a/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestFile.java b/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestFile.java index afc122c7235d..608c7dc57610 100644 --- a/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestFile.java +++ b/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestFile.java @@ -48,6 +48,7 @@ import java.util.Arrays; import java.util.List; import java.util.function.Function; +import java.util.function.Predicate; /** * This file includes several {@link ManifestEntry}s, representing the additional changes since last @@ -303,17 +304,41 @@ public long suggestedFileSize() { } public List readExpireFileEntries(String fileName) { + return readExpireFileEntries(fileName, null, entry -> true); + } + + /** + * Reads only expiring entries accepted by the supplied filters. + * + *

The bucket filter is evaluated by the Avro reader before nested data-file metadata is + * decoded. The entry filter then runs on a reusable projected view, before an {@link + * ExpireFileEntry} is materialized. The entry filter must not retain its argument. + */ + public List readExpireFileEntries( + String fileName, + @Nullable BucketFilter bucketFilter, + Predicate entryFilter) { List result = new ArrayList<>(); - try (CloseableIterator entries = - scan(fileName, EXPIRE_FILE_PROJECTION)) { - while (entries.hasNext()) { - result.add(ExpireFileEntry.from(entries.next())); + ProjectedManifestEntry entry = EXPIRE_FILE_PROJECTION.createEntry(); + try (ManifestAvroReader reader = scanAvroBlocks(fileName, null)) { + while (reader.hasNext()) { + ManifestAvroReader.RowIterator rows = + reader.next() + .toRows(EXPIRE_FILE_PROJECTION.projectedType(), null, bucketFilter); + while (rows.hasNext()) { + entry.replace(rows.next()); + if (entryFilter.test(entry)) { + result.add(ExpireFileEntry.from(entry)); + } + } } } catch (Exception e) { throw new RuntimeException( String.format( "Failed to scan expiring entries from manifest file '%s'.", fileName), e); + } finally { + entry.clear(); } return result; } diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/FileDeletionBase.java b/paimon-core/src/main/java/org/apache/paimon/operation/FileDeletionBase.java index ad544af0ea6c..1b8b245e9b39 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/FileDeletionBase.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/FileDeletionBase.java @@ -25,14 +25,17 @@ import org.apache.paimon.index.IndexFileHandler; import org.apache.paimon.index.IndexFileMeta; import org.apache.paimon.io.DataFilePathFactory; +import org.apache.paimon.manifest.BucketFilter; import org.apache.paimon.manifest.ExpireFileEntry; import org.apache.paimon.manifest.FileEntry; import org.apache.paimon.manifest.FileEntry.Identifier; import org.apache.paimon.manifest.FileKind; import org.apache.paimon.manifest.IndexManifestEntry; +import org.apache.paimon.manifest.ManifestBucketFilter; import org.apache.paimon.manifest.ManifestFile; import org.apache.paimon.manifest.ManifestFileMeta; import org.apache.paimon.manifest.ManifestList; +import org.apache.paimon.manifest.ProjectedManifestEntry; import org.apache.paimon.stats.StatsFileHandler; import org.apache.paimon.utils.DataFilePathFactories; import org.apache.paimon.utils.FileOperationThreadPool; @@ -56,7 +59,9 @@ import java.util.LinkedHashSet; import java.util.List; import java.util.Map; +import java.util.NavigableSet; import java.util.Set; +import java.util.TreeSet; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutionException; @@ -185,7 +190,7 @@ protected void recordDeletionBuckets(ExpireFileEntry entry) { } /** Plan data files referenced by DELETE entries in the snapshot's delta manifest list. */ - public List planDeletedInDeltaManifest(T snapshot, Predicate skipper) { + public DataFileDeletionPlan planDeletedInDeltaManifest(T snapshot) { String deltaManifestList = snapshot.deltaManifestList(); // data file path -> (original manifest entry, extra file paths) Map>> dataFileToDelete = new HashMap<>(); @@ -218,12 +223,22 @@ public List planDeletedInDeltaManifest(T snapshot, Predicate planDeletedInDeltaManifest(T snapshot, Predicate skipper) { + return dataFilesToDelete(planDeletedInDeltaManifest(snapshot), skipper); + } + + public List dataFilesToDelete( + DataFileDeletionPlan plan, Predicate skipper) { // apply skipper List actualDataFileToDelete = new ArrayList<>(); - dataFileToDelete.forEach( + plan.dataFileToDelete.forEach( (path, pair) -> { ExpireFileEntry entry = pair.getLeft(); // check whether we should skip the data file @@ -406,6 +421,38 @@ public Predicate createDataFileSkipperForTag(Snapshot tag) thro return entry -> containsDataFile(tagDataFiles, entry); } + /** + * Creates a tag skipper restricted to the files in the supplied deletion plans. + * + *

A tag may reference every data file in a large table. Building an index for all of them + * makes snapshot expiration use memory proportional to the table size (and, when tags are read + * concurrently, to the number of tags). Only files that are candidates for the current + * expiration batch can be deleted, so merge only matching tag entries instead. + */ + public Predicate createDataFileSkipperForTag( + Snapshot tag, Collection plans) throws Exception { + Map>> candidates = new HashMap<>(); + for (DataFileDeletionPlan plan : plans) { + for (Pair> pair : plan.dataFileToDelete.values()) { + addDataFile(candidates, pair.getLeft()); + } + } + if (candidates.isEmpty()) { + return entry -> false; + } + + Collection matchingEntries = + readMergedDataFiles( + manifestList.readDataManifests(tag), + createCandidateBucketFilter(candidates), + entry -> containsDataFile(candidates, entry)); + Map>> taggedCandidates = new HashMap<>(); + for (ExpireFileEntry entry : matchingEntries) { + addDataFile(taggedCandidates, entry); + } + return entry -> containsDataFile(taggedCandidates, entry); + } + /** * It is possible that a job was killed during expiration and some manifest files have been * deleted, so if the clean methods need to get manifests of a snapshot to be cleaned, we should @@ -430,10 +477,7 @@ protected void addMergedDataFiles( throws IOException { for (ExpireFileEntry entry : readMergedDataFiles(manifestList.readDataManifests(snapshot))) { - dataFiles - .computeIfAbsent(entry.partition(), p -> new HashMap<>()) - .computeIfAbsent(entry.bucket(), b -> new HashSet<>()) - .add(entry.fileName()); + addDataFile(dataFiles, entry); } } @@ -444,13 +488,79 @@ protected Collection readMergedDataFiles(List return map.values(); } + protected Collection readMergedDataFiles( + List manifests, + BucketFilter bucketFilter, + Predicate filter) + throws IOException { + Map map = new HashMap<>(); + FileEntry.mergeEntries( + ManifestReadThreadPool.sequentialBatchedExecute( + manifest -> { + if (!bucketFilter.mayContain(manifest)) { + return Collections.emptyList(); + } + return manifestFile.readExpireFileEntries( + manifest.fileName(), bucketFilter, filter); + }, + manifests, + manifestReadParallelism), + map); + return map.values(); + } + + private BucketFilter createCandidateBucketFilter( + Map>> candidates) { + NavigableSet candidateBuckets = new TreeSet<>(); + for (Map> buckets : candidates.values()) { + candidateBuckets.addAll(buckets.keySet()); + } + + ManifestBucketFilter filter = + new ManifestBucketFilter() { + @Override + public boolean test(BinaryRow partition, Integer bucket, Integer totalBuckets) { + Map> buckets = candidates.get(partition); + return buckets != null && buckets.containsKey(bucket); + } + + @Override + public boolean mayContain(int minBucket, int maxBucket, int totalBuckets) { + Integer firstCandidate = candidateBuckets.ceiling(minBucket); + return firstCandidate != null && firstCandidate <= maxBucket; + } + }; + return new BucketFilter(false, null, null, filter); + } + + private void addDataFile( + Map>> dataFiles, ExpireFileEntry entry) { + dataFiles + .computeIfAbsent(entry.partition(), p -> new HashMap<>()) + .computeIfAbsent(entry.bucket(), b -> new HashSet<>()) + .add(entry.fileName()); + } + protected boolean containsDataFile( Map>> dataFiles, ExpireFileEntry entry) { - Map> buckets = dataFiles.get(entry.partition()); + return containsDataFile(dataFiles, entry.partition(), entry.bucket(), entry.fileName()); + } + + private boolean containsDataFile( + Map>> dataFiles, ProjectedManifestEntry entry) { + return containsDataFile(dataFiles, entry.partition(), entry.bucket(), entry.fileName()); + } + + private boolean containsDataFile( + Map>> dataFiles, + BinaryRow partition, + int bucket, + String fileName) { + Map> buckets = dataFiles.get(partition); if (buckets != null) { - Set fileNames = buckets.get(entry.bucket()); + Set fileNames = buckets.get(bucket); if (fileNames != null) { - return fileNames.contains(entry.fileName()); + return fileNames.contains(fileName); } } return false; @@ -560,4 +670,19 @@ public void executeAll(Collection tasks) { throw new RuntimeException(e.getCause()); } } + + /** Candidate data files from one snapshot delta manifest. */ + public static class DataFileDeletionPlan { + + private final Map>> dataFileToDelete; + + private DataFileDeletionPlan( + Map>> dataFileToDelete) { + this.dataFileToDelete = dataFileToDelete; + } + + private static DataFileDeletionPlan empty() { + return new DataFileDeletionPlan(Collections.emptyMap()); + } + } } diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/SnapshotDeletion.java b/paimon-core/src/main/java/org/apache/paimon/operation/SnapshotDeletion.java index 00bb731d03a4..bbeb12462775 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/SnapshotDeletion.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/SnapshotDeletion.java @@ -70,8 +70,8 @@ public void cleanDeletedDataFiles(Snapshot snapshot, Predicate } @Override - public List planDeletedInDeltaManifest( - Snapshot snapshot, Predicate skipper) { + public List dataFilesToDelete( + DataFileDeletionPlan plan, Predicate skipper) { Predicate enriched = skipper; if (changelogDecoupled && !produceChangelog) { // Skip clean the 'APPEND' data files.If we do not have the file source information @@ -83,7 +83,7 @@ public List planDeletedInDeltaManifest( || (manifestEntry.fileSource().orElse(FileSource.APPEND) == FileSource.APPEND); } - return super.planDeletedInDeltaManifest(snapshot, enriched); + return super.dataFilesToDelete(plan, enriched); } @Override diff --git a/paimon-core/src/main/java/org/apache/paimon/table/ExpireSnapshotsImpl.java b/paimon-core/src/main/java/org/apache/paimon/table/ExpireSnapshotsImpl.java index b7cabb4d32f3..4b8008548921 100644 --- a/paimon-core/src/main/java/org/apache/paimon/table/ExpireSnapshotsImpl.java +++ b/paimon-core/src/main/java/org/apache/paimon/table/ExpireSnapshotsImpl.java @@ -24,6 +24,7 @@ import org.apache.paimon.consumer.ConsumerManager; import org.apache.paimon.fs.Path; import org.apache.paimon.manifest.ExpireFileEntry; +import org.apache.paimon.operation.FileDeletionBase.DataFileDeletionPlan; import org.apache.paimon.operation.SnapshotDeletion; import org.apache.paimon.options.ExpireConfig; import org.apache.paimon.tag.Tag; @@ -274,8 +275,7 @@ private void cleanDataFiles( List taggedSnapshots, long beginInclusiveId) throws ExecutionException, InterruptedException { - Map tagIdBySnapshotId = new HashMap<>(); - Map tags = new HashMap<>(); + Map tagBySnapshotId = new HashMap<>(); int tagIndex = -1; for (Snapshot snapshot : snapshotsIncludingEnd) { long id = snapshot.id(); @@ -286,16 +286,11 @@ private void cleanDataFiles( tagIndex = advancePreviousSnapshot(taggedSnapshots, tagIndex, id); if (tagIndex >= 0) { Snapshot tag = taggedSnapshots.get(tagIndex); - tagIdBySnapshotId.put(id, tag.id()); - tags.put(tag.id(), tag); + tagBySnapshotId.put(id, tag); } } - Map>> skippers = - collectTagSkippers(tags.values()); - Predicate deleteAll = entry -> false; - List>> futures = new ArrayList<>(); - int plannedSnapshots = 0; + List batch = new ArrayList<>(snapshotExpireBatchSize); for (Snapshot snapshot : snapshotsIncludingEnd) { long id = snapshot.id(); if (id == beginInclusiveId) { @@ -304,44 +299,76 @@ private void cleanDataFiles( if (LOG.isDebugEnabled()) { LOG.debug("Ready to delete merge tree files not used by snapshot #{}", id); } + batch.add(snapshot); + if (batch.size() >= snapshotExpireBatchSize) { + cleanDataFileBatch(batch, tagBySnapshotId); + } + } + cleanDataFileBatch(batch, tagBySnapshotId); + } + + private void cleanDataFileBatch(List snapshots, Map tagBySnapshotId) + throws ExecutionException, InterruptedException { + if (snapshots.isEmpty()) { + return; + } - Long tagId = tagIdBySnapshotId.get(id); + List> planFutures = new ArrayList<>(); + for (Snapshot snapshot : snapshots) { + planFutures.add( + CompletableFuture.supplyAsync( + () -> snapshotDeletion.planDeletedInDeltaManifest(snapshot), + fileExecutor)); + } + List plans = getAll(planFutures); + + Map tags = new HashMap<>(); + Map> plansByTag = new HashMap<>(); + for (int i = 0; i < snapshots.size(); i++) { + Snapshot tag = tagBySnapshotId.get(snapshots.get(i).id()); + if (tag != null) { + tags.put(tag.id(), tag); + plansByTag.computeIfAbsent(tag.id(), id -> new ArrayList<>()).add(plans.get(i)); + } + } + + Map>> skippers = + collectTagSkippers(tags, plansByTag); + Predicate deleteAll = entry -> false; + List paths = new ArrayList<>(); + for (int i = 0; i < snapshots.size(); i++) { + Snapshot snapshot = snapshots.get(i); + Snapshot tag = tagBySnapshotId.get(snapshot.id()); Optional> skipper = - tagId == null + tag == null ? Optional.of(deleteAll) - : skippers.getOrDefault(tagId, Optional.empty()); + : skippers.getOrDefault(tag.id(), Optional.empty()); if (!skipper.isPresent()) { LOG.info( "Skip cleaning data files of snapshot '{}' due to failed to build skipping set.", - id); + snapshot.id()); continue; } - futures.add( - CompletableFuture.supplyAsync( - () -> - snapshotDeletion.planDeletedInDeltaManifest( - snapshot, skipper.get()), - fileExecutor)); - if (++plannedSnapshots >= snapshotExpireBatchSize) { - cleanBatch(futures); - plannedSnapshots = 0; - } + paths.addAll(snapshotDeletion.dataFilesToDelete(plans.get(i), skipper.get())); } - cleanBatch(futures); + snapshotDeletion.cleanDataFiles(paths); + snapshots.clear(); } private Map>> collectTagSkippers( - Collection tags) throws ExecutionException, InterruptedException { + Map tags, Map> plansByTag) + throws ExecutionException, InterruptedException { Map>>> futures = new HashMap<>(); - for (Snapshot tag : tags) { + for (Snapshot tag : tags.values()) { futures.put( tag.id(), CompletableFuture.supplyAsync( () -> { try { return Optional.of( - snapshotDeletion.createDataFileSkipperForTag(tag)); + snapshotDeletion.createDataFileSkipperForTag( + tag, plansByTag.get(tag.id()))); } catch (Exception e) { LOG.info( "Failed to build data file skipping set for tag snapshot '{}'.", diff --git a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileTest.java b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileTest.java index 5172c305216b..24f84659d147 100644 --- a/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/manifest/ManifestFileTest.java @@ -1084,6 +1084,50 @@ void testReadExpireFileEntriesWithProjectedScan() { } } + @Test + void testReadExpireFileEntriesPushesFiltersIntoStreamingScan() { + List entries = Arrays.asList(gen.next(), gen.next(), gen.next()); + ManifestFile manifestFile = createManifestFile(tempDir.toString(), Long.MAX_VALUE); + ManifestFileMeta manifest = writeSingleManifest(manifestFile, entries); + + int[] exactFilterCalls = {0}; + List bucketRejected = + manifestFile.readExpireFileEntries( + manifest.fileName(), + new BucketFilter(false, null, bucket -> false, null), + entry -> { + exactFilterCalls[0]++; + return true; + }); + assertThat(bucketRejected).isEmpty(); + assertThat(exactFilterCalls[0]).isZero(); + + String selectedFile = entries.get(1).fileName(); + ProjectedManifestEntry[] reusableView = {null}; + List selected = + manifestFile.readExpireFileEntries( + manifest.fileName(), + null, + entry -> { + if (reusableView[0] == null) { + reusableView[0] = entry; + } else { + assertThat(entry).isSameAs(reusableView[0]); + } + exactFilterCalls[0]++; + return entry.fileName().equals(selectedFile); + }); + + assertThat(exactFilterCalls[0]).isEqualTo(entries.size()); + assertThat(selected) + .containsExactly( + ExpireFileEntry.from( + entries.stream() + .filter(entry -> entry.fileName().equals(selectedFile)) + .findFirst() + .orElseThrow(AssertionError::new))); + } + @Test void testScanProjectedManifestCreatesDistinctEntryWrappers() throws Exception { List entries = Arrays.asList(gen.next(), gen.next(), gen.next()); diff --git a/paimon-core/src/test/java/org/apache/paimon/operation/ExpireSnapshotsTest.java b/paimon-core/src/test/java/org/apache/paimon/operation/ExpireSnapshotsTest.java index 957e4b73177c..7fac929cdd32 100644 --- a/paimon-core/src/test/java/org/apache/paimon/operation/ExpireSnapshotsTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/operation/ExpireSnapshotsTest.java @@ -33,6 +33,7 @@ import org.apache.paimon.fs.local.LocalFileIO; import org.apache.paimon.io.DataFileMeta; import org.apache.paimon.io.DataFilePathFactory; +import org.apache.paimon.manifest.BucketFilter; import org.apache.paimon.manifest.ExpireFileEntry; import org.apache.paimon.manifest.FileKind; import org.apache.paimon.manifest.FileSource; @@ -40,7 +41,9 @@ import org.apache.paimon.manifest.ManifestFileMeta; import org.apache.paimon.manifest.ManifestIndexTestUtils; import org.apache.paimon.manifest.ManifestSidecar; +import org.apache.paimon.manifest.ProjectedManifestEntry; import org.apache.paimon.mergetree.compact.DeduplicateMergeFunction; +import org.apache.paimon.operation.FileDeletionBase.DataFileDeletionPlan; import org.apache.paimon.options.ExpireConfig; import org.apache.paimon.options.MemorySize; import org.apache.paimon.schema.FileSystemSchemaManager; @@ -74,6 +77,7 @@ import java.util.Arrays; import java.util.Collection; import java.util.Collections; +import java.util.HashMap; import java.util.HashSet; import java.util.List; import java.util.Map; @@ -966,6 +970,85 @@ public void testExpireWithTagsAndConcurrentPlanningKeepsTaggedSnapshotsReadable( assertSnapshot(tagManager.getOrThrow("tag6").trimToSnapshot(), allData, snapshotPositions); } + @Test + public void testExpireFiltersTagFilesByDeletionCandidates() throws Exception { + List candidateData = FileStoreTestUtils.partitionedData(5, gen, "0401", 8); + List unrelatedData = FileStoreTestUtils.partitionedData(5, gen, "0402", 8); + BinaryRow candidatePartition = gen.getPartition(candidateData.get(0)); + BinaryRow unrelatedPartition = gen.getPartition(unrelatedData.get(0)); + Map>> writers = new HashMap<>(); + writers.put( + candidatePartition, + Collections.singletonMap( + 0, + FileStoreTestUtils.writeData(store, candidateData, candidatePartition, 0))); + writers.put( + unrelatedPartition, + Collections.singletonMap( + 0, + FileStoreTestUtils.writeData(store, unrelatedData, unrelatedPartition, 0))); + FileStoreTestUtils.commitData(store, 0, writers); + + Snapshot taggedSnapshot = snapshotManager.latestSnapshot(); + TagManager tagManager = store.newTagManager(); + tagManager.createTag( + taggedSnapshot, + "tag1", + store.options().tagDefaultTimeRetained(), + Collections.emptyList(), + false); + + List delete = + store.newScan().plan().files().stream() + .filter(entry -> entry.partition().equals(candidatePartition)) + .map( + entry -> + ManifestEntry.create( + FileKind.DELETE, + entry.partition(), + entry.bucket(), + entry.totalBuckets(), + entry.file())) + .collect(Collectors.toList()); + try (FileStoreCommitImpl commit = store.newCommit()) { + commit.tryCommitOnce( + null, + delete, + Collections.emptyList(), + Collections.emptyList(), + 1, + null, + Collections.emptyMap(), + Snapshot.CommitKind.APPEND, + false, + taggedSnapshot, + true, + null); + } + + CandidateTrackingSnapshotDeletion snapshotDeletion = + new CandidateTrackingSnapshotDeletion(store); + ExpireSnapshotsImpl expire = + newExpireWithSnapshotDeletion(store, snapshotManager, snapshotDeletion); + expire.config( + ExpireConfig.builder() + .snapshotRetainMin(1) + .snapshotRetainMax(1) + .snapshotTimeRetain(Duration.ofMillis(Long.MAX_VALUE)) + .build()); + expire.expire(); + + assertThat(snapshotDeletion.candidateTagSkipperCalls()).isGreaterThan(0); + assertThat(snapshotDeletion.filteredTagReads()).isGreaterThan(0); + assertThat(snapshotDeletion.materializedTagEntries()).isOne(); + List taggedData = new ArrayList<>(candidateData); + taggedData.addAll(unrelatedData); + assertSnapshot( + tagManager.getOrThrow("tag1").trimToSnapshot(), + taggedData, + Collections.singletonList(taggedData.size())); + } + @Test public void testExpireReadsTagsConcurrentlyWithObjectStoreFileIO() throws Exception { TestFileStore slowStore = createSlowStore(); @@ -1567,12 +1650,11 @@ private void blockManifestPlans() { } @Override - public List planDeletedInDeltaManifest( - Snapshot snapshot, Predicate skipper) { + public DataFileDeletionPlan planDeletedInDeltaManifest(Snapshot snapshot) { if (blockDataFilePlans && shouldBlock(snapshot.id())) { dataFilePlans.awaitConcurrentCall(); } - return super.planDeletedInDeltaManifest(snapshot, skipper); + return super.planDeletedInDeltaManifest(snapshot); } @Override @@ -1608,6 +1690,71 @@ private int maxActiveManifestPlans() { } } + private static class CandidateTrackingSnapshotDeletion extends SnapshotDeletion { + + private final AtomicInteger candidateTagSkipperCalls = new AtomicInteger(); + private final AtomicInteger filteredTagReads = new AtomicInteger(); + private final AtomicInteger materializedTagEntries = new AtomicInteger(); + + private CandidateTrackingSnapshotDeletion(TestFileStore store) { + super( + store.fileIO(), + store.pathFactory(), + store.manifestFileFactory().create(), + store.manifestListFactory().create(), + store.newIndexFileHandler(), + store.newStatsFileHandler(), + store.options().changelogProducer() != CoreOptions.ChangelogProducer.NONE, + store.options().cleanEmptyDirectories(), + store.options().fileOperationThreadNum(), + store.options().scanManifestParallelism()); + } + + @Override + public Predicate createDataFileSkipperForTag(Snapshot tag) + throws Exception { + throw new AssertionError("Snapshot expiration must not index all files in a tag."); + } + + @Override + public Predicate createDataFileSkipperForTag( + Snapshot tag, Collection plans) throws Exception { + candidateTagSkipperCalls.incrementAndGet(); + return super.createDataFileSkipperForTag(tag, plans); + } + + @Override + protected Collection readMergedDataFiles(List manifests) + throws IOException { + throw new AssertionError("Snapshot expiration must not merge all files in a tag."); + } + + @Override + protected Collection readMergedDataFiles( + List manifests, + BucketFilter bucketFilter, + Predicate filter) + throws IOException { + filteredTagReads.incrementAndGet(); + Collection entries = + super.readMergedDataFiles(manifests, bucketFilter, filter); + materializedTagEntries.addAndGet(entries.size()); + return entries; + } + + private int candidateTagSkipperCalls() { + return candidateTagSkipperCalls.get(); + } + + private int filteredTagReads() { + return filteredTagReads.get(); + } + + private int materializedTagEntries() { + return materializedTagEntries.get(); + } + } + private static class CapturingSnapshotDeletion extends SnapshotDeletion { private final List> deleteBatches = new ArrayList<>();