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

Reply via email to