Skip to content

Commit 09a0258

Browse files
authored
[spark] Avoid scanning partition entries without partition_idle_time (#8932)
1 parent 6557de3 commit 09a0258

2 files changed

Lines changed: 50 additions & 24 deletions

File tree

paimon-spark/paimon-spark-common/src/main/java/org/apache/paimon/spark/procedure/CompactProcedure.java

Lines changed: 24 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -322,13 +322,17 @@ private void compactAwareBucketTable(
322322
if (partitionPredicate != null) {
323323
snapshotReader.withPartitionFilter(partitionPredicate);
324324
}
325+
boolean filterByPartitionIdleTime = partitionIdleTime != null;
325326
Set<BinaryRow> partitionToBeCompacted =
326-
getHistoryPartition(snapshotReader, partitionIdleTime);
327+
getPartitionsToCompact(snapshotReader, partitionIdleTime);
327328
List<Pair<byte[], Integer>> partitionBuckets =
328329
snapshotReader.bucketEntries().stream()
329330
.map(entry -> Pair.of(entry.partition(), entry.bucket()))
330331
.distinct()
331-
.filter(pair -> partitionToBeCompacted.contains(pair.getKey()))
332+
.filter(
333+
pair ->
334+
!filterByPartitionIdleTime
335+
|| partitionToBeCompacted.contains(pair.getKey()))
332336
.map(
333337
p ->
334338
Pair.of(
@@ -615,29 +619,25 @@ private static List<CommitMessage> deserializeCommitMessagesAndReleaseSerialized
615619
return messages;
616620
}
617621

618-
private Set<BinaryRow> getHistoryPartition(
622+
static Set<BinaryRow> getPartitionsToCompact(
619623
SnapshotReader snapshotReader, @Nullable Duration partitionIdleTime) {
620-
Set<Pair<BinaryRow, Long>> partitionInfo =
621-
snapshotReader.partitionEntries().stream()
622-
.map(
623-
partitionEntry ->
624-
Pair.of(
625-
partitionEntry.partition(),
626-
partitionEntry.lastFileCreationTime()))
627-
.collect(Collectors.toSet());
628-
if (partitionIdleTime != null) {
629-
long historyMilli =
630-
LocalDateTime.now()
631-
.minus(partitionIdleTime)
632-
.atZone(ZoneId.systemDefault())
633-
.toInstant()
634-
.toEpochMilli();
635-
partitionInfo =
636-
partitionInfo.stream()
637-
.filter(partition -> partition.getValue() <= historyMilli)
638-
.collect(Collectors.toSet());
639-
}
640-
return partitionInfo.stream().map(Pair::getKey).collect(Collectors.toSet());
624+
return partitionIdleTime == null
625+
? Collections.emptySet()
626+
: getHistoryPartition(snapshotReader, partitionIdleTime);
627+
}
628+
629+
private static Set<BinaryRow> getHistoryPartition(
630+
SnapshotReader snapshotReader, Duration partitionIdleTime) {
631+
long historyMilli =
632+
LocalDateTime.now()
633+
.minus(partitionIdleTime)
634+
.atZone(ZoneId.systemDefault())
635+
.toInstant()
636+
.toEpochMilli();
637+
return snapshotReader.partitionEntries().stream()
638+
.filter(partition -> partition.lastFileCreationTime() <= historyMilli)
639+
.map(PartitionEntry::partition)
640+
.collect(Collectors.toSet());
641641
}
642642

643643
private void sortCompactUnAwareBucketTable(

paimon-spark/paimon-spark-ut/src/test/scala/org/apache/paimon/spark/procedure/CompactProcedureTestBase.scala

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@ import org.apache.paimon.spark.PaimonSparkTestBase
2424
import org.apache.paimon.spark.utils.SparkProcedureUtils
2525
import org.apache.paimon.table.FileStoreTable
2626
import org.apache.paimon.table.source.DataSplit
27+
import org.apache.paimon.table.source.snapshot.SnapshotReader
2728

2829
import org.apache.spark.scheduler.{SparkListener, SparkListenerStageSubmitted}
2930
import org.apache.spark.sql.{Dataset, Row}
@@ -33,7 +34,9 @@ import org.assertj.core.api.Assertions
3334
import org.assertj.core.api.Assertions.assertThatThrownBy
3435
import org.scalatest.time.Span
3536

37+
import java.lang.reflect.{InvocationHandler, Method, Proxy}
3638
import java.util
39+
import java.util.concurrent.atomic.AtomicBoolean
3740

3841
import scala.collection.JavaConverters._
3942
import scala.util.Random
@@ -45,6 +48,29 @@ abstract class CompactProcedureTestBase extends PaimonSparkTestBase with StreamT
4548

4649
// ----------------------- Minor Compact -----------------------
4750

51+
test("Paimon Procedure: skip partition entries scan without partition idle time") {
52+
val partitionEntriesScanned = new AtomicBoolean(false)
53+
val snapshotReader = Proxy
54+
.newProxyInstance(
55+
classOf[SnapshotReader].getClassLoader,
56+
Array(classOf[SnapshotReader]),
57+
new InvocationHandler {
58+
override def invoke(proxy: Any, method: Method, args: Array[AnyRef]): AnyRef = {
59+
if (method.getName == "partitionEntries") {
60+
partitionEntriesScanned.set(true)
61+
}
62+
null
63+
}
64+
}
65+
)
66+
.asInstanceOf[SnapshotReader]
67+
68+
val partitions = CompactProcedure.getPartitionsToCompact(snapshotReader, null)
69+
70+
Assertions.assertThat(partitions.isEmpty).isTrue
71+
Assertions.assertThat(partitionEntriesScanned.get()).isFalse
72+
}
73+
4874
test("Paimon Procedure: compact aware bucket pk table with minor compact strategy") {
4975
withTable("T") {
5076
spark.sql(s"""

0 commit comments

Comments
 (0)