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]

Reply via email to