This is an automated email from the ASF dual-hosted git repository.

voonhous pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/hudi.git


The following commit(s) were added to refs/heads/master by this push:
     new a806d3968488 [HUDI-8970] Change scheduleandexecute in 
RunCompactionProcedure to compact all pending plans too (#12794)
a806d3968488 is described below

commit a806d3968488a05a75c756e4dd436fde2f717755
Author: voonhous <[email protected]>
AuthorDate: Sat Jul 18 19:56:39 2026 +0800

    [HUDI-8970] Change scheduleandexecute in RunCompactionProcedure to compact 
all pending plans too (#12794)
---
 .../command/procedures/HoodieProcedureUtils.scala  |  3 +-
 .../procedures/RunCompactionProcedure.scala        | 12 +++--
 .../hudi/procedure/TestCompactionProcedure.scala   | 54 ++++++++++++++++++++--
 3 files changed, 60 insertions(+), 9 deletions(-)

diff --git 
a/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/command/procedures/HoodieProcedureUtils.scala
 
b/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/command/procedures/HoodieProcedureUtils.scala
index f433c58179b5..30dbbb0f16a6 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/command/procedures/HoodieProcedureUtils.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/command/procedures/HoodieProcedureUtils.scala
@@ -80,8 +80,7 @@ object HoodieProcedureUtils {
   }
 
   /**
-   * scheduleAndExecute: schedule a new plan and then execute it, if no plan 
is generated during
-   * schedule, execute all pending plans
+   * scheduleAndExecute: schedule a new plan and then execute all pending 
plans regardless of plan scheduling outcome
    */
   case object ScheduleAndExecute extends Operation {
     override def value: String = "scheduleandexecute"
diff --git 
a/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/command/procedures/RunCompactionProcedure.scala
 
b/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/command/procedures/RunCompactionProcedure.scala
index 0f881cdc8547..772d8be1b5b7 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/command/procedures/RunCompactionProcedure.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/command/procedures/RunCompactionProcedure.scala
@@ -72,7 +72,7 @@ class RunCompactionProcedure extends BaseProcedure with 
ProcedureBuilder with Sp
       confs = confs ++ 
HoodieCLIUtils.extractOptions(getArgValueOrDefault(args, 
PARAMETERS(4)).get.asInstanceOf[String])
     }
     var specificInstants = getArgValueOrDefault(args, PARAMETERS(5))
-    val limit = getArgValueOrDefault(args, PARAMETERS(6))
+    val limit = getArgValueOrDefault(args, 
PARAMETERS(6)).asInstanceOf[Option[Int]]
 
     // For old version compatibility
     if (op.equals("run")) {
@@ -92,7 +92,7 @@ class RunCompactionProcedure extends BaseProcedure with 
ProcedureBuilder with Sp
       .toSeq.sortBy(f => f)
 
     var (filteredPendingCompactionInstants, operation) = 
HoodieProcedureUtils.filterPendingInstantsAndGetOperation(
-      pendingCompactionInstants, 
specificInstants.asInstanceOf[Option[String]], Option(op), 
limit.asInstanceOf[Option[Int]])
+      pendingCompactionInstants, 
specificInstants.asInstanceOf[Option[String]], Option(op), limit)
 
     var client: SparkRDDWriteClient[_] = null
     try {
@@ -108,10 +108,16 @@ class RunCompactionProcedure extends BaseProcedure with 
ProcedureBuilder with Sp
       if (operation.isSchedule) {
         val instantTime = 
client.scheduleCompaction(HOption.empty[java.util.Map[String, String]])
         instantTime.ifPresent(instant => {
-          filteredPendingCompactionInstants = Seq(instant)
+          // If schedule ONLY, return the instant scheduled else, return 
pending compaction instants too
+          filteredPendingCompactionInstants = if (operation.isExecute) 
filteredPendingCompactionInstants :+ instant else Seq(instant)
         })
       }
 
+      filteredPendingCompactionInstants = if (limit.isDefined) {
+        filteredPendingCompactionInstants.take(limit.get)
+      } else {
+        filteredPendingCompactionInstants
+      }
       logInfo(s"Compaction instants to run: 
${filteredPendingCompactionInstants.mkString(",")}.")
 
       if (operation.isExecute) {
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/procedure/TestCompactionProcedure.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/procedure/TestCompactionProcedure.scala
index fcc3cbdd9430..0624e370fdd4 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/procedure/TestCompactionProcedure.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/procedure/TestCompactionProcedure.scala
@@ -344,7 +344,7 @@ class TestCompactionProcedure extends 
HoodieSparkProcedureTestBase {
     }
   }
 
-  test("Test Call run_clustering with limit parameter") {
+  test("Test Call run_compaction with limit parameter") {
     withSQLConf("hoodie.compact.inline" -> "false", 
"hoodie.compact.inline.max.delta.commits" -> "1") {
       withTempDir { tmp =>
         val tableName = generateTableName
@@ -368,7 +368,7 @@ class TestCompactionProcedure extends 
HoodieSparkProcedureTestBase {
         val conf = new Configuration
         val metaClient = HoodieTestUtils.createMetaClient(new 
HadoopStorageConfiguration(conf), basePath)
 
-        assert(0 == 
metaClient.getActiveTimeline.getCompletedReplaceTimeline.getInstants.size())
+        
assert(metaClient.getActiveTimeline.getCompletedReplaceTimeline.getInstants.size()
 == 0)
         
assert(metaClient.getActiveTimeline.filterPendingClusteringTimeline().empty())
 
         spark.sql(s"insert into $tableName values(1, 'a1', 10, 1000)")
@@ -381,12 +381,58 @@ class TestCompactionProcedure extends 
HoodieSparkProcedureTestBase {
 
         spark.sql(s"call run_compaction(table => '$tableName', op => 
'schedule')")
         metaClient.reloadActiveTimeline();
-        assert(2 == 
metaClient.getActiveTimeline.filterPendingCompactionTimeline().getInstants.size())
+        
assert(metaClient.getActiveTimeline.filterPendingCompactionTimeline().getInstants.size()
 == 2)
 
         spark.sql(s"call run_compaction(table => '$tableName', op => 
'execute', limit => 1)");
 
         metaClient.reloadActiveTimeline();
-        assert(1 == 
metaClient.getActiveTimeline.filterPendingCompactionTimeline().getInstants.size())
+        
assert(metaClient.getActiveTimeline.filterPendingCompactionTimeline().getInstants.size()
 == 1)
+      }
+    }
+  }
+
+  test("Test Call run_compaction with scheduleandexecute op") {
+    withSQLConf("hoodie.compact.inline" -> "false", 
"hoodie.compact.inline.max.delta.commits" -> "1") {
+      withTempDir { tmp =>
+        val tableName = generateTableName
+        val basePath = s"${tmp.getCanonicalPath}/$tableName"
+        spark.sql(
+          s"""
+             |create table $tableName (
+             |  id int,
+             |  name string,
+             |  price double,
+             |  ts long
+             |) using hudi
+             | tblproperties (
+             |  type = 'mor',
+             |  primaryKey = 'id',
+             |  preCombineField = 'ts'
+             | )
+             | location '${basePath}'
+       """.stripMargin)
+
+        val conf = new Configuration
+        val metaClient = HoodieTestUtils.createMetaClient(new 
HadoopStorageConfiguration(conf), basePath)
+
+        
assert(metaClient.getActiveTimeline.getCompletedReplaceTimeline.getInstants.size()
 == 0)
+        
assert(metaClient.getActiveTimeline.filterPendingClusteringTimeline().empty())
+
+        spark.sql(s"insert into $tableName values(1, 'a1', 10, 1000)")
+        spark.sql(s"update $tableName set name = 'a2' where id = 1")
+
+        spark.sql(s"call run_compaction(table => '$tableName', op => 
'schedule')")
+
+        metaClient.reloadActiveTimeline();
+        
assert(metaClient.getActiveTimeline.filterPendingCompactionTimeline().getInstants.size()
 == 1)
+
+        spark.sql(s"insert into $tableName values(2, 'b1', 20, 3000)")
+        spark.sql(s"update $tableName set name = 'b3' where id = 2")
+        spark.sql(s"call run_compaction(table => '$tableName', op => 
'scheduleandexecute')");
+
+        // scheduleandexecute should run compact on all pending compactions
+        metaClient.reloadActiveTimeline();
+        
assert(metaClient.getActiveTimeline.filterPendingCompactionTimeline().getInstants.size()
 == 0)
       }
     }
   }

Reply via email to