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 93f6d34f53 [spark] Support skipping expired partitions during
compaction (#9021)
93f6d34f53 is described below
commit 93f6d34f53e1410984c1d63f2e7e860d9cdf078a
Author: sanshi <[email protected]>
AuthorDate: Wed Aug 5 23:18:25 2026 +0800
[spark] Support skipping expired partitions during compaction (#9021)
---
docs/docs/maintenance/dedicated-compaction.mdx | 25 +++++
.../paimon/spark/procedure/CompactProcedure.java | 21 ++++
.../spark/procedure/CompactProcedureTestBase.scala | 114 +++++++++++++++++++++
3 files changed, 160 insertions(+)
diff --git a/docs/docs/maintenance/dedicated-compaction.mdx
b/docs/docs/maintenance/dedicated-compaction.mdx
index 7695b09970..c2f6dc4fb1 100644
--- a/docs/docs/maintenance/dedicated-compaction.mdx
+++ b/docs/docs/maintenance/dedicated-compaction.mdx
@@ -489,6 +489,17 @@ CALL sys.compact(`table` => 'default.T', options =>
'compaction.skip-expired-par
</TabItem>
+<TabItem value="spark-sql" label="Spark SQL">
+
+Run the following sql:
+
+```sql
+-- skip expired partitions compact table
+CALL sys.compact(table => 'default.T', options =>
'compaction.skip-expired-partitions=true')
+```
+
+</TabItem>
+
<TabItem value="flink-action-jar" label="Flink Action Jar">
```bash
@@ -525,6 +536,20 @@ CALL sys.compact_database(
</TabItem>
+<TabItem value="spark-sql" label="Spark SQL">
+
+Run the following sql:
+
+```sql
+-- skip expired partitions compact database
+CALL sys.compact_database(
+ including_databases => 'default',
+ options => 'compaction.skip-expired-partitions=true'
+)
+```
+
+</TabItem>
+
<TabItem value="flink-action-jar" label="Flink Action Jar">
```bash
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 1e156ca42a..b73c0b189a 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
@@ -38,6 +38,7 @@ import org.apache.paimon.io.DataIncrement;
import org.apache.paimon.manifest.PartitionEntry;
import org.apache.paimon.operation.BaseAppendFileStoreWrite;
import org.apache.paimon.partition.PartitionPredicate;
+import org.apache.paimon.partition.PartitionValuesTimeExpireStrategy;
import org.apache.paimon.spark.SparkUtils;
import org.apache.paimon.spark.commands.PaimonSparkWriter;
import org.apache.paimon.spark.sort.TableSorter;
@@ -94,6 +95,7 @@ import java.util.Iterator;
import java.util.List;
import java.util.Map;
import java.util.Set;
+import java.util.function.Predicate;
import java.util.stream.Collectors;
import scala.collection.JavaConverters;
@@ -325,6 +327,7 @@ public class CompactProcedure extends BaseProcedure {
boolean filterByPartitionIdleTime = partitionIdleTime != null;
Set<BinaryRow> partitionToBeCompacted =
getPartitionsToCompact(snapshotReader, partitionIdleTime);
+ Predicate<BinaryRow> shouldCompactPartition =
nonExpiredPartitionPredicate(table);
List<Pair<byte[], Integer>> partitionBuckets =
snapshotReader.bucketEntries().stream()
.map(entry -> Pair.of(entry.partition(),
entry.bucket()))
@@ -333,6 +336,7 @@ public class CompactProcedure extends BaseProcedure {
pair ->
!filterByPartitionIdleTime
||
partitionToBeCompacted.contains(pair.getKey()))
+ .filter(pair ->
shouldCompactPartition.test(pair.getKey()))
.map(
p ->
Pair.of(
@@ -395,6 +399,23 @@ public class CompactProcedure extends BaseProcedure {
}
}
+ private static Predicate<BinaryRow>
nonExpiredPartitionPredicate(FileStoreTable table) {
+ CoreOptions options = table.coreOptions();
+ if (!options.compactionSkipExpiredPartitions()
+ || options.partitionExpireTime() == null
+ || !CoreOptions.PartitionExpireStrategy.VALUES_TIME
+ .toString()
+ .equals(options.partitionExpireStrategy())) {
+ return partition -> true;
+ }
+
+ LocalDateTime expireDateTime =
LocalDateTime.now().minus(options.partitionExpireTime());
+ PartitionValuesTimeExpireStrategy expireStrategy =
+ new PartitionValuesTimeExpireStrategy(
+ options, table.schema().logicalPartitionType());
+ return partition -> !expireStrategy.isExpired(expireDateTime,
partition);
+ }
+
private void compactUnAwareBucketTable(
FileStoreTable table,
@Nullable PartitionPredicate partitionPredicate,
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 6bc1a898bc..8fccd1b568 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
@@ -35,6 +35,8 @@ import org.assertj.core.api.Assertions.assertThatThrownBy
import org.scalatest.time.Span
import java.lang.reflect.{InvocationHandler, Method, Proxy}
+import java.time.LocalDate
+import java.time.format.DateTimeFormatter
import java.util
import java.util.concurrent.atomic.AtomicBoolean
@@ -658,6 +660,118 @@ abstract class CompactProcedureTestBase extends
PaimonSparkTestBase with StreamT
Row(5, "e", "p1") :: Row(6, "f", "p2") :: Nil)
}
+ test("Paimon Procedure: compact skips expired partitions for aware bucket
table") {
+ Seq(1, -1).foreach {
+ bucket =>
+ withClue(s"bucket=$bucket") {
+ withTable("T") {
+ createPartitionExpireTable(bucket, "values-time", endInputCheck =
false)
+ val table = loadTable("T")
+ val (expiredDt, activeDt) = writeExpiredAndActivePartitions()
+
+ spark.sql(
+ "CALL sys.compact(table => 'T', " +
+ "options => 'compaction.skip-expired-partitions=true')")
+
+
Assertions.assertThat(lastSnapshotCommand(table)).isEqualTo(CommitKind.COMPACT)
+ val fileCounts = partitionFileCounts(table)
+ Assertions.assertThat(fileCounts(expiredDt)).isEqualTo(2)
+ Assertions.assertThat(fileCounts(activeDt)).isEqualTo(1)
+ }
+ }
+ }
+ }
+
+ test("Paimon Procedure: compact does not skip expired partitions by
default") {
+ withTable("T") {
+ createPartitionExpireTable(1, "values-time", endInputCheck = false)
+ val table = loadTable("T")
+ val (expiredDt, activeDt) = writeExpiredAndActivePartitions()
+
+ spark.sql("CALL sys.compact(table => 'T')")
+
+ val fileCounts = partitionFileCounts(table)
+ Assertions.assertThat(fileCounts(expiredDt)).isEqualTo(1)
+ Assertions.assertThat(fileCounts(activeDt)).isEqualTo(1)
+ }
+ }
+
+ test("Paimon Procedure: compact skip expired partitions ignores update-time
strategy") {
+ withTable("T") {
+ createPartitionExpireTable(1, "update-time", endInputCheck = false)
+ val table = loadTable("T")
+ val (expiredDt, activeDt) = writeExpiredAndActivePartitions()
+
+ spark.sql(
+ "CALL sys.compact(table => 'T', " +
+ "options => 'compaction.skip-expired-partitions=true')")
+
+ val fileCounts = partitionFileCounts(table)
+ Assertions.assertThat(fileCounts(expiredDt)).isEqualTo(1)
+ Assertions.assertThat(fileCounts(activeDt)).isEqualTo(1)
+ }
+ }
+
+ test("Paimon Procedure: compact end input still expires skipped partitions")
{
+ withTable("T") {
+ createPartitionExpireTable(1, "values-time", endInputCheck = false)
+ val table = loadTable("T")
+ val (_, activeDt) = writeExpiredAndActivePartitions()
+
+ spark.sql(
+ "ALTER TABLE T SET TBLPROPERTIES (" +
+ "'end-input.check-partition-expire'='true')")
+
+ spark.sql(
+ "CALL sys.compact(table => 'T', " +
+ "options => 'compaction.skip-expired-partitions=true')")
+
+ val fileCounts = partitionFileCounts(table)
+ Assertions.assertThat(fileCounts.asJava).containsOnlyKeys(activeDt)
+ Assertions.assertThat(fileCounts(activeDt)).isEqualTo(1)
+ }
+ }
+
+ private def createPartitionExpireTable(
+ bucket: Int,
+ expirationStrategy: String,
+ endInputCheck: Boolean): Unit = {
+ val dynamicBucketOption =
+ if (bucket == -1) ", 'dynamic-bucket.initial-buckets'='1'" else ""
+ spark.sql(s"""
+ |CREATE TABLE T (id INT, value STRING, dt STRING)
+ |TBLPROPERTIES (
+ | 'primary-key'='id, dt',
+ | 'bucket'='$bucket',
+ | 'write-only'='true',
+ | 'partition.expiration-time'='7 d',
+ | 'partition.expiration-strategy'='$expirationStrategy',
+ | 'partition.timestamp-formatter'='yyyyMMdd',
+ | 'partition.expiration-check-interval'='999 d',
+ | 'end-input.check-partition-expire'='$endInputCheck'
+ | $dynamicBucketOption)
+ |PARTITIONED BY (dt)
+ |""".stripMargin)
+ }
+
+ private def writeExpiredAndActivePartitions(): (String, String) = {
+ val formatter = DateTimeFormatter.ofPattern("yyyyMMdd")
+ val expiredDt = LocalDate.now.minusDays(30).format(formatter)
+ val activeDt = LocalDate.now.format(formatter)
+ spark.sql(s"INSERT INTO T VALUES (1, 'old', '$expiredDt'), (1, 'new',
'$activeDt')")
+ spark.sql(s"INSERT INTO T VALUES (2, 'old', '$expiredDt'), (2, 'new',
'$activeDt')")
+ (expiredDt, activeDt)
+ }
+
+ private def partitionFileCounts(table: FileStoreTable): Map[String, Int] = {
+ table.newSnapshotReader.read.dataSplits.asScala
+ .groupBy(_.partition().getString(0).toString)
+ .map {
+ case (partition, splits) =>
+ partition -> splits.map(_.dataFiles().size()).sum
+ }
+ }
+
test("Paimon Procedure: compact with partition_idle_time for pk table") {
Seq(1, -1).foreach(
bucket => {