cloud-fan commented on code in PR #58077:
URL: https://github.com/apache/spark/pull/58077#discussion_r3926343023


##########
sql/core/src/main/scala/org/apache/spark/sql/execution/subquery.scala:
##########
@@ -165,14 +197,193 @@ case class InSubqueryExec(
     }
   }
 
+  // Invariant schema/ordering data for the multi-column evaluator, computed 
once after the result
+  // is available. @transient so that serialization (result=null) does not 
trigger evaluation.
+  @transient private lazy val multiColFieldTypes: Array[DataType] =
+    plan.output.map(_.dataType).toArray
+  @transient private lazy val multiColFieldOrderings: Array[Ordering[Any]] =
+    multiColFieldTypes.map(TypeUtils.getInterpretedOrdering)
+  // Struct-level ordering used to index fully non-null result rows in a 
TreeSet.
+  @transient private lazy val multiColRowOrdering: Ordering[InternalRow] =
+    
TypeUtils.getInterpretedOrdering(child.dataType).asInstanceOf[Ordering[InternalRow]]
+
+  // Split collected rows using the struct-level Catalyst ordering so that 
duplicates are
+  // deduplicated. Fully non-null rows go into multiColNonNullSet (TreeSet) 
for O(log n)
+  // membership tests. Null-containing rows are deduplicated via a temporary 
TreeSet and then
+  // stored as multiColNullRows (Array) for linear scanning; each distinct 
null-containing row
+  // is thus scanned at most once per outer row regardless of RHS duplicate 
multiplicity.
+  // See SPARK-58481.
+  @transient private lazy val (multiColNonNullSet, multiColNullRows) = {
+    val withNull = TreeSet.newBuilder[InternalRow](multiColRowOrdering)
+    val nonNull = TreeSet.newBuilder[InternalRow](multiColRowOrdering)
+    result.foreach { r =>
+      val row = r.asInstanceOf[InternalRow]
+      if (row.anyNull) withNull += row else nonNull += row
+    }
+    (nonNull.result(), withNull.result().toArray)
+  }
+
+  // Three-valued IN semantics for multi-column subqueries.
+  // Result rows are InternalRow objects; InSet's TreeSet uses Catalyst 
ordering, but membership
+  // cannot distinguish a definitively-false candidate from an indeterminate 
one.
+  //
+  // When the LHS struct has no null fields:
+  //   Fast path: O(log n) TreeSet lookup against fully non-null result rows 
for TRUE.
+  //   Slow path: linear scan over null-containing result rows only for 
potential UNKNOWN.
+  //
+  // When the LHS struct has at least one null field, the fast path cannot be 
used (a null LHS
+  // field produces UNKNOWN against any non-null RHS row whose non-null fields 
all match). Both
+  // sets of result rows are scanned linearly, stopping once UNKNOWN is 
established.
+  //
+  // Per-candidate three-valued logic: TRUE if every field matches; UNKNOWN if 
no field is
+  // definitively unequal but at least one comparison involves null; FALSE 
otherwise.
+  private def evalMultiColumn(inputRow: InternalRow): Any = {
+    // Current behavior (ANSI on, or legacyNullInEmptyBehavior=false): IN 
(empty set) is always
+    // FALSE without evaluating the LHS. Legacy behavior (ANSI off by default) 
returns NULL when
+    // the LHS is null and FALSE otherwise, requiring the LHS to be evaluated. 
Mirror InSet.eval's
+    // guard exactly: skip child.eval only when legacyNullInEmptyBehavior is 
false (SPARK-44550).
+    if (result.isEmpty && !legacyNullInEmptyBehavior) return false
+    val value = child.eval(inputRow)
+    if (value == null) return null
+    val inputStruct = value.asInstanceOf[InternalRow]
+    val fieldTypes = multiColFieldTypes
+    val orderings = multiColFieldOrderings
+    val numFields = fieldTypes.length
+    // Cache lazy accessors in locals so the loop bodies do not re-enter them 
on every iteration.
+    val nullRows = multiColNullRows
+    val nonNullSet = multiColNonNullSet
+
+    if (!inputStruct.anyNull) {

Review Comment:
   **Nit (P3):** Legacy mode must evaluate `child`, but once `value != null` 
and `result.isEmpty`, the result is definitely false. This path currently 
initializes both empty TreeSets and scans `inputStruct.anyNull` for every outer 
row. Please add `if (result.isEmpty) return false` immediately after the `value 
== null` guard; that preserves legacy evaluation semantics while bypassing the 
redundant index and per-field work.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/subquery.scala:
##########
@@ -165,14 +197,193 @@ case class InSubqueryExec(
     }
   }
 
+  // Invariant schema/ordering data for the multi-column evaluator, computed 
once after the result
+  // is available. @transient so that serialization (result=null) does not 
trigger evaluation.
+  @transient private lazy val multiColFieldTypes: Array[DataType] =
+    plan.output.map(_.dataType).toArray
+  @transient private lazy val multiColFieldOrderings: Array[Ordering[Any]] =
+    multiColFieldTypes.map(TypeUtils.getInterpretedOrdering)
+  // Struct-level ordering used to index fully non-null result rows in a 
TreeSet.
+  @transient private lazy val multiColRowOrdering: Ordering[InternalRow] =
+    
TypeUtils.getInterpretedOrdering(child.dataType).asInstanceOf[Ordering[InternalRow]]
+
+  // Split collected rows using the struct-level Catalyst ordering so that 
duplicates are
+  // deduplicated. Fully non-null rows go into multiColNonNullSet (TreeSet) 
for O(log n)
+  // membership tests. Null-containing rows are deduplicated via a temporary 
TreeSet and then
+  // stored as multiColNullRows (Array) for linear scanning; each distinct 
null-containing row
+  // is thus scanned at most once per outer row regardless of RHS duplicate 
multiplicity.
+  // See SPARK-58481.
+  @transient private lazy val (multiColNonNullSet, multiColNullRows) = {
+    val withNull = TreeSet.newBuilder[InternalRow](multiColRowOrdering)
+    val nonNull = TreeSet.newBuilder[InternalRow](multiColRowOrdering)
+    result.foreach { r =>
+      val row = r.asInstanceOf[InternalRow]
+      if (row.anyNull) withNull += row else nonNull += row
+    }
+    (nonNull.result(), withNull.result().toArray)
+  }
+
+  // Three-valued IN semantics for multi-column subqueries.
+  // Result rows are InternalRow objects; InSet's TreeSet uses Catalyst 
ordering, but membership
+  // cannot distinguish a definitively-false candidate from an indeterminate 
one.
+  //
+  // When the LHS struct has no null fields:
+  //   Fast path: O(log n) TreeSet lookup against fully non-null result rows 
for TRUE.
+  //   Slow path: linear scan over null-containing result rows only for 
potential UNKNOWN.
+  //
+  // When the LHS struct has at least one null field, the fast path cannot be 
used (a null LHS
+  // field produces UNKNOWN against any non-null RHS row whose non-null fields 
all match). Both
+  // sets of result rows are scanned linearly, stopping once UNKNOWN is 
established.
+  //
+  // Per-candidate three-valued logic: TRUE if every field matches; UNKNOWN if 
no field is
+  // definitively unequal but at least one comparison involves null; FALSE 
otherwise.
+  private def evalMultiColumn(inputRow: InternalRow): Any = {
+    // Current behavior (ANSI on, or legacyNullInEmptyBehavior=false): IN 
(empty set) is always
+    // FALSE without evaluating the LHS. Legacy behavior (ANSI off by default) 
returns NULL when
+    // the LHS is null and FALSE otherwise, requiring the LHS to be evaluated. 
Mirror InSet.eval's
+    // guard exactly: skip child.eval only when legacyNullInEmptyBehavior is 
false (SPARK-44550).
+    if (result.isEmpty && !legacyNullInEmptyBehavior) return false
+    val value = child.eval(inputRow)
+    if (value == null) return null
+    val inputStruct = value.asInstanceOf[InternalRow]
+    val fieldTypes = multiColFieldTypes
+    val orderings = multiColFieldOrderings
+    val numFields = fieldTypes.length
+    // Cache lazy accessors in locals so the loop bodies do not re-enter them 
on every iteration.
+    val nullRows = multiColNullRows
+    val nonNullSet = multiColNonNullSet
+
+    if (!inputStruct.anyNull) {
+      // Fast path: indexed lookup among fully non-null candidates.
+      if (nonNullSet.contains(inputStruct)) return true
+      // No null-containing candidates: no path to UNKNOWN, result is FALSE.
+      if (nullRows.isEmpty) return false
+      // Materialize LHS fields once before the candidate scans to avoid 
repeated get() calls
+      // inside the per-candidate loop.
+      val inputFields = Array.tabulate(numFields)(i => inputStruct.get(i, 
fieldTypes(i)))
+      // Slow path: scan null-containing candidates for potential UNKNOWN.
+      // Stop early once hasUnknown is set: the indexed lookup already ruled 
out TRUE,
+      // and every row here contains NULL, so no later candidate can improve 
UNKNOWN to TRUE.
+      var hasUnknown = false
+      var i = 0
+      while (i < nullRows.length && !hasUnknown) {
+        val candidate = nullRows(i)
+        var fieldIdx = 0
+        var candidateIsUnknown = false
+        var candidateIsFalse = false
+        while (fieldIdx < numFields && !candidateIsFalse) {
+          val candidateField = candidate.get(fieldIdx, fieldTypes(fieldIdx))
+          if (candidateField == null) {
+            candidateIsUnknown = true
+          } else if (orderings(fieldIdx).compare(inputFields(fieldIdx), 
candidateField) != 0) {
+            candidateIsFalse = true
+          }
+          fieldIdx += 1
+        }
+        if (!candidateIsFalse && candidateIsUnknown) hasUnknown = true
+        i += 1
+      }
+      if (hasUnknown) null else false
+    } else {
+      // LHS has at least one null field: must scan both result sets, stopping 
once UNKNOWN
+      // is established (a null LHS field can produce UNKNOWN against any 
non-null RHS row
+      // whose other fields all match).
+      // No candidates at all: result is FALSE (no match possible).
+      if (nullRows.isEmpty && nonNullSet.isEmpty) return false
+      // Materialize LHS fields once before the scans to avoid repeated get() 
calls.
+      val inputFields = Array.tabulate(numFields)(i => inputStruct.get(i, 
fieldTypes(i)))
+      var hasUnknown = false
+      // Scan null-containing result rows first.
+      var i = 0
+      while (i < nullRows.length && !hasUnknown) {
+        val candidate = nullRows(i)
+        var fieldIdx = 0
+        var candidateIsUnknown = false
+        var candidateIsFalse = false
+        while (fieldIdx < numFields && !candidateIsFalse) {
+          val inputField = inputFields(fieldIdx)
+          if (inputField == null) {
+            // LHS field is null: UNKNOWN regardless of the RHS value; skip 
candidate.get.
+            candidateIsUnknown = true
+          } else {
+            val candidateField = candidate.get(fieldIdx, fieldTypes(fieldIdx))
+            if (candidateField == null) {
+              candidateIsUnknown = true
+            } else if (orderings(fieldIdx).compare(inputField, candidateField) 
!= 0) {
+              candidateIsFalse = true
+            }
+          }
+          fieldIdx += 1
+        }
+        if (!candidateIsFalse && candidateIsUnknown) hasUnknown = true
+        i += 1
+      }
+      // Scan non-null rows: a null LHS comparison is UNKNOWN unless a 
non-null field differs.
+      val nonNullIter = nonNullSet.iterator
+      while (nonNullIter.hasNext && !hasUnknown) {
+        val candidate = nonNullIter.next()
+        var fieldIdx = 0
+        var candidateIsUnknown = false
+        var candidateIsFalse = false
+        while (fieldIdx < numFields && !candidateIsFalse) {
+          val inputField = inputFields(fieldIdx)
+          if (inputField == null) {
+            candidateIsUnknown = true
+          } else if (orderings(fieldIdx).compare(
+              inputField, candidate.get(fieldIdx, fieldTypes(fieldIdx))) != 0) 
{
+            candidateIsFalse = true
+          }
+          fieldIdx += 1
+        }
+        if (!candidateIsFalse && candidateIsUnknown) hasUnknown = true
+      }
+      if (hasUnknown) null else false
+    }
+  }
+
   override def eval(input: InternalRow): Any = {
     prepareResult()
-    if (isResultUnavailable) true else inSet.eval(input)
+    if (isResultUnavailable) {
+      true
+    } else if (plan.output.length > 1) {
+      evalMultiColumn(input)
+    } else {
+      inSet.eval(input)
+    }
   }
 
   override def doGenCode(ctx: CodegenContext, ev: ExprCode): ExprCode = {
     prepareResult()
-    if (isResultUnavailable) Literal.TrueLiteral.doGenCode(ctx, ev) else 
inSet.doGenCode(ctx, ev)
+    if (isResultUnavailable) {
+      Literal.TrueLiteral.doGenCode(ctx, ev)
+    } else if (plan.output.length > 1) {
+      // Multi-column: per-candidate three-valued comparison cannot be 
expressed with InSet's
+      // generated code.  Fall back to the interpreted path via eval().
+      // Register any Nondeterministic descendants (e.g. rand() in the LHS) 
for partition-level
+      // initialization, mirroring CodegenFallback's protocol.
+      val resultIdx = ctx.references.length
+      ctx.references += this

Review Comment:
   **Non-blocking (P2):** `ctx.references += this` puts the whole 
`InSubqueryExec` into `WholeStageCodegenEvaluatorFactory.references`. Although 
`result` is transient, `plan` is not, so task serialization retains the 
`BaseSubqueryExec` and its physical subtree after the compact result has 
already been broadcast. A large file-backed subquery can therefore bloat task 
closures or exceed task-size limits. Please register a lightweight evaluator 
containing only the LHS expression, output types/orderings, and broadcast 
result, so the driver plan does not cross the executor boundary.
   
   **Recommended change:** Reference a lightweight serializable multi-column 
evaluator from generated code instead of the full InSubqueryExec.
   
   **Why this works:** Separate the child, collected result, schema/orderings, 
and multi-column evaluation logic into executor state that excludes 
BaseSubqueryExec; register that object while preserving explicit 
Nondeterministic initialization.
   
   **Scope:** SQL core subquery execution and focused codegen or serialization 
coverage.
   
   **Compatibility:** Keep SQL three-valued results, legacy empty-RHS behavior, 
the broadcast lifecycle, Dynamic Partition Pruning behavior, and 
nondeterministic initialization unchanged.
   
   **Risks:** Separating evaluator state could omit required output ordering 
information. Moving the child expression could accidentally bypass 
partition-level initialization for nondeterministic descendants.
   
   **Constraints:** Do not serialize BaseSubqueryExec or its driver SparkPlan 
into task references. Keep the single-column InSet path unchanged.
   
   **Success:** Whole-stage evaluator references contain only compact subquery 
value/evaluator state and no BaseSubqueryExec, while multi-column and 
nondeterministic codegen behavior remains unchanged.



##########
sql/core/src/test/scala/org/apache/spark/sql/SubquerySuite.scala:
##########
@@ -2678,4 +2681,395 @@ class SubquerySuite extends SharedSparkSession
 
     assert(exposedAttribute.exprId == outerReferenceAttribute.exprId)
   }
+
+  test("SPARK-58481: InSubqueryExec nullable correctly accounts for subquery 
output nullability") {
+    // 5 NOT IN (99, NULL) is UNKNOWN, not TRUE or FALSE.  A join condition 
that is not TRUE
+    // matches no rows, so a FULL OUTER JOIN must emit null-padded rows for 
every row in each
+    // side -- 3 + 3 = 6 null-padded rows -- not the full cross product (9 
rows).
+    //
+    // The predicate 5 NOT IN (...) is a constant that references neither join 
input, so
+    // RewritePredicateSubquery returns the original Join unchanged and 
PlanSubqueries always
+    // constructs InSubqueryExec here regardless of any optimizer config.
+    withTable("t0", "t1", "t3") {
+      sql("CREATE TABLE t0(c0 INT) USING PARQUET")
+      sql("INSERT INTO t0 VALUES (1), (2), (3)")
+      sql("CREATE TABLE t1(c0 INT) USING PARQUET")
+      sql("INSERT INTO t1 VALUES (10), (20), (30)")
+      sql("CREATE TABLE t3(c0 INT) USING PARQUET")
+      sql("INSERT INTO t3 VALUES (99), (CAST(NULL AS INT))")
+
+      val query = sql(
+        "SELECT t0.c0, t1.c0 FROM t1 FULL OUTER JOIN t0" +
+        " ON (5 NOT IN (SELECT t3.c0 FROM t3))")
+      // Verify the physical plan contains InSubqueryExec so the nullability 
fix is
+      // actually exercised and not silently bypassed by an optimizer path.
+      // Use AdaptiveSparkPlanHelper.find to traverse AQE wrapper nodes, and 
recurse
+      // into expression subtrees via Expression.exists to find the wrapped 
InSubqueryExec.
+      assert(
+        find(query.queryExecution.executedPlan) { p =>
+          p.expressions.exists(_.exists {
+            case _: InSubqueryExec => true
+            case _ => false
+          })
+        }.isDefined,
+        "Expected InSubqueryExec in the physical plan")
+      // Unmatched t1 rows null-pad t0; unmatched t0 rows null-pad t1.
+      checkAnswer(query, Seq(
+        Row(null, 10), Row(null, 20), Row(null, 30),
+        Row(1, null), Row(2, null), Row(3, null)))
+    }
+  }
+
+  test("SPARK-58481: multi-column IN subquery with nullable non-head output is 
nullable") {
+    // Disable the optimizer's join-condition IN rewrite so the query 
exercises InSubqueryExec.
+    // Use VALUES-derived temp views: their nullability is inferred from the 
literals (no NULL
+    // literal => non-nullable), rather than declared and then widened. 
Parquet file-source
+    // analysis applies dataSchema.asNullable regardless of DDL NOT NULL, 
which would defeat
+    // the nullability control this test relies on.
+    withSQLConf(
+      
"spark.sql.optimizer.optimizeUncorrelatedInSubqueriesInJoinCondition.enabled" 
-> "false"
+    ) {
+      // Case A: NULL after a definitive match (null in non-head position 
after matching head).
+      // RHS: (99,99) and (1,NULL).
+      //   (1,1) vs (99,99): first field 1!=99 => FALSE.
+      //   (1,1) vs (1,NULL): first fields equal, second null => UNKNOWN.
+      //   Overall for (1,1): UNKNOWN => NOT IN = null-padded.
+      //   (2,2) vs both: all FALSE => NOT IN = TRUE => joins with both rhs 
rows.
+      withTempView("lhs", "rhs") {
+        sql("CREATE TEMPORARY VIEW lhs AS SELECT * FROM VALUES (1, 1), (2, 2) 
AS t(a, b)")
+        sql(
+          """CREATE TEMPORARY VIEW rhs AS
+            |SELECT * FROM VALUES (99, 99), (1, CAST(NULL AS INT)) AS t(a, 
b)""".stripMargin)
+        checkAnswer(
+          sql(
+            """SELECT lhs.a, rhs.a FROM lhs FULL OUTER JOIN rhs
+              |ON ((lhs.a, lhs.b) NOT IN (SELECT a, b FROM 
rhs))""".stripMargin),
+          Seq(Row(1, null), Row(2, 99), Row(2, 1)))
+      }
+
+      // Case B: NULL in head position followed by a definitive mismatch in a 
later field.
+      // RHS: (NULL, 99).
+      //   (1,1) vs (NULL,99): first field null => UNKNOWN so far; second 
field 1!=99 => FALSE.
+      //   A later definitive mismatch must override the earlier UNKNOWN: 
result is FALSE,
+      //   NOT IN = TRUE. A field-order regression would leave (1,1) as 
UNKNOWN instead.
+      withTempView("lhs2", "rhs2") {
+        sql("CREATE TEMPORARY VIEW lhs2 AS SELECT * FROM VALUES (1, 1) AS t(a, 
b)")
+        sql(
+          """CREATE TEMPORARY VIEW rhs2 AS
+            |SELECT * FROM VALUES (CAST(NULL AS INT), 99) AS t(a, 
b)""".stripMargin)
+        // (1,1) NOT IN ((NULL,99)): second field 1!=99 makes the candidate 
FALSE =>
+        // NOT IN = TRUE => inner join returns the single matching row.
+        checkAnswer(
+          sql(
+            """SELECT lhs2.a FROM lhs2 JOIN (SELECT 1 AS a)
+              |ON ((lhs2.a, lhs2.b) NOT IN (SELECT a, b FROM 
rhs2))""".stripMargin),
+          Seq(Row(1)))
+      }
+
+      // Case C: UNKNOWN candidate followed by an exact-match candidate => IN 
= TRUE.
+      // RHS: (1,NULL) and (1,1).
+      //   (1,1) vs (1,NULL): first fields equal, second null => UNKNOWN.
+      //   (1,1) vs (1,1): exact match => TRUE.
+      //   The exact match must dominate the UNKNOWN: IN = TRUE, NOT IN = 
FALSE.
+      //   A candidate-order regression would short-circuit on UNKNOWN and 
miss the TRUE.
+      withTempView("lhs3", "rhs3") {
+        sql("CREATE TEMPORARY VIEW lhs3 AS SELECT * FROM VALUES (1, 1) AS t(a, 
b)")
+        sql(
+          """CREATE TEMPORARY VIEW rhs3 AS
+            |SELECT * FROM VALUES (1, CAST(NULL AS INT)), (1, 1) AS t(a, 
b)""".stripMargin)
+        // (1,1) IN ((1,NULL),(1,1)): exact match exists => IN = TRUE => inner 
join returns row.
+        checkAnswer(
+          sql(
+            """SELECT lhs3.a FROM lhs3 JOIN (SELECT 1 AS a)
+              |ON ((lhs3.a, lhs3.b) IN (SELECT a, b FROM 
rhs3))""".stripMargin),
+          Seq(Row(1)))
+      }
+    }
+  }
+
+  test("SPARK-58481: multi-column IN subquery uses Catalyst ordering for 
BinaryType fields") {
+    // Object.equals on Array[Byte] compares by identity, not value; Catalyst 
ordering compares
+    // by content. A multi-column IN where one field is BinaryType would 
incorrectly return FALSE
+    // (no match) with JVM equality even when the bytes are equal. Use an 
inner join to keep the
+    // assertion simple: the join condition is TRUE iff the IN match succeeds.
+    withSQLConf(
+      
"spark.sql.optimizer.optimizeUncorrelatedInSubqueriesInJoinCondition.enabled" 
-> "false"
+    ) {
+      withTable("lbin", "rbin") {
+        sql("CREATE TABLE lbin(id INT NOT NULL, b BINARY NOT NULL) USING 
PARQUET")
+        sql("INSERT INTO lbin VALUES (1, X'01')")
+        sql("CREATE TABLE rbin(id INT NOT NULL, b BINARY NOT NULL) USING 
PARQUET")
+        sql("INSERT INTO rbin VALUES (1, X'01')")
+        // (1, 0x01) IN ((1, 0x01)) must be TRUE; the join should return one 
row.
+        checkAnswer(
+          sql(
+            """SELECT lbin.id FROM lbin JOIN rbin
+              |ON ((lbin.id, lbin.b) IN (SELECT id, b FROM 
rbin))""".stripMargin),
+          Seq(Row(1)))
+      }
+    }
+  }
+
+  test("SPARK-58481: multi-column NOT IN with nullable LHS and non-nullable 
RHS is nullable") {
+    // CreateNamedStruct.nullable is always false, so child.nullable would 
return false for a
+    // multi-column LHS even when individual fields are nullable. The 
generated NOT IN code
+    // would then suppress null handling and turn UNKNOWN into TRUE, producing 
wrong results.
+    withSQLConf(
+      
"spark.sql.optimizer.optimizeUncorrelatedInSubqueriesInJoinCondition.enabled" 
-> "false"
+    ) {
+      // Case A: RHS has a null-containing row. LHS (NULL, 2) vs RHS (99, 2): 
second fields
+      // equal, first field null => UNKNOWN; exercises the nullRows scan in 
evalMultiColumn.
+      // Fixture: lhs.a is nullable; rhs columns are NOT NULL (VALUES-derived 
to avoid Parquet
+      // dataSchema.asNullable widening that would defeat the RHS 
non-nullability control).
+      withTable("lhs") {
+        withTempView("rhs") {
+          sql("CREATE TABLE lhs(a INT, b INT NOT NULL) USING PARQUET")
+          sql("INSERT INTO lhs VALUES (1, 1), (NULL, 2)")
+          // rhs as VALUES view: both columns inferred non-nullable from 
all-literal rows.
+          sql("CREATE TEMPORARY VIEW rhs AS SELECT * FROM VALUES (99, 2) AS 
t(a, b)")
+          // (NULL, 2) NOT IN ((99,2)): second fields match, first is null => 
UNKNOWN
+          //   => join condition not TRUE => (NULL,2) is null-padded: 
Row(null, null).
+          // (1, 1)   NOT IN ((99,2)): first field 1!=99 => FALSE => NOT IN = 
TRUE
+          //   => (1,1) joins with rhs(99,2): Row(1, 99). rhs matched; no 
null-padded rhs.
+          checkAnswer(
+            sql(
+              """SELECT lhs.a, rhs.a FROM lhs FULL OUTER JOIN rhs
+                |ON ((lhs.a, lhs.b) NOT IN (SELECT a, b FROM 
rhs))""".stripMargin),
+            Seq(Row(1, 99), Row(null, null)))
+        }
+      }
+
+      // Case B: RHS is fully non-null (goes into nonNullSet, not nullRows). 
LHS (NULL, 1) vs
+      // RHS (99, 2): first field null, but second field 1 != 2 is a 
definitive mismatch that
+      // makes the candidate FALSE. An implementation that stops at the first 
UNKNOWN instead
+      // of continuing to check later fields would return UNKNOWN here instead 
of FALSE, causing
+      // NOT IN to be UNKNOWN and the inner join to return no rows.
+      withTempView("lhs_nn", "rhs_nn") {
+        sql("CREATE TEMPORARY VIEW lhs_nn AS " +
+          "SELECT * FROM VALUES (CAST(NULL AS INT), 1) AS t(a, b)")
+        sql("CREATE TEMPORARY VIEW rhs_nn AS " +
+          "SELECT * FROM VALUES (99, 2) AS t(a, b)")
+        // (NULL, 1) NOT IN ((99, 2)): second field 1 != 2 => candidate FALSE 
=> NOT IN TRUE.
+        // Inner join returns the single lhs row.
+        checkAnswer(
+          sql(
+            """SELECT lhs_nn.b FROM lhs_nn JOIN (SELECT 1 AS x)
+              |ON ((lhs_nn.a, lhs_nn.b) NOT IN (SELECT a, b FROM rhs_nn))
+              |""".stripMargin),
+          Seq(Row(1)))
+      }
+    }
+  }
+
+  test("SPARK-58481: LEGACY_IN_SUBQUERY_NULLABILITY suppresses RHS-only 
nullability") {
+    // Legacy mode suppresses only RHS-derived nullability (plan.output 
nullable).
+    // A non-nullable scalar LHS (Literal 5) has lhsNullable=false; with RHS 
suppressed,
+    // nullable=false. The generated code omits null handling and NOT IN on a 
subquery that
+    // returns NULL evaluates to TRUE -- the pre-fix single-column behaviour 
the flag preserves.
+    // Note: intentionally codegen-specific. The interpreted path correctly 
returns UNKNOWN
+    // regardless of nullable (6 rows); the assertion of 9 verifies codegen 
ran.
+    withSQLConf(
+      SQLConf.LEGACY_IN_SUBQUERY_NULLABILITY.key -> "true",
+      
"spark.sql.optimizer.optimizeUncorrelatedInSubqueriesInJoinCondition.enabled" 
-> "false"
+    ) {
+      withTable("t0", "t1", "t3") {
+        sql("CREATE TABLE t0(c0 INT) USING PARQUET")
+        sql("INSERT INTO t0 VALUES (1), (2), (3)")
+        sql("CREATE TABLE t1(c0 INT) USING PARQUET")
+        sql("INSERT INTO t1 VALUES (10), (20), (30)")
+        sql("CREATE TABLE t3(c0 INT) USING PARQUET")
+        sql("INSERT INTO t3 VALUES (99), (CAST(NULL AS INT))")
+
+        // Legacy: lhsNullable=false (Literal 5), rhsNullable suppressed => 
nullable=false.
+        // Generated code suppresses null; 5 NOT IN (99, NULL) evaluates to 
TRUE.
+        // FULL OUTER JOIN condition is TRUE => full cross product of 3 x 3 = 
9 rows.
+        assert(sql(
+          "SELECT t0.c0, t1.c0 FROM t1 FULL OUTER JOIN t0 ON (5 NOT IN (SELECT 
t3.c0 FROM t3))")
+          .count() === 9)
+      }
+    }
+  }
+
+  test("SPARK-58481: LEGACY_IN_SUBQUERY_NULLABILITY preserves nullable LHS 
fields " +
+      "for multi-column IN subqueries") {
+    // Legacy mode suppresses RHS nullability but preserves LHS field 
nullability.
+    // With a nullable LHS field, lhsNullable=true even in legacy mode, so 
nullable=true.
+    // Generated NOT IN code propagates UNKNOWN correctly; result is identical 
to non-legacy.
+    withSQLConf(
+      SQLConf.LEGACY_IN_SUBQUERY_NULLABILITY.key -> "true",
+      
"spark.sql.optimizer.optimizeUncorrelatedInSubqueriesInJoinCondition.enabled" 
-> "false"
+    ) {
+      withTable("lhs", "rhs") {
+        sql("CREATE TABLE lhs(a INT, b INT NOT NULL) USING PARQUET")
+        sql("INSERT INTO lhs VALUES (1, 1), (NULL, 2)")
+        sql("CREATE TABLE rhs(a INT NOT NULL, b INT NOT NULL) USING PARQUET")
+        sql("INSERT INTO rhs VALUES (99, 2)")
+        // (NULL, 2) NOT IN ((99,2)): first field null => UNKNOWN => 
null-padded: Row(null, null).
+        // (1, 1)   NOT IN ((99,2)): 1!=99 => FALSE => NOT IN=TRUE => joins: 
Row(1, 99).
+        checkAnswer(
+          sql(
+            """SELECT lhs.a, rhs.a FROM lhs FULL OUTER JOIN rhs
+              |ON ((lhs.a, lhs.b) NOT IN (SELECT a, b FROM 
rhs))""".stripMargin),
+          Seq(Row(1, 99), Row(null, null)))
+      }
+    }
+  }
+
+  test("SPARK-58481: nondeterministic LHS in a Project IN subquery initializes 
via codegen") {
+    // RewritePredicateSubquery does not rewrite Project expressions, and 
CheckAnalysis permits

Review Comment:
   **Non-blocking (P2):** This premise is reversed: `Project` is a `UnaryNode`, 
and `RewritePredicateSubquery` sends any matching unary node through 
`handleUnaryNode`, which replaces the IN expression with an `ExistenceJoin`. 
This query can pass without constructing `InSubqueryExec`, so it does not 
verify `doGenCode` or nondeterministic initialization. Please remove this 
projection test and its end-to-end coverage claim; the focused 
direct-construction `GeneratePredicate` test below is the regression that 
actually executes this codegen branch.



##########
sql/core/src/test/scala/org/apache/spark/sql/SubquerySuite.scala:
##########
@@ -2678,4 +2681,395 @@ class SubquerySuite extends SharedSparkSession
 
     assert(exposedAttribute.exprId == outerReferenceAttribute.exprId)
   }
+
+  test("SPARK-58481: InSubqueryExec nullable correctly accounts for subquery 
output nullability") {
+    // 5 NOT IN (99, NULL) is UNKNOWN, not TRUE or FALSE.  A join condition 
that is not TRUE
+    // matches no rows, so a FULL OUTER JOIN must emit null-padded rows for 
every row in each
+    // side -- 3 + 3 = 6 null-padded rows -- not the full cross product (9 
rows).
+    //
+    // The predicate 5 NOT IN (...) is a constant that references neither join 
input, so
+    // RewritePredicateSubquery returns the original Join unchanged and 
PlanSubqueries always
+    // constructs InSubqueryExec here regardless of any optimizer config.
+    withTable("t0", "t1", "t3") {
+      sql("CREATE TABLE t0(c0 INT) USING PARQUET")
+      sql("INSERT INTO t0 VALUES (1), (2), (3)")
+      sql("CREATE TABLE t1(c0 INT) USING PARQUET")
+      sql("INSERT INTO t1 VALUES (10), (20), (30)")
+      sql("CREATE TABLE t3(c0 INT) USING PARQUET")
+      sql("INSERT INTO t3 VALUES (99), (CAST(NULL AS INT))")
+
+      val query = sql(
+        "SELECT t0.c0, t1.c0 FROM t1 FULL OUTER JOIN t0" +
+        " ON (5 NOT IN (SELECT t3.c0 FROM t3))")
+      // Verify the physical plan contains InSubqueryExec so the nullability 
fix is
+      // actually exercised and not silently bypassed by an optimizer path.
+      // Use AdaptiveSparkPlanHelper.find to traverse AQE wrapper nodes, and 
recurse
+      // into expression subtrees via Expression.exists to find the wrapped 
InSubqueryExec.
+      assert(
+        find(query.queryExecution.executedPlan) { p =>
+          p.expressions.exists(_.exists {
+            case _: InSubqueryExec => true
+            case _ => false
+          })
+        }.isDefined,
+        "Expected InSubqueryExec in the physical plan")
+      // Unmatched t1 rows null-pad t0; unmatched t0 rows null-pad t1.
+      checkAnswer(query, Seq(
+        Row(null, 10), Row(null, 20), Row(null, 30),
+        Row(1, null), Row(2, null), Row(3, null)))
+    }
+  }
+
+  test("SPARK-58481: multi-column IN subquery with nullable non-head output is 
nullable") {
+    // Disable the optimizer's join-condition IN rewrite so the query 
exercises InSubqueryExec.
+    // Use VALUES-derived temp views: their nullability is inferred from the 
literals (no NULL
+    // literal => non-nullable), rather than declared and then widened. 
Parquet file-source
+    // analysis applies dataSchema.asNullable regardless of DDL NOT NULL, 
which would defeat
+    // the nullability control this test relies on.
+    withSQLConf(
+      
"spark.sql.optimizer.optimizeUncorrelatedInSubqueriesInJoinCondition.enabled" 
-> "false"
+    ) {
+      // Case A: NULL after a definitive match (null in non-head position 
after matching head).
+      // RHS: (99,99) and (1,NULL).
+      //   (1,1) vs (99,99): first field 1!=99 => FALSE.
+      //   (1,1) vs (1,NULL): first fields equal, second null => UNKNOWN.
+      //   Overall for (1,1): UNKNOWN => NOT IN = null-padded.
+      //   (2,2) vs both: all FALSE => NOT IN = TRUE => joins with both rhs 
rows.
+      withTempView("lhs", "rhs") {
+        sql("CREATE TEMPORARY VIEW lhs AS SELECT * FROM VALUES (1, 1), (2, 2) 
AS t(a, b)")
+        sql(
+          """CREATE TEMPORARY VIEW rhs AS
+            |SELECT * FROM VALUES (99, 99), (1, CAST(NULL AS INT)) AS t(a, 
b)""".stripMargin)
+        checkAnswer(
+          sql(
+            """SELECT lhs.a, rhs.a FROM lhs FULL OUTER JOIN rhs
+              |ON ((lhs.a, lhs.b) NOT IN (SELECT a, b FROM 
rhs))""".stripMargin),
+          Seq(Row(1, null), Row(2, 99), Row(2, 1)))
+      }
+
+      // Case B: NULL in head position followed by a definitive mismatch in a 
later field.
+      // RHS: (NULL, 99).
+      //   (1,1) vs (NULL,99): first field null => UNKNOWN so far; second 
field 1!=99 => FALSE.
+      //   A later definitive mismatch must override the earlier UNKNOWN: 
result is FALSE,
+      //   NOT IN = TRUE. A field-order regression would leave (1,1) as 
UNKNOWN instead.
+      withTempView("lhs2", "rhs2") {
+        sql("CREATE TEMPORARY VIEW lhs2 AS SELECT * FROM VALUES (1, 1) AS t(a, 
b)")
+        sql(
+          """CREATE TEMPORARY VIEW rhs2 AS
+            |SELECT * FROM VALUES (CAST(NULL AS INT), 99) AS t(a, 
b)""".stripMargin)
+        // (1,1) NOT IN ((NULL,99)): second field 1!=99 makes the candidate 
FALSE =>
+        // NOT IN = TRUE => inner join returns the single matching row.
+        checkAnswer(
+          sql(
+            """SELECT lhs2.a FROM lhs2 JOIN (SELECT 1 AS a)
+              |ON ((lhs2.a, lhs2.b) NOT IN (SELECT a, b FROM 
rhs2))""".stripMargin),
+          Seq(Row(1)))
+      }
+
+      // Case C: UNKNOWN candidate followed by an exact-match candidate => IN 
= TRUE.
+      // RHS: (1,NULL) and (1,1).
+      //   (1,1) vs (1,NULL): first fields equal, second null => UNKNOWN.
+      //   (1,1) vs (1,1): exact match => TRUE.
+      //   The exact match must dominate the UNKNOWN: IN = TRUE, NOT IN = 
FALSE.
+      //   A candidate-order regression would short-circuit on UNKNOWN and 
miss the TRUE.
+      withTempView("lhs3", "rhs3") {
+        sql("CREATE TEMPORARY VIEW lhs3 AS SELECT * FROM VALUES (1, 1) AS t(a, 
b)")
+        sql(
+          """CREATE TEMPORARY VIEW rhs3 AS
+            |SELECT * FROM VALUES (1, CAST(NULL AS INT)), (1, 1) AS t(a, 
b)""".stripMargin)
+        // (1,1) IN ((1,NULL),(1,1)): exact match exists => IN = TRUE => inner 
join returns row.
+        checkAnswer(
+          sql(
+            """SELECT lhs3.a FROM lhs3 JOIN (SELECT 1 AS a)
+              |ON ((lhs3.a, lhs3.b) IN (SELECT a, b FROM 
rhs3))""".stripMargin),
+          Seq(Row(1)))
+      }
+    }
+  }
+
+  test("SPARK-58481: multi-column IN subquery uses Catalyst ordering for 
BinaryType fields") {
+    // Object.equals on Array[Byte] compares by identity, not value; Catalyst 
ordering compares
+    // by content. A multi-column IN where one field is BinaryType would 
incorrectly return FALSE
+    // (no match) with JVM equality even when the bytes are equal. Use an 
inner join to keep the
+    // assertion simple: the join condition is TRUE iff the IN match succeeds.
+    withSQLConf(
+      
"spark.sql.optimizer.optimizeUncorrelatedInSubqueriesInJoinCondition.enabled" 
-> "false"
+    ) {
+      withTable("lbin", "rbin") {
+        sql("CREATE TABLE lbin(id INT NOT NULL, b BINARY NOT NULL) USING 
PARQUET")
+        sql("INSERT INTO lbin VALUES (1, X'01')")
+        sql("CREATE TABLE rbin(id INT NOT NULL, b BINARY NOT NULL) USING 
PARQUET")
+        sql("INSERT INTO rbin VALUES (1, X'01')")
+        // (1, 0x01) IN ((1, 0x01)) must be TRUE; the join should return one 
row.
+        checkAnswer(
+          sql(
+            """SELECT lbin.id FROM lbin JOIN rbin
+              |ON ((lbin.id, lbin.b) IN (SELECT id, b FROM 
rbin))""".stripMargin),
+          Seq(Row(1)))
+      }
+    }
+  }
+
+  test("SPARK-58481: multi-column NOT IN with nullable LHS and non-nullable 
RHS is nullable") {
+    // CreateNamedStruct.nullable is always false, so child.nullable would 
return false for a
+    // multi-column LHS even when individual fields are nullable. The 
generated NOT IN code
+    // would then suppress null handling and turn UNKNOWN into TRUE, producing 
wrong results.
+    withSQLConf(
+      
"spark.sql.optimizer.optimizeUncorrelatedInSubqueriesInJoinCondition.enabled" 
-> "false"
+    ) {
+      // Case A: RHS has a null-containing row. LHS (NULL, 2) vs RHS (99, 2): 
second fields

Review Comment:
   **Nit (P3):** The RHS is `(99, 2)`, so it is fully non-null and is stored in 
`multiColNonNullSet`. The UNKNOWN result comes from scanning `nonNullSet` with 
a nullable LHS, not from `nullRows`. Please update these two sentences to say 
`fully non-null RHS` and `nonNullSet scan`; a genuinely null-containing case 
can carry the `nullRows` coverage claim.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to