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