Aleksandr Efimov has posted comments on this change. ( http://gerrit.cloudera.org:8080/24592 )
Change subject: IMPALA-15189: Support HBO for SortNode cardinality ...................................................................... Patch Set 5: (2 comments) Went through PS5. The key layout (TIES, order directions, per-partition limits) and the operand qualification for sorts over joins look right, and HboKeyStringTest covers the identity distinctions that matter (ordering column, ASC/DESC, ties, partitioned top-n over a join). Two comments on the distributed local-TopN side. Note on the base: PS5 sits on 24426 PS12, and PS13/14 of 24426 reworked the JoinNode helpers this patch extends (buildHboOperandQualifierMap dropped its statsType/strategy params), so the upcoming rebase will need to reconcile buildHboInputQualifierMap with that refactor. http://gerrit.cloudera.org:8080/#/c/24592/5/fe/src/main/java/org/apache/impala/planner/DistributedPlanner.java File fe/src/main/java/org/apache/impala/planner/DistributedPlanner.java: http://gerrit.cloudera.org:8080/#/c/24592/5/fe/src/main/java/org/apache/impala/planner/DistributedPlanner.java@1294 PS5, Line 1294: lowerTopN.setTopNMergeInput(true); The recompute drops the HBO cardinality but not hasHboCard_ - nothing ever clears that flag (PlanNode.java:1100 is its only write), so when init() picked up a match, the lower TopN's explain keeps annotating the freshly recomputed estimate with '(from HBO)'. The ties childSortNode marked at line 1367 hits the same thing through the computeStats() at the end of createOrderByFragment(). Clearing the flag in setTopNMergeInput(), or resetting it at the top of tryUpdateCardinalityFromHbo(), would keep the annotation truthful. On the TODO: reusing the final TopN's HBO number here would understate the local output - sum over instances of min(limit, instance_rows) >= min(limit, total_rows) - so recomputing looks right to me. http://gerrit.cloudera.org:8080/#/c/24592/5/fe/src/main/java/org/apache/impala/planner/DistributedPlanner.java@1390 PS5, Line 1390: childSortNode.computeStats(ctx_.getRootAnalyzer()); In the non-ties path childSortNode is not marked as merge input, so this computeStats() keeps it HBO-tracked, and at toThrift it writes stats under the same key the logical sort used during single-node planning (for offset=0 the limit propagates unchanged). What gets stored is the cross-instance sum of per-instance limited outputs: with I instances all reaching the limit, num_rows = I*N recorded under the LIMIT:N key. Reads survive today only because tryUpdateCardinalityFromHbo() caps the matched value at the node's own limit, so cardinality_ lands back on N - but the stored runs are inflated, and anything consuming the matched row count without a cap (a future stats type, or surfacing the raw matched value) inherits I*N. Two follow-ups: (1) is leaving this path unmarked intentional? The commit message's 'skip local TopN since HBO tracks the summed cross-instance cardinality' rationale seems to apply here as much as to the ties/partitioned paths. (2) With offset > 0 the rewrite to (limit+offset, offset=0) means the LIMIT:N|OFFSET:M key that single-node planning looks up is never written by distributed executions, so offset queries can only match here, after the rewrite - worth a comment or a test if that is by design. -- To view, visit http://gerrit.cloudera.org:8080/24592 To unsubscribe, visit http://gerrit.cloudera.org:8080/settings Gerrit-Project: Impala-ASF Gerrit-Branch: master Gerrit-MessageType: comment Gerrit-Change-Id: Ib829887a91593bee124d56e661e714575fe3be97 Gerrit-Change-Number: 24592 Gerrit-PatchSet: 5 Gerrit-Owner: Quanlong Huang <[email protected]> Gerrit-Reviewer: Aleksandr Efimov <[email protected]> Gerrit-Reviewer: Impala Public Jenkins <[email protected]> Gerrit-Comment-Date: Mon, 17 Aug 2026 08:45:21 +0000 Gerrit-HasComments: Yes
