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

Aleksandr Efimov reassigned IMPALA-15326:
-----------------------------------------

    Assignee: Aleksandr Efimov

> Exchange queue estimate divides bytes by the sender count
> ---------------------------------------------------------
>
>                 Key: IMPALA-15326
>                 URL: https://issues.apache.org/jira/browse/IMPALA-15326
>             Project: IMPALA
>          Issue Type: Bug
>          Components: Frontend
>            Reporter: Aleksandr Efimov
>            Assignee: Aleksandr Efimov
>            Priority: Major
>
> {{ExchangeNode.estimateTotalQueueByteSize()}} caps the queue estimate with 
> the bytes that arrive divided by {{getNumNodes()}}:
> {code:java}
> long totalBytesToReceive = (long) Math.ceil(getAvgRowSize() * 
> getCardinality());
> // Assuming no skew in distribution during data shuffling.
> long bytesToReceivePerExchNode = isBroadcastExchange() ? totalBytesToReceive
>     : totalBytesToReceive / getNumNodes();
> {code}
> https://github.com/apache/impala/blob/7a3cb82b49d845bc5ef597e598f589593dd13549/fe/src/main/java/org/apache/impala/planner/ExchangeNode.java#L339-L355
> On an ExchangeNode {{getNumNodes()}} comes from the child - 
> {{PlanNode.computeStats()}} does {{numNodes_ = getChild(0).numNodes_}} - so 
> it counts the *senders*. The receivers are counted by {{getNumReceivers()}}, 
> which this class already uses when it computes processing cost. For an 
> UNPARTITIONED merging exchange every row lands on the single receiver while 
> the bytes are divided by the sender count.
> Numbers from card-limit-offset.test, with the IMPALA-15324 fix in place so 
> the arriving row count is right: 7200 rows of 8B are 57600B, divided by 3 
> senders gives 19200B, and with the deferred term (3 x 1024 x 12B = 36864B) 
> the node estimate is 56064B = 54.75KB. Dividing by the receiver count instead 
> gives 94464B = 92.25KB.
> Related and visible in the same expression: a limit that applies per instance 
> is counted once for the whole fragment, so a TOP-N with LIMIT=100 under a 
> 10-instance fragment reports cardinality 100 while up to 1000 rows can reach 
> the receiver's queues. Both make the cap smaller than what arrives.
> IMPALA-13610 is about the other direction in the same method - the planner 
> scaling the queue size by the number of senders where the backend has one 
> shared limit - so a fix here should probably settle both.
> Found while reviewing the fix for IMPALA-15324.



--
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