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 => {

Reply via email to