Josh Rosen created SPARK-58481:
----------------------------------
Summary: FULL OUTER JOIN with a NOT IN subquery condition returns
the cross product
Key: SPARK-58481
URL: https://issues.apache.org/jira/browse/SPARK-58481
Project: Spark
Issue Type: Bug
Components: SQL
Affects Versions: 4.0.0
Reporter: Josh Rosen
The following query returns incorrect results in Spark 3.4.0 and all later
versions (I haven't tested earlier versions):
{code:sql}
CREATE TABLE t0(c0 INT) USING PARQUET;
INSERT INTO t0 VALUES (1), (2), (3);
CREATE TABLE t1(c0 INT) USING PARQUET;
INSERT INTO t1 VALUES (10), (20), (30);
CREATE TABLE t3(c0 INT) USING PARQUET;
INSERT INTO t3 VALUES (99), (CAST(NULL AS INT));
SELECT t0.c0, t1.c0 FROM t1 FULL OUTER JOIN t0 ON (5 NOT IN (SELECT t3.c0 FROM
t3));
{code}
This returns 9 rows, the full cross product of tables {{t0}} and {{{}t1{}}}.
The correct answer is 6 rows, because {{5 NOT IN (99, NULL)}} evaluates to
UNKNOWn and a join condition that is not TRUE matches no rows, so a FULL OUTER
JOIN must emit 3 + 3 = 6 null-padded rows.
Setting {{spark.sql.codegen.factoryMode=NO_CODEGEN}} causes Spark to return the
correct 6 rows.
Controls on the same tables:
||ON condition||rows||
|{{ON false}}|6|
|{{ON CAST(NULL AS BOOLEAN)}}|6|
|{{ON true}}|9|
|{{ON (5 NOT IN (SELECT c0 FROM t3))}}|9, should be 6|
h3. Suspected root cause
The cause is a nullability declaration that does not match the values the
expression can produce:
* {{InSubqueryExec.nullable}} is
[defined|https://github.com/apache/spark/blob/4dcdd4dba544210898180a410b8b30e0b85157cb/sql/core/src/main/scala/org/apache/spark/sql/execution/subquery.scala#L128]
as {{child.nullable}} .
* In our repro query, the child is the literal {{{}5{}}}, which is not
nullable. But {{x IN (subquery)}} evaluates to UNKNOWN whenever no match is
found and the subquery result contains NULL, and the declaration does not
account for that.
* {{{}InSet{}}}, which {{InSubqueryExec}} delegates its own {{eval}} and
{{doGenCode}} to,[ derives it correctly
|https://github.com/apache/spark/blame/4dcdd4dba544210898180a410b8b30e0b85157cb/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/predicates.scala#L669]as
{{child.nullable || hasNull}} , so the wrapper and its delegate disagree.
* {{Not}} inherits the incorrect declaration.
{{UnaryExpression.nullSafeCodeGen}} therefore takes its non-nullable branch,
emitting {{isNull = FalseLiteral}} and computing {{!value }}unconditionally.
The generated {{InSet}} code does the right thing: it leaves {{value}} at its
default of {{false}} and raises {{{}isNull{}}}. Nothing reads {{{}isNull{}}},
so the condition is computed as {{{}!false{}}}, which is TRUE.
* The interpreted path checks the runtime value for null regardless of the
declared nullability, which is why the two evaluators disagree.
This is long-standing rather than a recent regression. 3.4.0 reproduces
identically at 9 rows
versus 6. It is not ANSI-dependent either: 3.4.0 defaults to
{{spark.sql.ansi.enabled=false}} and
behaves the same.
h3. Potential fix
Derive nullability from the subquery output as well as the child, for example
{{{}child.nullable || plan.output.head.nullable{}}}. The delegate's
{{child.nullable || hasNull }}cannot be reused directly because {{hasNull}}
needs the executed result; the expression above is the planning-time
over-approximation. No codegen change is needed, since {{{}InSet{}}}'s
generated code already models the NULL case.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]