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 09a0258c7f [spark] Avoid scanning partition entries without 
partition_idle_time (#8932)
09a0258c7f is described below

commit 09a0258c7f2be5ce4ec3b430230e7d46a3be05b2
Author: sanshi <[email protected]>
AuthorDate: Mon Aug 3 19:11:50 2026 +0800

    [spark] Avoid scanning partition entries without partition_idle_time (#8932)
---
 .../paimon/spark/procedure/CompactProcedure.java   | 48 +++++++++++-----------
 .../spark/procedure/CompactProcedureTestBase.scala | 26 ++++++++++++
 2 files changed, 50 insertions(+), 24 deletions(-)

diff --git 
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CompactProcedure.java
 
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CompactProcedure.java
index 59a92cd269..27dd435aba 100644
--- 
a/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CompactProcedure.java
+++ 
b/paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CompactProcedure.java
@@ -322,13 +322,17 @@ public class CompactProcedure extends BaseProcedure {
         if (partitionPredicate != null) {
             snapshotReader.withPartitionFilter(partitionPredicate);
         }
+        boolean filterByPartitionIdleTime = partitionIdleTime != null;
         Set<BinaryRow> partitionToBeCompacted =
-                getHistoryPartition(snapshotReader, partitionIdleTime);
+                getPartitionsToCompact(snapshotReader, partitionIdleTime);
         List<Pair<byte[], Integer>> partitionBuckets =
                 snapshotReader.bucketEntries().stream()
                         .map(entry -> Pair.of(entry.partition(), 
entry.bucket()))
                         .distinct()
-                        .filter(pair -> 
partitionToBeCompacted.contains(pair.getKey()))
+                        .filter(
+                                pair ->
+                                        !filterByPartitionIdleTime
+                                                || 
partitionToBeCompacted.contains(pair.getKey()))
                         .map(
                                 p ->
                                         Pair.of(
@@ -615,29 +619,25 @@ public class CompactProcedure extends BaseProcedure {
         return messages;
     }
 
-    private Set<BinaryRow> getHistoryPartition(
+    static Set<BinaryRow> getPartitionsToCompact(
             SnapshotReader snapshotReader, @Nullable Duration 
partitionIdleTime) {
-        Set<Pair<BinaryRow, Long>> partitionInfo =
-                snapshotReader.partitionEntries().stream()
-                        .map(
-                                partitionEntry ->
-                                        Pair.of(
-                                                partitionEntry.partition(),
-                                                
partitionEntry.lastFileCreationTime()))
-                        .collect(Collectors.toSet());
-        if (partitionIdleTime != null) {
-            long historyMilli =
-                    LocalDateTime.now()
-                            .minus(partitionIdleTime)
-                            .atZone(ZoneId.systemDefault())
-                            .toInstant()
-                            .toEpochMilli();
-            partitionInfo =
-                    partitionInfo.stream()
-                            .filter(partition -> partition.getValue() <= 
historyMilli)
-                            .collect(Collectors.toSet());
-        }
-        return 
partitionInfo.stream().map(Pair::getKey).collect(Collectors.toSet());
+        return partitionIdleTime == null
+                ? Collections.emptySet()
+                : getHistoryPartition(snapshotReader, partitionIdleTime);
+    }
+
+    private static Set<BinaryRow> getHistoryPartition(
+            SnapshotReader snapshotReader, Duration partitionIdleTime) {
+        long historyMilli =
+                LocalDateTime.now()
+                        .minus(partitionIdleTime)
+                        .atZone(ZoneId.systemDefault())
+                        .toInstant()
+                        .toEpochMilli();
+        return snapshotReader.partitionEntries().stream()
+                .filter(partition -> partition.lastFileCreationTime() <= 
historyMilli)
+                .map(PartitionEntry::partition)
+                .collect(Collectors.toSet());
     }
 
     private void sortCompactUnAwareBucketTable(
diff --git 
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CompactProcedureTestBase.scala
 
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CompactProcedureTestBase.scala
index 76218f19ef..6bc1a898bc 100644
--- 
a/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CompactProcedureTestBase.scala
+++ 
b/paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CompactProcedureTestBase.scala
@@ -24,6 +24,7 @@ import org.apache.paimon.spark.PaimonSparkTestBase
 import org.apache.paimon.spark.utils.SparkProcedureUtils
 import org.apache.paimon.table.FileStoreTable
 import org.apache.paimon.table.source.DataSplit
+import org.apache.paimon.table.source.snapshot.SnapshotReader
 
 import org.apache.spark.scheduler.{SparkListener, SparkListenerStageSubmitted}
 import org.apache.spark.sql.{Dataset, Row}
@@ -33,7 +34,9 @@ import org.assertj.core.api.Assertions
 import org.assertj.core.api.Assertions.assertThatThrownBy
 import org.scalatest.time.Span
 
+import java.lang.reflect.{InvocationHandler, Method, Proxy}
 import java.util
+import java.util.concurrent.atomic.AtomicBoolean
 
 import scala.collection.JavaConverters._
 import scala.util.Random
@@ -45,6 +48,29 @@ abstract class CompactProcedureTestBase extends 
PaimonSparkTestBase with StreamT
 
   // ----------------------- Minor Compact -----------------------
 
+  test("Paimon Procedure: skip partition entries scan without partition idle 
time") {
+    val partitionEntriesScanned = new AtomicBoolean(false)
+    val snapshotReader = Proxy
+      .newProxyInstance(
+        classOf[SnapshotReader].getClassLoader,
+        Array(classOf[SnapshotReader]),
+        new InvocationHandler {
+          override def invoke(proxy: Any, method: Method, args: 
Array[AnyRef]): AnyRef = {
+            if (method.getName == "partitionEntries") {
+              partitionEntriesScanned.set(true)
+            }
+            null
+          }
+        }
+      )
+      .asInstanceOf[SnapshotReader]
+
+    val partitions = CompactProcedure.getPartitionsToCompact(snapshotReader, 
null)
+
+    Assertions.assertThat(partitions.isEmpty).isTrue
+    Assertions.assertThat(partitionEntriesScanned.get()).isFalse
+  }
+
   test("Paimon Procedure: compact aware bucket pk table with minor compact 
strategy") {
     withTable("T") {
       spark.sql(s"""

Reply via email to