[ 
https://issues.apache.org/jira/browse/IMPALA-15324?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Aleksandr Efimov reassigned IMPALA-15324:
-----------------------------------------

    Assignee: Aleksandr Efimov

> Exchange memory estimates ignore the rows skipped by OFFSET
> -----------------------------------------------------------
>
>                 Key: IMPALA-15324
>                 URL: https://issues.apache.org/jira/browse/IMPALA-15324
>             Project: IMPALA
>          Issue Type: Bug
>          Components: Frontend
>            Reporter: Aleksandr Efimov
>            Assignee: Aleksandr Efimov
>            Priority: Major
>
> Both parts of {{ExchangeNode}}'s memory estimate are sized from 
> {{getCardinality()}}, which {{computeStats()}} has already reduced by the 
> offset:
> {code:java}
> // estimateDeferredRPCQueueSize()
> if (getCardinality() > 0) rowBatchSize = Math.min(rowBatchSize, 
> getCardinality());
> // estimateTotalQueueByteSize()
> long totalBytesToReceive = (long) Math.ceil(getAvgRowSize() * 
> getCardinality());
> {code}
> https://github.com/apache/impala/blob/c5b6c6e04c2ec70dfe8655e88dfa89bcce212565/fe/src/main/java/org/apache/impala/planner/ExchangeNode.java#L324-L333
> https://github.com/apache/impala/blob/c5b6c6e04c2ec70dfe8655e88dfa89bcce212565/fe/src/main/java/org/apache/impala/planner/ExchangeNode.java#L339-L355
> The queues see the rows before the offset is applied. The sender-side sort 
> applies no offset - it returns {{limit + offset}} rows when there is a limit 
> ([DistributedPlanner.java#L1414-L1421|https://github.com/apache/impala/blob/c5b6c6e04c2ec70dfe8655e88dfa89bcce212565/fe/src/main/java/org/apache/impala/planner/DistributedPlanner.java#L1414-L1421])
>  - and the receiver skips the first offset rows itself while reading 
> ([exchange-node.cc#L253|https://github.com/apache/impala/blob/c5b6c6e04c2ec70dfe8655e88dfa89bcce212565/be/src/exec/exchange-node.cc#L253]).
>  So both estimates are sized from a row count that never arrives at the node. 
> {{SortNode}} does add it back, through 
> [SortInfo.estimateTopNMaterializedSize(cardinality, 
> offset)|https://github.com/apache/impala/blob/c5b6c6e04c2ec70dfe8655e88dfa89bcce212565/fe/src/main/java/org/apache/impala/analysis/SortInfo.java#L304-L307].
> The 0 cardinality that the current ordering produces hides this: a 0 makes 
> {{estimateDeferredRPCQueueSize()}} skip the cap and keep a whole row batch. 
> IMPALA-15278 fixes the ordering and the estimate follows the cardinality down 
> - in {{tpcds_cpu_cost/ddl.test}} a MERGING-EXCHANGE with {{offset: 5}}, 
> {{limit: 1}} over {{01:TOP-N [LIMIT=6]}} on 4 instances goes from 
> mem-estimate=372.00KB to 16.00KB, sized for one row while up to six rows per 
> sender come through. It lands on MIN_ESTIMATE_BYTES there and the real 
> traffic is well below that, so that plan is fine; the gap grows with the 
> offset.
> Fix: use the input cardinality (cardinality + offset) in both methods - 
> Quanlong's suggestion on the review.
> Found while reviewing IMPALA-15278 (https://gerrit.cloudera.org/c/24718/).



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to