David Mollitor created SPARK-59356:
--------------------------------------

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


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