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

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

> Estimate selectivity of StartsWith/EndsWith/Contains in FilterEstimation
> ------------------------------------------------------------------------
>
>                 Key: SPARK-59356
>                 URL: https://issues.apache.org/jira/browse/SPARK-59356
>             Project: Spark
>          Issue Type: Improvement
>          Components: SQL
>    Affects Versions: 4.1.0
>            Reporter: David Mollitor
>            Priority: Minor
>              Labels: pull-request-available
>
> h3. What
> The cost-based optimizer's {{FilterEstimation}} lists the string predicates
> {{StartsWith}} / {{EndsWith}} / {{Contains}} (and {{{}Like{}}}) as "not 
> supported yet", so they fall
> through to a fixed default selectivity that ignores the column statistics 
> Spark already
> collects. This estimates {{StartsWith}} / {{EndsWith}} / {{Contains}} from 
> two existing stats –
> {{nullCount}} and {{maxLen}} – so filtered cardinalities, and the join plans 
> that depend on
> them (join order, build side, broadcast eligibility), are more accurate.
> These are exactly the predicates the LIKE simplifications produce (e.g. {{col 
> LIKE 'abc%'}} -> {{{}StartsWith(col, 'abc'){}}}), so today those rewrites 
> help pushdown/pruning but not the CBO's cardinality estimates.
> h3. The estimate
> For a null-intolerant string predicate on a column with statistics:
> {code:java}
> selectivity(StartsWith(col, prefix)) =
>     0                        if prefix.codePointCount > maxLen   // no value 
> is long enough
>     0                        if nullPercent == 1                 // column is 
> all null
>     (1 - nullPercent) * f    otherwise
> {code}
>  * {{{}nullCount{}}}: {{{}StartsWith{}}}/{{{}EndsWith{}}}/{{{}Contains{}}} 
> are null-intolerant, so a null input yields null and never matches in a 
> predicate. Only the non-null rows can match, so selectivity is bounded by {{1 
> - nullPercent}} (and is 0 for an all-null column). {{FilterEstimation}} 
> already computes {{nullPercent}} for {{{}IsNull{}}}/{{{}IsNotNull{}}}; this 
> reuses it.
>  * {{{}maxLen{}}}: a value must have at least {{prefix.codePointCount}} 
> characters to start with (or end with / contain) {{{}prefix{}}}. Spark stores 
> {{maxLen}} as the maximum *character* length ({{{}Max(Length(col)){}}}), so 
> when the operand has more code points than {{{}maxLen{}}}, no row can match 
> and selectivity is 0.
> h3. Correctness
>  * Both bounds are sound and need no collation gate: {{maxLen}} is code-point 
> based, and a match consumes at least {{prefix.codePointCount}} code points 
> under any collation (including UTF8_LCASE, whose case folding is 1:1 at the 
> code-point level); {{nullCount}} reflects rows that can never match a 
> null-intolerant predicate.
>  * This affects estimation only. Statistics can be stale, so a computed 
> selectivity of 0 means "treat as empty for planning" – it must not drop the 
> predicate or replace the input with an empty relation. Selectivity is clamped 
> into {{{}FilterEstimation{}}}'s usual range.
> h3. Scope
>  * {{{}StartsWith{}}}, {{{}EndsWith{}}}, {{{}Contains{}}}. General {{Like}} 
> selectivity is out of scope (most useful {{LIKE}} patterns are already 
> rewritten to these before estimation), and {{EqualTo}} is already estimated 
> via NDV (though the {{maxLen}} zero-case could extend to it later).
>  * Applies only when the CBO is enabled and column statistics are present.
> h3. Does this PR introduce any user-facing change?
> No. Query results are unchanged; this only refines cost-based cardinality 
> estimates, and only under `spark.sql.cbo.enabled` with column stats.



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