Josh Rosen created SPARK-58383:
----------------------------------

             Summary: "ReplaceExceptWithFilter binds the RHS filter by column 
name, returning wrong EXCEPT results when an LHS alias reuses a base column name
                 Key: SPARK-58383
                 URL: https://issues.apache.org/jira/browse/SPARK-58383
             Project: Spark
          Issue Type: Improvement
          Components: SQL
    Affects Versions: 4.0.0
            Reporter: Josh Rosen


This is a report of a longstanding (but likely difficult to hit) correctness 
bug present in Spark 2.3.0+ and all later versions.

The 
[{{ReplaceExceptWithFilter}}|https://github.com/apache/spark/blob/e8d5e4a962b434821ede3cf9cc3a046a8eead8a8/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/ReplaceExceptWithFilter.scala#L29-L44]
 optimizer rule conditionally rewrites *{{left EXCEPT right}}* to  
*{{Distinct(Filter(NOT coalesce(cond, false), left))}}* in cases where the RHS 
is the same relation and projection as the LHS plus an additional filter: the 
RHS filter is negated and applied to the LHS, then a Distinct is applied to 
produce set-valued output.

The bug: the filter condition references base columns below the RHS projection, 
but is remapped against the LHS output by name, so when an LHS alias reuses a 
base column's name, the injected filter predicate binds to the wrong column.

For example, consider the following repro (which requires a non-VALUES input 
source to prevent LocalRelation optimizations):
{code:java}
CREATE TABLE repro_t(id INT, val INT) USING parquet;
INSERT INTO repro_t VALUES (1, 10), (2, 20);

SELECT val AS id FROM repro_t; 
-- outputs: [10], [20]

SELECT val AS v FROM repro_t WHERE id = 1;
-- outputs: [10]

-- If we combine these queries with except:

(SELECT val AS id FROM repro_t)
EXCEPT 
(SELECT val AS v FROM repro_t WHERE id = 1); 
-- outputs: [10],[20] WRONG{code}
Note the RHS block contains no alias named {{{}id{}}}, so its {{WHERE id = 1}} 
unambiguously references the base column under any resolution semantics; the 
name collision is only with the LHS output's alias, a namespace the RHS 
filter's scoping never sees.

The filter is ported to the LHS column _by name_ and gets negated and applied 
to {{val}} (which the LHS aliases as {{{}id{}}}):
{code:java}
== Analyzed Logical Plan ==
id: int
Except false
:- Project [val#12251 AS id#12239]
:  +- SubqueryAlias workspace.default.repro_t
:     +- Relation workspace.default.repro_t[id#12250,val#12251] parquet
+- Project [val#12253 AS v#12240]
   +- Filter (id#12252 = 1)
      +- SubqueryAlias workspace.default.repro_t
         +- Relation workspace.default.repro_t[id#12252,val#12253] parquet

== Optimized Logical Plan ==
Aggregate [id#12239], [id#12239]
+- Project [val#12251 AS id#12239]
   +- Filter NOT coalesce((val#12251 = 1), false)
      +- Relation workspace.default.repro_t[id#12250,val#12251] parquet {code}
{*}Workaround{*}: {{SET spark.sql.optimizer.replaceExceptWithFilter=false}}

{*}Related tickets{*}: SPARK-23274 added a name presence guard 

*Potential fix:*  I think we need to translate the filter into the LHS output 
by exprId/expression equivalence through the RHS project list, and fall back to 
the anti-join when a reference cannot be traced.



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