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]

Reply via email to