KiteSoar commented on code in PR #20099:
URL: https://github.com/apache/hudi/pull/20099#discussion_r4172586737


##########
hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/spark/sql/hudi/procedure/TestShowCleansProcedures.scala:
##########
@@ -481,126 +479,123 @@ class TestShowCleansProcedures extends 
HoodieSparkProcedureTestBase {
     }
   }
 
-  test("Test show_clean_plans with an archived clean instant") {
-    withSQLConf("hoodie.clean.automatic" -> "false", 
"hoodie.archive.automatic" -> "false") {
-      withTempDir { tmp =>
-        val tableName = generateTableName
-        val tablePath = tmp.getCanonicalPath
-        spark.sql(
-          s"""
-             |create table $tableName (
-             | id int,
-             | name string,
-             | price double,
-             | ts long
-             | ) using hudi
-             | location '$tablePath'
-             | tblproperties (
-             |   primaryKey = 'id',
-             |   type = 'cow',
-             |   preCombineField = 'ts',
-             |   hoodie.metadata.enable = 'false'
-             | )
-             |""".stripMargin)
+  Seq(6, 9).foreach { tableVersion =>
+    test(s"Test show_clean_plans with an archived clean instant - table 
version $tableVersion") {
+      withSQLConf("hoodie.clean.automatic" -> "false", 
"hoodie.archive.automatic" -> "false",
+        "hoodie.write.table.version" -> tableVersion.toString) {
+        withTempDir { tmp =>
+          val tableName = generateTableName
+          val tablePath = tmp.getCanonicalPath
+          spark.sql(
+            s"""
+               |create table $tableName (
+               | id int,
+               | name string,
+               | price double,
+               | ts long
+               | ) using hudi
+               | location '$tablePath'
+               | tblproperties (
+               |   primaryKey = 'id',
+               |   type = 'cow',
+               |   preCombineField = 'ts',
+               |   hoodie.metadata.enable = 'false',
+               |   'hoodie.write.table.version' = '$tableVersion'
+               | )
+               |""".stripMargin)
 
-        def rowCount(procedure: String, showArchived: Boolean): Int =
-          spark.sql(s"call $procedure(table => '$tableName', showArchived => 
$showArchived)").collect().length
-
-        // Six write commits with a clean in the middle: the first clean sits 
before the last
-        // commit that archival will move, so it gets archived along with 
those commits, while
-        // the second clean stays on the active timeline.
-        spark.sql(s"insert into $tableName values(1, 'a1', 10, 1000)")
-        spark.sql(s"update $tableName set price = 11 where id = 1")
-        spark.sql(s"update $tableName set price = 12 where id = 1")
-        spark.sql(s"call run_clean(table => '$tableName', retain_commits => 
1)").collect()
-        spark.sql(s"update $tableName set price = 13 where id = 1")
-        spark.sql(s"update $tableName set price = 14 where id = 1")
-        spark.sql(s"update $tableName set price = 15 where id = 1")
-        spark.sql(s"call run_clean(table => '$tableName', retain_commits => 
1)").collect()
-
-        // Nothing has been archived yet, so showArchived => true here only 
exercises the merge of
-        // the active timeline with an empty archived one. All three 
procedures succeed and see
-        // both cleans, which is what makes the failures asserted after 
archival attributable to
-        // the archived instants themselves rather than to the merged read 
path. Two rows for
-        // show_cleans_metadata as well, since that procedure emits one row 
per partition per
-        // clean and this table is not partitioned.
-        assertResult(2)(rowCount("show_cleans", showArchived = true))
-        assertResult(2)(rowCount("show_cleans_metadata", showArchived = true))
-        assertResult(2)(rowCount("show_clean_plans", showArchived = true))
-
-        spark.sql(s"call archive_commits(table => '$tableName', min_commits => 
2, max_commits => 3," +
-          " retain_commits => 1, enable_metadata => false)").collect()
-
-        // Precondition: archival must have split the two cleans across the 
two timelines,
-        // otherwise the archived branch below is never exercised and the test 
passes vacuously.
-        val metaClient = createMetaClient(spark, tablePath)
-        val archivedCleans = metaClient.getArchivedTimeline.getCleanerTimeline
-          .getInstants.asScala.map(_.requestedTime).toSeq
-        val activeCleans = metaClient.getActiveTimeline.getCleanerTimeline
-          .getInstants.asScala.map(_.requestedTime).toSeq
-        assert(archivedCleans.length == 1, s"expected exactly 1 archived 
clean, got $archivedCleans")
-        assert(activeCleans.length == 1, s"expected exactly 1 active clean, 
got $activeCleans")
-        assert(archivedCleans.head.compareTo(activeCleans.head) < 0,
-          "the archived clean must be the older of the two")
-
-        // Sibling-procedure controls, pinning what the other two procedures 
do with the very same
-        // archived clean. On the active timeline all three agree: one row for 
the one active
-        // clean. With showArchived => true they diverge, and neither sibling 
is correct today.
-        assertResult(1)(rowCount("show_cleans", showArchived = false))
-        assertResult(1)(rowCount("show_cleans_metadata", showArchived = false))
-        assertResult(1)(rowCount("show_clean_plans", showArchived = false))
-
-        // The other half of #19639, which covers all three clean procedures. 
show_cleans and its
-        // show_cleans_metadata variant do route to getArchivedTimeline, but 
the archived instants
-        // carry no content there, so readCleanMetadata cannot deserialize 
them and the call fails
-        // outright rather than degrading to a partial row. Same 
missing-archived-content cause as
-        // the all-null plan rows asserted below, just a harsher symptom. 
Pinned as observed.
-        Seq("show_cleans", "show_cleans_metadata").foreach { procedure =>
-          val e = intercept[IOException](rowCount(procedure, showArchived = 
true))
-          assert(e.getMessage.contains(archivedCleans.head),
-            s"$procedure over the archived timeline should fail on the 
archived clean" +
-              s" ${archivedCleans.head}, but failed with: ${e.getMessage}")
-        }
+          def rowCount(procedure: String, showArchived: Boolean): Int =
+            spark.sql(s"call $procedure(table => '$tableName', showArchived => 
$showArchived)").collect().length
 
-        val plans = spark.sql(s"call show_clean_plans(table => '$tableName', 
showArchived => true)").collect()
-        assertResult(2)(plans.length)
-        val planTimes = plans.map(_.getString(0)).mkString(", ")
-        val activePlan = plans.find(_.getString(0) == activeCleans.head)
-          .getOrElse(fail(s"no plan row for the active clean 
${activeCleans.head}, got plan_times: $planTimes"))
-        val archivedPlan = plans.find(_.getString(0) == archivedCleans.head)
-          .getOrElse(fail(s"no plan row for the archived clean 
${archivedCleans.head}, got plan_times: $planTimes"))
-
-        // Fields that come from the cleaner plan itself, resolved by name off 
the row schema so
-        // that a change in output-schema ordering cannot silently re-point 
these assertions.
-        // extra_metadata is left out: it is null for both rows, so it does 
not discriminate.
-        val planFields = Seq(
-          "earliest_instant_to_retain",
-          "last_completed_commit_timestamp",
-          "policy",
-          "version",
-          "total_partitions_to_clean",
-          "total_partitions_to_delete")
-
-        // The active clean plan is read correctly.
-        assertResult("COMPLETED")(activePlan.getString(1))
-        assertResult("clean")(activePlan.getString(2))
-        planFields.foreach { name =>
-          assert(!activePlan.isNullAt(activePlan.fieldIndex(name)), s"active 
clean plan should have a non-null $name")
-        }
+          // V1 archives cleans using a separate count threshold, so generate 
enough cleans
+          // to force archival with both timeline formats.
+          spark.sql(s"insert into $tableName values(1, 'a1', 10, 1000)")
+          spark.sql(s"update $tableName set price = 11 where id = 1")
+          spark.sql(s"update $tableName set price = 12 where id = 1")
+          (0 until 6).foreach { i =>
+            if (i > 0) {
+              spark.sql(s"update $tableName set price = ${12 + i} where id = 
1")
+            }
+            spark.sql(s"call run_clean(table => '$tableName', retain_commits 
=> 1, " +
+              s"options => 
'hoodie.write.table.version=$tableVersion')").collect()
+          }
 
-        // Known limitation, see #19639: getCleanerPlans collects the archived 
clean instants but
-        // then hands every instant to processCleanPlan against the ACTIVE 
timeline, so the
-        // archived instant's .clean.requested file is not found, the read 
falls back to
-        // createErrorRow and every plan field comes back null. Only 
plan_time/state/action
-        // survive, because those are taken from the instant and not from the 
plan. A fix has to
-        // read the archived instant's own content and would flip the null 
assertions below to
-        // non-null; merely routing to getArchivedTimeline the way the sibling 
ShowCleansProcedure
-        // does is not enough, since that path throws today (pinned above).
-        assertResult("COMPLETED")(archivedPlan.getString(1))
-        assertResult("clean")(archivedPlan.getString(2))
-        planFields.foreach { name =>
-          assert(archivedPlan.isNullAt(archivedPlan.fieldIndex(name)),
-            s"archived clean plan is expected to return a null $name today 
(#19639)")
+          // This table is not partitioned, so all three procedures return one 
row per clean.
+          assertResult(6)(rowCount("show_cleans", showArchived = true))
+          assertResult(6)(rowCount("show_cleans_metadata", showArchived = 
true))
+          assertResult(6)(rowCount("show_clean_plans", showArchived = true))
+
+          val procedures = Seq("show_clean_plans", "show_cleans", 
"show_cleans_metadata")
+          val beforeArchive = procedures.map { procedure =>
+            procedure -> spark.sql(s"call $procedure(table => '$tableName', 
showArchived => true)").collect().toSeq
+          }.toMap
+
+          spark.sql(s"call archive_commits(table => '$tableName', min_commits 
=> 2, max_commits => 3," +
+            s" retain_commits => 1, enable_metadata => false, options => 
'hoodie.write.table.version=$tableVersion')").collect()
+
+          // Precondition: archival must have split the two cleans across the 
two timelines,

Review Comment:
   Done.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to