KiteSoar commented on code in PR #20099:
URL: https://github.com/apache/hudi/pull/20099#discussion_r4172549999
##########
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:
Done in b231fe4c5612. Archived payloads are now loaded lazily in
descending-time batches, with only the current batch retained. The
partition-row iterator stops at the row limit and continues past empty cleans.
Added zero-row/multi-row partition-metadata tests, including limit = 3
returning two rows from one clean and one from the next, with payload-load
assertions.
##########
hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/command/procedures/ShowCleansPlanProcedure.scala:
##########
@@ -184,21 +184,17 @@ class ShowCleansPlanProcedure extends BaseProcedure with
ProcedureBuilder with S
}
private def getCleanerPlans(metaClient: HoodieTableMetaClient, limit: Int,
showArchived: Boolean): Seq[Row] = {
- val activeCleanInstants =
getSortedCleanInstants(metaClient.getActiveTimeline)
- .take(limit)
-
- val cleanInstants = if (showArchived) {
- val archivedCleanInstants =
getSortedCleanInstants(metaClient.getArchivedTimeline)
- .take(limit)
- (activeCleanInstants ++ archivedCleanInstants)
- .sortWith((a, b) => a.requestedTime() > b.requestedTime())
- .take(limit)
+ val activeTimeline = metaClient.getActiveTimeline
+ val activeCleanInstants =
getSortedCleanInstants(activeTimeline).take(limit)
+ val activeRows = activeCleanInstants.map(processCleanPlan(metaClient,
activeTimeline, _))
+
+ if (showArchived) {
+ val archivedTimeline =
ShowCleansProcedure.getArchivedCleanTimeline(metaClient, loadPlans = true,
limit = limit)
+ val archivedCleanInstants = getSortedCleanInstants(archivedTimeline)
+ val archivedRows =
archivedCleanInstants.map(processCleanPlan(metaClient, archivedTimeline, _))
+ (activeRows ++ archivedRows).sortWith((a, b) => a.getString(0) >
b.getString(0)).take(limit)
Review Comment:
Done in b231fe4c5612. Active and archived instant descriptors are merged and
the final top N is selected before any plan payload is read. The archived
helper accepts those selected instants. Tests verify that active plans filling
the limit load no archived payloads and that a mixed result loads only the
selected archived plan.
--
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]