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]