Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -322,13 +322,17 @@ private void compactAwareBucketTable(
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(
Expand Down Expand Up @@ -603,29 +607,25 @@ private static List<CommitMessage> deserializeCommitMessagesAndReleaseSerialized
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(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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}
Expand All @@ -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
Expand All @@ -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"""
Expand Down
Loading