David Mollitor created SPARK-59348:
--------------------------------------
Summary: Rewrite match-all LIKE '%' to IsNotNull in predicates
Key: SPARK-59348
URL: https://issues.apache.org/jira/browse/SPARK-59348
Project: Spark
Issue Type: Improvement
Components: SQL
Affects Versions: 4.1.0
Reporter: David Mollitor
h3. What
{{col LIKE '%'}} (a pattern of only unescaped {{{}%{}}}) matches every non-null
value, so a filter{{{}WHERE col LIKE '%'{}}} keeps exactly the non-null rows.
Today Spark evaluates it as a per-row regex ({{{}.*{}}}) and pushes nothing to
the data source. This rewrites such patterns, in null-rejecting predicate
positions, to {{{}IsNotNull(col){}}}.
{code:java}
Aggregate [count(1)] Aggregate [count(1)]
+- Filter (col LIKE '%') ==> +- Filter (isnotnull(col))
+- Relation t +- Relation t
{code}
h3. Why are the changes needed?
* {{IsNotNull}} is far cheaper than compiling and running a regex per row.
* {{IsNotNull}} pushes down: Parquet/ORC skip row groups via null-count
statistics; the raw {{LIKE}} regex does not.
* If {{col}} is non-nullable, {{IsNotNull(col)}} then folds to {{true}} and
downstream
{{{}PruneFilters{}}}/{{{}NullPropagation{}}} drop the filter entirely.
The pattern shows up in practice from dynamically-generated SQL (an empty
search box becoming {{{}LIKE '%'{}}}).
h3. Correctness
{{col LIKE '%'}} returns {{true}} for a non-null value and {{null}} for
{{{}null{}}}, whereas
{{IsNotNull(col)}} returns {{false}} for {{{}null{}}}. They are equivalent
o{_}nly in null-rejecting{_}
{_}predicate positions{_}, where {{null}} is treated as {{{}false{}}}. The
rewrite is therefore
implemented in {{{}ReplaceNullWithFalseInPredicate{}}}, which already applies
this kind of null-as-false simplification (see its existing {{Not(In(...))}}
handling) and which:
* targets exactly the null-rejecting plan nodes ({{{}Filter{}}}, {{Join}}
condition, {{{}MergeIntoTable{}}}, {{{}DeleteFromTable{}}},
{{{}UpdateTable{}}}, {{{}ReplaceData{}}}, {{{}WriteDelta{}}}); and recurses
only through null-preserving operators ({{{}And{}}}, {{{}Or{}}}, {{{}If{}}},
{{{}CaseWhen{}}}) and *not* into {{Not}} – so {{NOT (col LIKE '%')}} and {{col
LIKE '%'}} in a projection are left untouched, where {{null}} vs {{false}}
would be observable.
The rewrite fires only for a literal pattern of one or more {{%}} with the
escape character not being {{%}} (so the {{%}} are wildcards, not escaped
literals); {{{}LIKE ''{}}}, prefixes, and {{_}}
patterns are excluded. No collation gate is needed – {{'%'}} matches every
non-null value under any collation.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]