This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new dd5ce17c26 [core] Optimize manifest entry scans (#8899)
dd5ce17c26 is described below

commit dd5ce17c26394899ceeb129dbb868a56089d75d3
Author: Jingsong Lee <[email protected]>
AuthorDate: Wed Jul 29 11:36:42 2026 +0800

    [core] Optimize manifest entry scans (#8899)
---
 .../java/org/apache/paimon/manifest/FileEntry.java | 47 ++++++++---
 .../org/apache/paimon/manifest/ManifestFile.java   | 45 +++++++++-
 .../apache/paimon/manifest/ManifestFileTest.java   | 98 ++++++++++++++++++++++
 3 files changed, 176 insertions(+), 14 deletions(-)

diff --git 
a/paimon-core/src/main/java/org/apache/paimon/manifest/FileEntry.java 
b/paimon-core/src/main/java/org/apache/paimon/manifest/FileEntry.java
index 62f3f685c9..f619f02d60 100644
--- a/paimon-core/src/main/java/org/apache/paimon/manifest/FileEntry.java
+++ b/paimon-core/src/main/java/org/apache/paimon/manifest/FileEntry.java
@@ -20,11 +20,13 @@ package org.apache.paimon.manifest;
 
 import org.apache.paimon.data.BinaryRow;
 import org.apache.paimon.data.InternalRow;
+import org.apache.paimon.utils.CloseableIterator;
 import org.apache.paimon.utils.FileStorePathFactory;
 import org.apache.paimon.utils.Filter;
 
 import javax.annotation.Nullable;
 
+import java.util.ArrayList;
 import java.util.Arrays;
 import java.util.Collection;
 import java.util.Iterator;
@@ -237,16 +239,41 @@ public interface FileEntry {
             ManifestFile manifestFile,
             List<ManifestFileMeta> manifestFiles,
             @Nullable Integer manifestReadParallelism) {
-        return readDeletedEntries(
-                m ->
-                        manifestFile.read(
-                                m.fileName(),
-                                m.fileSize(),
-                                deletedFilter(),
-                                Filter.alwaysTrue(),
-                                SimpleFileEntry::from),
-                manifestFiles,
-                manifestReadParallelism);
+        manifestFiles =
+                manifestFiles.stream()
+                        .filter(file -> file.numDeletedFiles() > 0)
+                        .collect(Collectors.toList());
+        Function<ManifestFileMeta, List<Identifier>> processor =
+                manifest -> {
+                    List<Identifier> identifiers =
+                            new ArrayList<>((int) 
Math.min(manifest.numDeletedFiles(), 1 << 20));
+                    try (CloseableIterator<BinaryManifestEntry> entries =
+                            manifestFile.scan(
+                                    manifest.fileName(),
+                                    manifest.fileSize(),
+                                    
BinaryManifestEntry.DELETE_ENTRY_PROJECTION)) {
+                        while (entries.hasNext()) {
+                            BinaryManifestEntry entry = entries.next();
+                            if (entry.isDelete()) {
+                                identifiers.add(entry.identifier());
+                            }
+                        }
+                    } catch (Exception e) {
+                        throw new RuntimeException(
+                                String.format(
+                                        "Failed to scan deleted entries from 
manifest file '%s'.",
+                                        manifest.fileName()),
+                                e);
+                    }
+                    return identifiers;
+                };
+        Iterator<Identifier> identifiers =
+                randomlyExecuteSequentialReturn(processor, manifestFiles, 
manifestReadParallelism);
+        Set<Identifier> result = ConcurrentHashMap.newKeySet();
+        while (identifiers.hasNext()) {
+            result.add(identifiers.next());
+        }
+        return result;
     }
 
     static <T extends FileEntry> Set<Identifier> readDeletedEntries(
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 c5ee9d9c0c..4f1aba75ba 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
@@ -26,6 +26,7 @@ import org.apache.paimon.format.FormatWriterFactory;
 import org.apache.paimon.format.SimpleStatsCollector;
 import org.apache.paimon.fs.FileIO;
 import org.apache.paimon.fs.Path;
+import org.apache.paimon.io.DataFileMeta;
 import org.apache.paimon.io.RollingFileWriter;
 import org.apache.paimon.io.RollingFileWriterImpl;
 import org.apache.paimon.io.SingleFileWriter;
@@ -48,6 +49,7 @@ import javax.annotation.Nullable;
 import java.io.IOException;
 import java.io.UncheckedIOException;
 import java.util.ArrayList;
+import java.util.Arrays;
 import java.util.Collections;
 import java.util.List;
 import java.util.function.Function;
@@ -58,6 +60,8 @@ import java.util.function.Function;
  */
 public class ManifestFile extends ObjectsFile<ManifestEntry> {
 
+    private static final Projection EXPIRE_FILE_PROJECTION = 
createExpireFileProjection();
+
     private final SchemaManager schemaManager;
     private final RowType partitionType;
     private final FileFormat fileFormat;
@@ -206,14 +210,47 @@ public class ManifestFile extends 
ObjectsFile<ManifestEntry> {
     }
 
     public List<ExpireFileEntry> readExpireFileEntries(String fileName, 
@Nullable Long fileSize) {
-        List<ManifestEntry> entries = read(fileName, fileSize);
-        List<ExpireFileEntry> result = new ArrayList<>(entries.size());
-        for (ManifestEntry entry : entries) {
-            result.add(ExpireFileEntry.from(entry));
+        List<ExpireFileEntry> result = new ArrayList<>();
+        try (CloseableIterator<BinaryManifestEntry> entries =
+                scan(fileName, fileSize, EXPIRE_FILE_PROJECTION)) {
+            while (entries.hasNext()) {
+                result.add(ExpireFileEntry.from(entries.next()));
+            }
+        } catch (Exception e) {
+            throw new RuntimeException(
+                    String.format(
+                            "Failed to scan expiring entries from manifest 
file '%s'.", fileName),
+                    e);
         }
         return result;
     }
 
+    private static Projection createExpireFileProjection() {
+        RowType manifestType = ManifestEntry.MANIFEST_ROW_TYPE;
+        return Projection.create(
+                new RowType(
+                        false,
+                        Arrays.asList(
+                                manifestType.getField(ManifestEntry.KIND),
+                                manifestType.getField(ManifestEntry.PARTITION),
+                                manifestType.getField(ManifestEntry.BUCKET),
+                                
manifestType.getField(ManifestEntry.TOTAL_BUCKETS),
+                                manifestType
+                                        .getField(ManifestEntry.FILE)
+                                        .newType(
+                                                DataFileMeta.SCHEMA.project(
+                                                        DataFileMeta.FILE_NAME,
+                                                        DataFileMeta.ROW_COUNT,
+                                                        DataFileMeta.MIN_KEY,
+                                                        DataFileMeta.MAX_KEY,
+                                                        DataFileMeta.LEVEL,
+                                                        
DataFileMeta.EXTRA_FILES,
+                                                        
DataFileMeta.EMBEDDED_FILE_INDEX,
+                                                        
DataFileMeta.FILE_SOURCE,
+                                                        
DataFileMeta.EXTERNAL_PATH,
+                                                        
DataFileMeta.FIRST_ROW_ID)))));
+    }
+
     /**
      * Write several {@link ManifestEntry}s into manifest files.
      *
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 e99cf205c1..fef9f7b9ac 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
@@ -43,6 +43,7 @@ import java.io.IOException;
 import java.util.ArrayList;
 import java.util.Arrays;
 import java.util.List;
+import java.util.Set;
 import java.util.UUID;
 import java.util.concurrent.ThreadLocalRandom;
 import java.util.concurrent.atomic.AtomicInteger;
@@ -128,6 +129,103 @@ public class ManifestFileTest {
         assertThat(creationTimesFound).isPositive();
     }
 
+    @Test
+    void testReadDeletedEntriesWithProjectedScan() {
+        ManifestEntry first = gen.next();
+        ManifestEntry second = gen.next();
+        DataFileMeta firstFile =
+                first.file()
+                        .copy(Arrays.asList("extra-1", "extra-2"))
+                        .copy(new byte[] {1, 2})
+                        .newExternalPath("external/first")
+                        .newFirstRowId(10L);
+        DataFileMeta secondFile =
+                second.file()
+                        .copy(Arrays.asList("extra-3"))
+                        .copy(new byte[] {3, 4})
+                        .newExternalPath("external/second")
+                        .newFirstRowId(20L);
+        ManifestEntry firstAdd =
+                ManifestEntry.create(
+                        FileKind.ADD,
+                        first.partition(),
+                        first.bucket(),
+                        first.totalBuckets(),
+                        firstFile);
+        ManifestEntry firstDelete =
+                ManifestEntry.create(
+                        FileKind.DELETE,
+                        first.partition(),
+                        first.bucket(),
+                        first.totalBuckets(),
+                        firstFile);
+        ManifestEntry secondAdd =
+                ManifestEntry.create(
+                        FileKind.ADD,
+                        second.partition(),
+                        second.bucket(),
+                        second.totalBuckets(),
+                        secondFile);
+        ManifestEntry secondDelete =
+                ManifestEntry.create(
+                        FileKind.DELETE,
+                        second.partition(),
+                        second.bucket(),
+                        second.totalBuckets(),
+                        secondFile);
+        ManifestFile manifestFile = createManifestFile(tempDir.toString(), 
Long.MAX_VALUE);
+        ManifestFileMeta firstManifest =
+                writeSingleManifest(manifestFile, Arrays.asList(firstAdd, 
firstDelete));
+        ManifestFileMeta secondManifest =
+                writeSingleManifest(manifestFile, Arrays.asList(secondAdd, 
secondDelete));
+
+        Set<FileEntry.Identifier> deleted =
+                FileEntry.readDeletedEntries(
+                        manifestFile, Arrays.asList(firstManifest, 
secondManifest), 2);
+
+        assertThat(deleted)
+                .containsExactlyInAnyOrder(firstDelete.identifier(), 
secondDelete.identifier());
+    }
+
+    @Test
+    void testReadExpireFileEntriesWithProjectedScan() {
+        ManifestEntry source = gen.next();
+        DataFileMeta file =
+                source.file()
+                        .copy(Arrays.asList("extra-1", "extra-2"))
+                        .copy(new byte[] {1, 2})
+                        .newExternalPath("external/data-file")
+                        .newFirstRowId(10L);
+        List<ManifestEntry> entries =
+                Arrays.asList(
+                        ManifestEntry.create(
+                                FileKind.ADD,
+                                source.partition(),
+                                source.bucket(),
+                                source.totalBuckets(),
+                                file),
+                        ManifestEntry.create(
+                                FileKind.DELETE,
+                                source.partition(),
+                                source.bucket(),
+                                source.totalBuckets(),
+                                file));
+        ManifestFile manifestFile = createManifestFile(tempDir.toString(), 
Long.MAX_VALUE);
+        ManifestFileMeta manifest = writeSingleManifest(manifestFile, entries);
+
+        List<ExpireFileEntry> actual =
+                manifestFile.readExpireFileEntries(manifest.fileName(), 
manifest.fileSize());
+        List<ExpireFileEntry> expected =
+                
entries.stream().map(ExpireFileEntry::from).collect(Collectors.toList());
+
+        assertThat(actual).containsExactlyElementsOf(expected);
+        for (int i = 0; i < actual.size(); i++) {
+            assertThat(actual.get(i).embeddedIndex())
+                    .containsExactly(expected.get(i).embeddedIndex());
+            
assertThat(actual.get(i).fileSource()).isEqualTo(expected.get(i).fileSource());
+        }
+    }
+
     @Test
     void testScanProjectedManifestEntries() throws Exception {
         List<ManifestEntry> entries = Arrays.asList(gen.next(), gen.next(), 
gen.next());

Reply via email to