voonhous commented on code in PR #20099:
URL: https://github.com/apache/hudi/pull/20099#discussion_r4118483702
##########
hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/command/procedures/ShowCleansProcedure.scala:
##########
@@ -256,6 +266,48 @@ object ShowCleansProcedure {
val NAME = "show_cleans"
def builder: Supplier[ProcedureBuilder] = () => new
ShowCleansProcedure(false)
+
+ private[procedures] def getArchivedCleanTimeline(metaClient:
HoodieTableMetaClient,
+ loadPlans: Boolean,
+ limit: Int = Int.MaxValue):
HoodieTimeline = {
+ val cleanInstants =
metaClient.getArchivedTimeline.getCleanerTimeline.filterCompletedInstants
+ .getReverseOrderedInstants.iterator().asScala.take(limit).toList
+ val contents = new ConcurrentHashMap[String, Array[Byte]]()
+ val factory = metaClient.getTableFormat.getTimelineFactory
+ if (cleanInstants.nonEmpty) {
+ val legacy = metaClient.getTimelineLayoutVersion.getVersion <
TimelineLayoutVersion.VERSION_2
+ val actionField = if (legacy) "actionType" else "action"
+ val contentField = if (legacy) {
+ if (loadPlans) "hoodieCleanerPlan" else "hoodieCleanMetadata"
+ } else {
+ if (loadPlans) "plan" else "metadata"
+ }
Review Comment:
**minor:** not blocking, but this puts both archive formats' field names and
payload decoding into a Spark procedure, and hudi-cli has the same bug on V1:
`CLIUtils.java:57-60` loads details through `ArchivedTimelineV1.java:346`,
which stores clean payloads as JSON that `readCleanMetadata` cannot parse.
Would it make sense to host this helper in hudi-common (a static on
`HoodieArchivedTimeline` or `TimelineUtils`) so `CleansCommand` can share it,
and to reuse `ArchivedTimelineV2.ACTION_ARCHIVED_META_FIELD` /
`PLAN_ARCHIVED_META_FIELD` / `METADATA_ARCHIVED_META_FIELD` instead of the
literals?
##########
hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/command/procedures/ShowCleansProcedure.scala:
##########
@@ -165,10 +172,13 @@ class ShowCleansProcedure(includePartitionMetadata:
Boolean) extends BaseProcedu
getCleans(metaClient.getActiveTimeline, limit)
}
val finalResults = if (showArchived) {
+ val archivedCleanLimit = if (includePartitionMetadata) Int.MaxValue else
limit
Review Comment:
**minor:** not blocking. `Int.MaxValue` loads the metadata of every archived
clean regardless of `limit`. Many rows per instant argues for fewer instants,
not more; the only case that needs more than `limit` instants is a 0-row empty
clean, which requires `hoodie.write.empty.clean.interval.hours` > 0 (default
-1, `HoodieCleanConfig.java:254`). Could we pass `limit` here, or extend in
batches only while the row count is below `limit`?
##########
hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/command/procedures/ShowCleansProcedure.scala:
##########
@@ -256,6 +266,48 @@ object ShowCleansProcedure {
val NAME = "show_cleans"
def builder: Supplier[ProcedureBuilder] = () => new
ShowCleansProcedure(false)
+
+ private[procedures] def getArchivedCleanTimeline(metaClient:
HoodieTableMetaClient,
+ loadPlans: Boolean,
+ limit: Int = Int.MaxValue):
HoodieTimeline = {
+ val cleanInstants =
metaClient.getArchivedTimeline.getCleanerTimeline.filterCompletedInstants
+ .getReverseOrderedInstants.iterator().asScala.take(limit).toList
+ val contents = new ConcurrentHashMap[String, Array[Byte]]()
+ val factory = metaClient.getTableFormat.getTimelineFactory
+ if (cleanInstants.nonEmpty) {
+ val legacy = metaClient.getTimelineLayoutVersion.getVersion <
TimelineLayoutVersion.VERSION_2
+ val actionField = if (legacy) "actionType" else "action"
+ val contentField = if (legacy) {
+ if (loadPlans) "hoodieCleanerPlan" else "hoodieCleanMetadata"
+ } else {
+ if (loadPlans) "plan" else "metadata"
+ }
+ val loadMode = if (loadPlans) HoodieArchivedTimeline.LoadMode.PLAN else
HoodieArchivedTimeline.LoadMode.METADATA
+ factory.createArchivedTimelineLoader().loadInstants(metaClient,
+ new HoodieArchivedTimeline.ClosedClosedTimeRangeFilter(
+ cleanInstants.last.requestedTime(),
cleanInstants.head.requestedTime()),
+ loadMode,
+ record => HoodieTimeline.CLEAN_ACTION ==
record.get(actionField).toString,
Review Comment:
**blocker:** on V1 tables the clean-only `commitsFilter` trips the loader's
per-file early exit, so older archived cleans load no payload and the #19639
symptoms return (IOException from `show_cleans` / `show_cleans_metadata`, null
rows from `show_clean_plans`). Evidence: `ArchivedTimelineLoaderV1.java:139`
applies the filter before counting in-range instants and `:157-163` breaks at
the first file that yields 0 after one yielded > 0; every archival run writes a
new `.commits_.archive.N` (`HoodieLogFormatWriter.java:98-99`) and cleans
archive on their own threshold (`TimelineArchiverV1.java:205-213`), so a
commit-only file between two clean-bearing ones is routine, and
`show_cleans_metadata` (`Int.MaxValue`) always spans it. The new test archives
once, so it only ever sees one file. Could we pass `null` as the range filter
when `legacy` (the break is guarded by `filter != null`) and apply
`isInRange(instantTime)` in the consumer, plus a V1 case with three archival
runs and a commi
t-only middle file?
<details><summary>Why `record => true` is not enough</summary>
The exit assumes archive files are time-ordered, but cleans lag commits: a
newer file can hold only commits above the range while an older file still
holds in-range cleans, so counting every action still breaks early. Dropping
the filter on V1 costs one extra full pass, which `getArchivedTimeline` already
does today for V1.
</details>
##########
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,
+ // otherwise the archived branch below is never exercised and the
test passes vacuously.
+ val metaClient = createMetaClient(spark, tablePath)
+
assertResult(tableVersion)(metaClient.getTableConfig.getTableVersion.versionCode())
+ val archivedCleans =
metaClient.getArchivedTimeline.getCleanerTimeline.filterCompletedInstants
+ .getInstants.asScala.map(_.requestedTime).toSeq
+ val activeCleans =
metaClient.getActiveTimeline.getCleanerTimeline.filterCompletedInstants
+ .getInstants.asScala.map(_.requestedTime).toSeq
+ assert(archivedCleans.nonEmpty, s"expected archived cleans, got
$archivedCleans")
+ assert(activeCleans.nonEmpty, s"expected active cleans, got
$activeCleans")
+ assert(archivedCleans.max.compareTo(activeCleans.min) < 0,
+ "archived cleans must be older than active cleans")
+
+ // The active-only view must exclude the archived clean for all
three procedures.
+ assertResult(activeCleans.size)(rowCount("show_cleans", showArchived
= false))
+ assertResult(activeCleans.size)(rowCount("show_cleans_metadata",
showArchived = false))
+ assertResult(activeCleans.size)(rowCount("show_clean_plans",
showArchived = false))
+
+ procedures.foreach { procedure =>
+ val allRows = spark.sql(s"call $procedure(table => '$tableName',
showArchived => true)").collect().toSeq
+ assertResult(beforeArchive(procedure))(allRows)
+ val limitedRows = spark.sql(s"call $procedure(table =>
'$tableName', showArchived => true, limit => 1)").collect().toSeq
+ assertResult(allRows.take(1))(limitedRows)
Review Comment:
**minor:** not blocking. Both procedures sort the merged rows descending
before `take(limit)`, so `limit => 1` always returns the newest active clean
and the archived `take(limit)` pushdown is never exercised; taking the wrong
end would still pass. Could we use `limit => activeCleans.size + 1` and assert
the last row is `archivedCleans.max`?
##########
hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/command/procedures/ShowCleansProcedure.scala:
##########
@@ -256,6 +266,48 @@ object ShowCleansProcedure {
val NAME = "show_cleans"
def builder: Supplier[ProcedureBuilder] = () => new
ShowCleansProcedure(false)
+
+ private[procedures] def getArchivedCleanTimeline(metaClient:
HoodieTableMetaClient,
+ loadPlans: Boolean,
+ limit: Int = Int.MaxValue):
HoodieTimeline = {
Review Comment:
**nit:** feel free to ignore. Both callers pass `limit`, so this default is
the one line codecov reports as uncovered.
```suggestion
limit: Int):
HoodieTimeline = {
```
##########
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 =>
Review Comment:
**nit:** feel free to ignore. `HoodieTableVersion.current()` is TEN, so this
drops the default version the replaced test ran on (9 and 10 share layout V2,
so the branch coverage is the same). Needs `import
org.apache.hudi.common.table.HoodieTableVersion`.
```suggestion
Seq(6, HoodieTableVersion.current().versionCode()).foreach { tableVersion
=>
```
##########
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()
Review Comment:
**nit:** feel free to ignore. The table version is set three ways
(`withSQLConf`, the tblproperty, and `options` here). `run_clean` already picks
the tblproperty up through `HoodieCLIUtils.getWriteParameters`
(`HoodieCLIUtils.scala:72-75`); only `archive_commits` builds its config from
`options` alone. Could we drop the `withSQLConf` entry and this `options`
argument, keeping the tblproperty plus the `archive_commits` `options`?
##########
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:
**nit:** feel free to ignore. The test now creates six cleans, so "the two
cleans" is stale.
```suggestion
// Precondition: archival must have split the cleans across the
active and archived timelines,
```
--
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]