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..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 @@ -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"""