[ 
https://issues.apache.org/jira/browse/IMPALA-10430?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=17264707#comment-17264707
 ] 

Csaba Ringhofer commented on IMPALA-10430:
------------------------------------------

Thanks [~amansinha] for bringing my attention to IMPALA-10317 ! It already 
implements the  'joined' counter mentioned in this ticket.

I think that tracking original / duplicate ratio would be a further 
improvement, because it is easier to find a default that is unlikely to cancel 
legitimate queries. E.g. is 10 billion rows too much? Generally not, especially 
in a bigger cluster,  but if 99% of those are duplicates, than it is very 
likely that it is not intended, or could be optimized to avoid such huge number 
of duplicates.

> Allow limiting duplicate rows from joins
> ----------------------------------------
>
>                 Key: IMPALA-10430
>                 URL: https://issues.apache.org/jira/browse/IMPALA-10430
>             Project: IMPALA
>          Issue Type: Improvement
>          Components: Backend
>            Reporter: Csaba Ringhofer
>            Priority: Minor
>
> Small mistakes in join conditions can lead to huge number of duplicates and 
> turn a simple query into a massive resource hog until it is cancelled (or 
> runs out of memory). An even worse case is if this happens in an INSERT and 
> data is being dumped to HDFS until some quote is hit.
> An idea for a safety net is to track the  number of "original" rows (that 
> come from scanning a table) and "joined" rows (that come from joins), and 
> abort the query if joined/original is above a certain ratio. 
> E.g. the following query options could control this:
> max_joined_original_ratio
> min_rows_to_enforce_joined_original_ratio (to allow intentional cross joins 
> for smaller datesets)
> So the query could be cancelled if 
> max(max_joined_original_ratio * original, 
> min_rows_to_enforce_joined_original_ratio ) < joined
> We could find default values that would rarely limit sane queries, e.g.
> max_joined_original_ratio = 100
> min_rows_to_enforce_joined_original_ratio=10,000,000,000
> If someone knows that the query should have no duplicates at all then this 
> could be enforced with
> max_joined_original_ratio = 1,
> min_rows_to_enforce_joined_original_ratio=0
> Tracking 'original' and 'joined' could be done several ways, but it is 
> important that it should be done "globally", not per join node , because 
> consecutive joins'  "duplicate factors" are multiplied, e.g. five joins that 
> all  have a ratio of 10x would result in 10000x global ratio. A possibility 
> is to use profile counters in scan and join nodes and sum them in the 
> coordinator.
> Apart from safety from runaway queries this could also help the planner by 
> having enforced limits, especially in the max_joined_original_ratio = 1 case.



--
This message was sent by Atlassian Jira
(v8.3.4#803005)

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

Reply via email to