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)
}
}
}