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

ASF GitHub Bot updated SPARK-59541:
-----------------------------------
    Labels: pull-request-available  (was: )

> Coalesce OR-connected equality and IN predicates on the same column into a 
> single IN
> ------------------------------------------------------------------------------------
>
>                 Key: SPARK-59541
>                 URL: https://issues.apache.org/jira/browse/SPARK-59541
>             Project: Spark
>          Issue Type: Improvement
>          Components: SQL
>    Affects Versions: 4.1.0
>            Reporter: David Mollitor
>            Priority: Minor
>              Labels: pull-request-available
>
> h2. Background
> Spark leaves OR-connected membership tests on the same column un-merged:
>  * {{x = 1 OR x = 2 OR x = 3}} stays a nested {{Or(EqualTo, ...)}} - it is 
> never folded to {{{}IN{}}}.
>  * {{x IN (1, 2) OR x IN (3, 4)}} stays two separate {{In}} nodes.
>  * {{x = 1 OR x IN (2, 3)}} stays {{{}Or(EqualTo, In){}}}.
> {{OptimizeIn}} only processes a single {{In}} node at a time, so none of 
> these are combined today.
> h2. Proposal
> A new logical optimizer rule {{CombineDisjunctiveInPredicates}} that 
> coalesces OR-connected {{EqualTo}} and {{In}} disjuncts on the same 
> deterministic, non-foldable subject into one {{{}In{}}}, letting the adjacent 
> {{OptimizeIn}} finish the job (dedup, single-element -> {{{}EqualTo{}}}, 
> large list -> {{{}InSet{}}}). {{EqualTo(x, c)}} is treated as a one-element 
> membership {{{}In(x, [c]){}}}. Examples:
>  * {{x = 1 OR x = 2 OR x = 3}} -> {{x IN (1, 2, 3)}}
>  * {{x IN (1, 2) OR x IN (3, 4)}} -> {{x IN (1, 2, 3, 4)}}
>  * {{x = 1 OR x IN (2, 3)}} -> {{x IN (1, 2, 3)}}
> The rule runs immediately before {{OptimizeIn}} and is gated by an internal 
> flag {{spark.sql.optimizer.combineDisjunctiveInPredicates.enabled}} (default 
> true).
> h2. Why are the changes needed
> A single {{{}In{}}}/{{{}InSet{}}} is a strictly more optimized representation 
> than an OR chain of the same predicates:
>  * Single subject evaluation: {{In.eval}} evaluates the subject expression 
> once; an {{Or}} of N {{{}EqualTo{}}}/{{{}In{}}} nodes re-evaluates it N times.
>  * Set dispatch: once merged and above the conversion threshold, 
> {{OptimizeIn}} turns the {{In}} into {{{}InSet{}}}, which uses a {{switch}} 
> (compact ints) - dispatch that OR-connected nodes never get.
>  * Avoids the whole-stage-codegen method-size cliff: a very large OR chain 
> overflows the generated method-size limit and falls back to interpreted 
> execution; the merged {{InSet}} stays compact.
> Microbenchmark over a {{Long}} column, 10M rows, non-matching values (worst 
> case), rule off (raw OR) vs rule on:
> ||N||raw OR (ms)||merged IN/InSet (ms)||
> |20|39|40|
> |100|148|54|
> |500|65142 (interpreted fallback)|82|
> Parity at small N and a large win as N grows.
> h2. Correctness
>  * Determinism guard (the sole correctness gate): the subject must be 
> deterministic - {{In}} evaluates it once, the OR form once per disjunct - so 
> {{rand() = 1 OR rand() = 2}} is left unchanged.
>  * Three-valued logic: {{x IN A OR x IN B}} is equivalent to {{x IN (A union 
> B)}} under nulls (both NULL when the subject is NULL; both FALSE when 
> non-null and unmatched with no null member; both NULL when unmatched with a 
> null member), matching {{{}In.eval{}}}.
>  * Exclusions: {{EqualNullSafe}} (<=>) is never matched; structs are excluded 
> on subject and members; foldable-only "subjects" are not grouped (avoids the 
> degenerate {{1 IN (x, y)}} from {{{}x = 1 OR y = 1{}}}).
>  * Fixed-point safety: fires only on groups with >= 2 candidate disjuncts, 
> collapsing to one {{In}} that cannot re-match; returns the {{Or}} unchanged 
> otherwise.
> h2. Caveat: Parquet row-group pruning precision
> Merging a large (> {{{}spark.sql.optimizer.inSetConversionThreshold{}}}) 
> OR-of-equalities into {{IN}} > {{InSet}} can make Parquet row-group pruning 
> coarser: a raw OR-of{{{}EqualTo{}}} pushes as {{FilterApi.or(eq, ...)}} 
> (per-value), whereas the merged form pushes as {{FilterApi.in}} (a min/max 
> envelope over the set). The internal flag allows disabling the rule where 
> this matters.



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