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 077f058631 [core] Save the runtime memory while reading all the 
partitions or buckets (#8745)
077f058631 is described below

commit 077f058631650dccbcd02c19722e097ca6d55265
Author: weijie <[email protected]>
AuthorDate: Tue Jul 21 14:06:26 2026 +0800

    [core] Save the runtime memory while reading all the partitions or buckets 
(#8745)
---
 .../apache/paimon/manifest/ManifestEntryCache.java |  14 +-
 .../paimon/operation/AbstractFileStoreScan.java    |  10 +-
 .../java/org/apache/paimon/utils/ObjectsCache.java |  16 +-
 .../apache/paimon/utils/SimpleObjectsCache.java    |  16 +-
 .../FileStoreScanPartitionBucketEntryTest.java     | 300 +++++++++++++++++++++
 5 files changed, 332 insertions(+), 24 deletions(-)

diff --git 
a/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestEntryCache.java 
b/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestEntryCache.java
index 09afdcc3ac..75af92c75a 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestEntryCache.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/manifest/ManifestEntryCache.java
@@ -112,8 +112,10 @@ public class ManifestEntryCache extends ObjectsCache<Path, 
ManifestEntry, Manife
     }
 
     @Override
-    protected List<ManifestEntry> readFromSegments(
-            ManifestEntrySegments manifestSegments, Filters<ManifestEntry> 
filters)
+    protected <R> List<R> readFromSegments(
+            ManifestEntrySegments manifestSegments,
+            Filters<ManifestEntry> filters,
+            Function<ManifestEntry, R> convertor)
             throws IOException {
         PartitionPredicate partitionFilter = null;
         BucketFilter bucketFilter = null;
@@ -160,12 +162,16 @@ public class ManifestEntryCache extends 
ObjectsCache<Path, ManifestEntry, Manife
         }
 
         // read manifest entries from segments with per record filter
-        List<ManifestEntry> result = new ArrayList<>();
+        List<R> result = new ArrayList<>();
         InternalRowSerializer formatSerializer = this.formatSerializer.get();
         for (Segments subSegments : segmentsList) {
             result.addAll(
                     SimpleObjectsCache.readFromSegments(
-                            formatSerializer, projectedSerializer, 
subSegments, filters));
+                            formatSerializer,
+                            projectedSerializer,
+                            subSegments,
+                            filters,
+                            convertor));
         }
         return result;
     }
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreScan.java
 
b/paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreScan.java
index a78beae34d..a9ef5902ec 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreScan.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreScan.java
@@ -367,7 +367,10 @@ public abstract class AbstractFileStoreScan implements 
FileStoreScan {
         List<ManifestFileMeta> manifests = readManifests().filteredManifests;
         Map<BinaryRow, PartitionEntry> partitions = new ConcurrentHashMap<>();
         Consumer<ManifestFileMeta> processor =
-                m -> 
PartitionEntry.merge(PartitionEntry.merge(readManifest(m)), partitions);
+                m ->
+                        PartitionEntry.merge(
+                                readManifest(m, 
PartitionEntry::fromManifestEntry, null, null),
+                                partitions);
         randomlyOnlyExecute(getExecutorService(parallelism), processor, 
manifests);
         return partitions.values().stream()
                 .filter(p -> p.fileCount() > 0)
@@ -379,7 +382,10 @@ public abstract class AbstractFileStoreScan implements 
FileStoreScan {
         List<ManifestFileMeta> manifests = readManifests().filteredManifests;
         Map<Pair<BinaryRow, Integer>, BucketEntry> buckets = new 
ConcurrentHashMap<>();
         Consumer<ManifestFileMeta> processor =
-                m -> BucketEntry.merge(BucketEntry.merge(readManifest(m)), 
buckets);
+                m ->
+                        BucketEntry.merge(
+                                readManifest(m, 
BucketEntry::fromManifestEntry, null, null),
+                                buckets);
         randomlyOnlyExecute(getExecutorService(parallelism), processor, 
manifests);
         return buckets.values().stream()
                 .filter(p -> p.fileCount() > 0)
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/utils/ObjectsCache.java 
b/paimon-core/src/main/java/org/apache/paimon/utils/ObjectsCache.java
index bfc06d4b63..d978fdbbd1 100644
--- a/paimon-core/src/main/java/org/apache/paimon/utils/ObjectsCache.java
+++ b/paimon-core/src/main/java/org/apache/paimon/utils/ObjectsCache.java
@@ -28,7 +28,6 @@ import javax.annotation.Nullable;
 import javax.annotation.concurrent.ThreadSafe;
 
 import java.io.IOException;
-import java.util.ArrayList;
 import java.util.List;
 import java.util.function.Function;
 
@@ -77,7 +76,7 @@ public abstract class ObjectsCache<K, V, S extends Segments> {
             if (cacheMetrics != null) {
                 cacheMetrics.increaseHitObject();
             }
-            return convert(readFromSegments(segments, filters), convertor);
+            return readFromSegments(segments, filters, convertor);
         } else {
             if (cacheMetrics != null) {
                 cacheMetrics.increaseMissedObject();
@@ -88,7 +87,7 @@ public abstract class ObjectsCache<K, V, S extends Segments> {
             if (fileSize <= cache.maxElementSize()) {
                 segments = createSegments(key, fileSize);
                 cache.put(key, segments);
-                return convert(readFromSegments(segments, filters), convertor);
+                return readFromSegments(segments, filters, convertor);
             } else {
                 return readFromIterator(
                         reader.apply(key, fileSize),
@@ -100,15 +99,8 @@ public abstract class ObjectsCache<K, V, S extends 
Segments> {
         }
     }
 
-    private <R> List<R> convert(List<V> values, Function<V, R> convertor) {
-        List<R> result = new ArrayList<>(values.size());
-        for (V v : values) {
-            result.add(convertor.apply(v));
-        }
-        return result;
-    }
-
-    protected abstract List<V> readFromSegments(S segments, Filters<V> 
filters) throws IOException;
+    protected abstract <R> List<R> readFromSegments(
+            S segments, Filters<V> filters, Function<V, R> convertor) throws 
IOException;
 
     protected abstract S createSegments(K k, @Nullable Long fileSize);
 
diff --git 
a/paimon-core/src/main/java/org/apache/paimon/utils/SimpleObjectsCache.java 
b/paimon-core/src/main/java/org/apache/paimon/utils/SimpleObjectsCache.java
index b3ba2db167..727f30c396 100644
--- a/paimon-core/src/main/java/org/apache/paimon/utils/SimpleObjectsCache.java
+++ b/paimon-core/src/main/java/org/apache/paimon/utils/SimpleObjectsCache.java
@@ -37,6 +37,7 @@ import java.io.EOFException;
 import java.io.IOException;
 import java.util.ArrayList;
 import java.util.List;
+import java.util.function.Function;
 
 /** Cache records to {@link SegmentsCache} by compacted serializer. */
 @ThreadSafe
@@ -58,8 +59,10 @@ public class SimpleObjectsCache<K, V> extends 
ObjectsCache<K, V, Segments> {
     }
 
     @Override
-    protected List<V> readFromSegments(Segments segments, Filters<V> filters) 
throws IOException {
-        return readFromSegments(formatSerializer.get(), projectedSerializer, 
segments, filters);
+    protected <R> List<R> readFromSegments(
+            Segments segments, Filters<V> filters, Function<V, R> convertor) 
throws IOException {
+        return readFromSegments(
+                formatSerializer.get(), projectedSerializer, segments, 
filters, convertor);
     }
 
     @Override
@@ -81,13 +84,14 @@ public class SimpleObjectsCache<K, V> extends 
ObjectsCache<K, V, Segments> {
         }
     }
 
-    public static <V> List<V> readFromSegments(
+    public static <V, R> List<R> readFromSegments(
             InternalRowSerializer formatSerializer,
             ObjectSerializer<V> projectedSerializer,
             Segments segments,
-            Filters<V> filters)
+            Filters<V> filters,
+            Function<V, R> convertor)
             throws IOException {
-        List<V> entries = new ArrayList<>();
+        List<R> entries = new ArrayList<>();
         RandomAccessInputView view = createInputView(segments);
         BinaryRow binaryRow = new BinaryRow(formatSerializer.getArity());
         Filter<InternalRow> readFilter = filters.readFilter();
@@ -98,7 +102,7 @@ public class SimpleObjectsCache<K, V> extends 
ObjectsCache<K, V, Segments> {
                 if (readFilter.test(binaryRow)) {
                     V v = projectedSerializer.fromRow(binaryRow);
                     if (readVFilter.test(v)) {
-                        entries.add(v);
+                        entries.add(convertor.apply(v));
                     }
                 }
             } catch (EOFException e) {
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/operation/FileStoreScanPartitionBucketEntryTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/operation/FileStoreScanPartitionBucketEntryTest.java
new file mode 100644
index 0000000000..8ff5023b16
--- /dev/null
+++ 
b/paimon-core/src/test/java/org/apache/paimon/operation/FileStoreScanPartitionBucketEntryTest.java
@@ -0,0 +1,300 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.paimon.operation;
+
+import org.apache.paimon.data.GenericRow;
+import org.apache.paimon.manifest.BucketEntry;
+import org.apache.paimon.manifest.PartitionEntry;
+import org.apache.paimon.table.sink.BatchTableWrite;
+import org.apache.paimon.table.sink.CommitMessage;
+import org.apache.paimon.table.sink.TableCommitImpl;
+import org.apache.paimon.table.source.snapshot.ScannerTestBase;
+import org.apache.paimon.utils.Pair;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.Collections;
+import java.util.List;
+import java.util.Map;
+import java.util.stream.Collectors;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/**
+ * Tests for {@link AbstractFileStoreScan#readPartitionEntries} and {@link
+ * AbstractFileStoreScan#readBucketEntries}.
+ *
+ * <p>These methods are exercised through {@code 
snapshotReader.partitionEntries()} and {@code
+ * snapshotReader.bucketEntries()} which directly delegate to {@code 
scan.readPartitionEntries()}
+ * and {@code scan.readBucketEntries()}.
+ */
+public class FileStoreScanPartitionBucketEntryTest extends ScannerTestBase {
+
+    @Test
+    public void testReadPartitionEntriesSinglePartition() throws Exception {
+        // Write data to a single partition (pt=1)
+        BatchTableWrite write = table.newWrite(commitUser);
+        for (int i = 0; i < 10; i++) {
+            write.write(GenericRow.of(1, i, (long) i));
+        }
+        List<CommitMessage> messages = write.prepareCommit();
+        TableCommitImpl commit = table.newCommit(commitUser);
+        commit.commit(messages);
+        write.close();
+        commit.close();
+
+        List<PartitionEntry> entries = snapshotReader.partitionEntries();
+        assertThat(entries).hasSize(1);
+        PartitionEntry entry = entries.get(0);
+        assertThat(entry.partition().getInt(0)).isEqualTo(1);
+        assertThat(entry.recordCount()).isEqualTo(10);
+        assertThat(entry.fileCount()).isEqualTo(1);
+        assertThat(entry.totalBuckets()).isEqualTo(1);
+        assertThat(entry.fileSizeInBytes()).isGreaterThan(0);
+    }
+
+    @Test
+    public void testReadPartitionEntriesMultiplePartitions() throws Exception {
+        // Write data to multiple partitions
+        BatchTableWrite write = table.newWrite(commitUser);
+        int[] partitions = {1, 2, 3};
+        int[] recordCounts = {5, 10, 15};
+
+        for (int p = 0; p < partitions.length; p++) {
+            for (int i = 0; i < recordCounts[p]; i++) {
+                write.write(GenericRow.of(partitions[p], i, (long) i));
+            }
+        }
+        List<CommitMessage> messages = write.prepareCommit();
+        TableCommitImpl commit = table.newCommit(commitUser);
+        commit.commit(messages);
+        write.close();
+        commit.close();
+
+        List<PartitionEntry> entries = snapshotReader.partitionEntries();
+        assertThat(entries).hasSize(3);
+
+        Map<Integer, PartitionEntry> entryMap =
+                entries.stream().collect(Collectors.toMap(e -> 
e.partition().getInt(0), e -> e));
+
+        for (int p = 0; p < partitions.length; p++) {
+            PartitionEntry entry = entryMap.get(partitions[p]);
+            assertThat(entry).isNotNull();
+            assertThat(entry.recordCount()).isEqualTo(recordCounts[p]);
+            assertThat(entry.fileCount()).isEqualTo(1);
+        }
+    }
+
+    @Test
+    public void testReadPartitionEntriesMultipleFilesSamePartition() throws 
Exception {
+        // Write data in two commits to the same partition, creating multiple 
files
+        BatchTableWrite write1 = table.newWrite(commitUser);
+        for (int i = 0; i < 5; i++) {
+            write1.write(GenericRow.of(1, i, (long) i));
+        }
+        List<CommitMessage> messages1 = write1.prepareCommit();
+        TableCommitImpl commit1 = table.newCommit(commitUser);
+        commit1.commit(messages1);
+        write1.close();
+
+        BatchTableWrite write2 = table.newWrite(commitUser);
+        for (int i = 5; i < 15; i++) {
+            write2.write(GenericRow.of(1, i, (long) i));
+        }
+        List<CommitMessage> messages2 = write2.prepareCommit();
+        TableCommitImpl commit2 = table.newCommit(commitUser);
+        commit2.commit(messages2);
+        write2.close();
+        commit2.close();
+
+        List<PartitionEntry> entries = snapshotReader.partitionEntries();
+        assertThat(entries).hasSize(1);
+        PartitionEntry entry = entries.get(0);
+        assertThat(entry.partition().getInt(0)).isEqualTo(1);
+        assertThat(entry.recordCount()).isEqualTo(15);
+        assertThat(entry.fileCount()).isEqualTo(2);
+    }
+
+    @Test
+    public void testReadPartitionEntriesWithPartitionFilter() throws Exception 
{
+        // Write data to partitions 1, 2, 3
+        BatchTableWrite write = table.newWrite(commitUser);
+        for (int pt = 1; pt <= 3; pt++) {
+            for (int i = 0; i < 5; i++) {
+                write.write(GenericRow.of(pt, i, (long) i));
+            }
+        }
+        List<CommitMessage> messages = write.prepareCommit();
+        TableCommitImpl commit = table.newCommit(commitUser);
+        commit.commit(messages);
+        write.close();
+        commit.close();
+
+        // Create a new snapshot reader with partition filter to only read 
partition 2
+        List<PartitionEntry> entries =
+                table.newSnapshotReader()
+                        
.withPartitionFilter(Collections.singletonList(binaryRow(2)))
+                        .partitionEntries();
+
+        assertThat(entries).hasSize(1);
+        assertThat(entries.get(0).partition().getInt(0)).isEqualTo(2);
+        assertThat(entries.get(0).recordCount()).isEqualTo(5);
+    }
+
+    @Test
+    public void testReadPartitionEntriesAppendOnlyTable() throws Exception {
+        // Create an append-only table (no primary keys)
+        createAppendOnlyTable();
+
+        BatchTableWrite write = table.newWrite(commitUser);
+        for (int pt = 1; pt <= 3; pt++) {
+            for (int i = 0; i < 10; i++) {
+                write.write(GenericRow.of(pt, i, (long) i));
+            }
+        }
+        List<CommitMessage> messages = write.prepareCommit();
+        TableCommitImpl commit = table.newCommit(commitUser);
+        commit.commit(messages);
+        write.close();
+        commit.close();
+
+        List<PartitionEntry> entries = snapshotReader.partitionEntries();
+        assertThat(entries).hasSize(3);
+
+        Map<Integer, PartitionEntry> entryMap =
+                entries.stream().collect(Collectors.toMap(e -> 
e.partition().getInt(0), e -> e));
+
+        for (int pt = 1; pt <= 3; pt++) {
+            PartitionEntry entry = entryMap.get(pt);
+            assertThat(entry).isNotNull();
+            assertThat(entry.recordCount()).isEqualTo(10);
+        }
+    }
+
+    @Test
+    public void testReadBucketEntriesSinglePartition() throws Exception {
+        // Write data to a single partition with 1 bucket
+        BatchTableWrite write = table.newWrite(commitUser);
+        for (int i = 0; i < 10; i++) {
+            write.write(GenericRow.of(1, i, (long) i));
+        }
+        List<CommitMessage> messages = write.prepareCommit();
+        TableCommitImpl commit = table.newCommit(commitUser);
+        commit.commit(messages);
+        write.close();
+        commit.close();
+
+        List<BucketEntry> entries = snapshotReader.bucketEntries();
+        assertThat(entries).hasSize(1);
+        BucketEntry entry = entries.get(0);
+        assertThat(entry.partition().getInt(0)).isEqualTo(1);
+        assertThat(entry.bucket()).isEqualTo(0);
+        assertThat(entry.recordCount()).isEqualTo(10);
+        assertThat(entry.fileCount()).isEqualTo(1);
+    }
+
+    @Test
+    public void testReadBucketEntriesMultiplePartitions() throws Exception {
+        // Write data to multiple partitions
+        BatchTableWrite write = table.newWrite(commitUser);
+        for (int pt = 1; pt <= 3; pt++) {
+            for (int i = 0; i < 5; i++) {
+                write.write(GenericRow.of(pt, i, (long) i));
+            }
+        }
+        List<CommitMessage> messages = write.prepareCommit();
+        TableCommitImpl commit = table.newCommit(commitUser);
+        commit.commit(messages);
+        write.close();
+        commit.close();
+
+        List<BucketEntry> entries = snapshotReader.bucketEntries();
+        assertThat(entries).hasSize(3);
+
+        Map<Pair<Integer, Integer>, BucketEntry> entryMap =
+                entries.stream()
+                        .collect(
+                                Collectors.toMap(
+                                        e -> Pair.of(e.partition().getInt(0), 
e.bucket()), e -> e));
+
+        for (int pt = 1; pt <= 3; pt++) {
+            BucketEntry entry = entryMap.get(Pair.of(pt, 0));
+            assertThat(entry).isNotNull();
+            assertThat(entry.recordCount()).isEqualTo(5);
+            assertThat(entry.fileCount()).isEqualTo(1);
+        }
+    }
+
+    @Test
+    public void testReadBucketEntriesMultipleFilesSameBucket() throws 
Exception {
+        // Write data in two commits to the same partition/bucket
+        BatchTableWrite write1 = table.newWrite(commitUser);
+        for (int i = 0; i < 5; i++) {
+            write1.write(GenericRow.of(1, i, (long) i));
+        }
+        List<CommitMessage> messages1 = write1.prepareCommit();
+        TableCommitImpl commit1 = table.newCommit(commitUser);
+        commit1.commit(messages1);
+        write1.close();
+
+        BatchTableWrite write2 = table.newWrite(commitUser);
+        for (int i = 5; i < 15; i++) {
+            write2.write(GenericRow.of(1, i, (long) i));
+        }
+        List<CommitMessage> messages2 = write2.prepareCommit();
+        TableCommitImpl commit2 = table.newCommit(commitUser);
+        commit2.commit(messages2);
+        write2.close();
+        commit2.close();
+
+        List<BucketEntry> entries = snapshotReader.bucketEntries();
+        assertThat(entries).hasSize(1);
+        BucketEntry entry = entries.get(0);
+        assertThat(entry.partition().getInt(0)).isEqualTo(1);
+        assertThat(entry.bucket()).isEqualTo(0);
+        assertThat(entry.recordCount()).isEqualTo(15);
+        assertThat(entry.fileCount()).isEqualTo(2);
+    }
+
+    @Test
+    public void testReadBucketEntriesAppendOnlyTable() throws Exception {
+        // Create an append-only table (no primary keys)
+        createAppendOnlyTable();
+
+        BatchTableWrite write = table.newWrite(commitUser);
+        for (int pt = 1; pt <= 3; pt++) {
+            for (int i = 0; i < 10; i++) {
+                write.write(GenericRow.of(pt, i, (long) i));
+            }
+        }
+        List<CommitMessage> messages = write.prepareCommit();
+        TableCommitImpl commit = table.newCommit(commitUser);
+        commit.commit(messages);
+        write.close();
+        commit.close();
+
+        List<BucketEntry> entries = snapshotReader.bucketEntries();
+        assertThat(entries).hasSize(3);
+
+        for (BucketEntry entry : entries) {
+            assertThat(entry.recordCount()).isEqualTo(10);
+            assertThat(entry.fileCount()).isEqualTo(1);
+        }
+    }
+}

Reply via email to