danny0405 commented on code in PR #20099:
URL: https://github.com/apache/hudi/pull/20099#discussion_r4118335784
##########
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:
[P2] Select the final top N instants before loading and decoding their
plans. This now processes up to limit active plans plus limit archived plans
before taking the final limit. If the active timeline has limit newer clean
instants, every archived plan payload is loaded and decoded only to be
discarded; cleaner plans can contain large file lists, so the extra archive
reads and allocations can be significant. A simpler flow would merge instant
descriptors while retaining their source timeline, sort/take(limit), load
archived payloads only for selected archived instants, and then build rows.
Having the helper accept selected instants instead of its own limit would also
avoid repeated selection and sorting. The new limit => 1 result assertion would
still pass with these unnecessary reads, so a payload-load assertion would help
protect this behavior.
##########
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:
[P2] Bound archived metadata loading by the rows actually needed. For
show_cleans_metadata, Int.MaxValue causes getArchivedCleanTimeline to copy
every archived clean's metadata into the contents map before
getCleansWithPartitionMetadata applies the row limit. Thus even limit => 1
loads and retains the full archived clean metadata history, including when
newer active rows already fill the result. On long-lived tables this can cause
substantial archive I/O and driver-memory pressure, potentially OOM. Could this
load descending-time batches and stop once enough partition rows have been
produced? Merely passing limit as the instant count would not preserve behavior
because a clean can have empty partition metadata and contribute no rows. A
partitioned-table test with multi-row and zero-row cleans, plus an assertion on
payloads loaded, would cover this.
--
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]