From 82c9836a0b4bd6e35ddda8b54960693639494ea2 Mon Sep 17 00:00:00 2001 From: sanshi <1715734693@qq.com> Date: Thu, 30 Jul 2026 17:33:18 +0800 Subject: [PATCH 1/3] [spark] Avoid scanning partition entries without partition idle time --- .../spark/procedure/CompactProcedure.java | 48 +++++++++---------- .../procedure/CompactProcedureTestBase.scala | 25 ++++++++++ 2 files changed, 49 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 13a71d773c90..3c47563b1a4c 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 @@ private void compactAwareBucketTable( if (partitionPredicate != null) { snapshotReader.withPartitionFilter(partitionPredicate); } + boolean filterByPartitionIdleTime = partitionIdleTime != null; Set partitionToBeCompacted = - getHistoryPartition(snapshotReader, partitionIdleTime); + getPartitionsToCompact(snapshotReader, partitionIdleTime); List> 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( @@ -603,29 +607,25 @@ private static List deserializeCommitMessagesAndReleaseSerialized return messages; } - private Set getHistoryPartition( + static Set getPartitionsToCompact( SnapshotReader snapshotReader, @Nullable Duration partitionIdleTime) { - Set> 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 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 76218f19efda..7584049b9098 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,28 @@ 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""" From d8f07bc26962fcc547c87216b9b7cef150770ea8 Mon Sep 17 00:00:00 2001 From: sanshi <1715734693@qq.com> Date: Thu, 30 Jul 2026 19:19:39 +0800 Subject: [PATCH 2/3] fix code format --- .../paimon/spark/procedure/CompactProcedureTestBase.scala | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) 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 7584049b9098..6bc1a898bc44 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 @@ -61,7 +61,8 @@ abstract class CompactProcedureTestBase extends PaimonSparkTestBase with StreamT } null } - }) + } + ) .asInstanceOf[SnapshotReader] val partitions = CompactProcedure.getPartitionsToCompact(snapshotReader, null) From 3d6edb14e9e72abb9cb837eed5a319ea8221dc11 Mon Sep 17 00:00:00 2001 From: sanshi <1715734693@qq.com> Date: Thu, 30 Jul 2026 21:59:17 +0800 Subject: [PATCH 3/3] trigger ci