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


##########
sql/core/src/main/scala/org/apache/spark/sql/execution/subquery.scala:
##########
@@ -188,6 +267,161 @@ case class InSubqueryExec(
     copy(child = newChild)
 }
 
+/**
+ * Lightweight, serializable evaluator for multi-column IN subqueries.
+ * Holds only what executors need: the LHS expression, collected result rows, 
field metadata,
+ * and the legacy empty-list behavior flag. Excludes plan (BaseSubqueryExec) 
so that
+ * task-closure serialization does not carry the driver's physical plan 
subtree.
+ * All @transient lazy fields are rebuilt from the serialized fieldTypes and 
resultRows
+ * after deserialization, mirroring InSubqueryExec's own pattern. See 
SPARK-58481.
+ */
+private[sql] class MultiColumnInSubqueryEvaluator(
+    val child: Expression,
+    val fieldTypes: Array[DataType],
+    val resultRows: Array[InternalRow],
+    val legacyNullInEmptyBehavior: Boolean) extends Serializable {
+
+  @transient private lazy val fieldOrderings: Array[Ordering[Any]] =
+    fieldTypes.map(TypeUtils.getInterpretedOrdering)
+
+  @transient private lazy val rowOrdering: 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 nonNullSet (TreeSet) for O(log 
n)
+  // membership tests. Null-containing rows are deduplicated via a temporary 
TreeSet and
+  // then stored as nullRows (Array) for linear scanning; each distinct 
null-containing row
+  // is thus scanned at most once per outer row regardless of RHS duplicate 
multiplicity.
+  @transient private lazy val (nonNullSet, nullRows) = {
+    val withNull = TreeSet.newBuilder[InternalRow](rowOrdering)
+    val nonNull = TreeSet.newBuilder[InternalRow](rowOrdering)
+    resultRows.foreach { r => if (r.anyNull) withNull += r else nonNull += r }
+    (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 set.
+  //
+  // 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.
+  def eval(inputRow: InternalRow): Any = {
+    // Current behavior (ANSI on, or legacyNullInEmptyBehavior=false): IN 
(empty set) is

Review Comment:
   **Nit (P3):** ANSI controls only the default here; an explicit 
spark.sql.legacy.nullInEmptyListBehavior=true still evaluates the LHS when ANSI 
is enabled, as the test below demonstrates. Please describe this branch as 
legacyNullInEmptyBehavior=false and mention ANSI only as the defaulting rule.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/subquery.scala:
##########
@@ -165,14 +197,61 @@ case class InSubqueryExec(
     }
   }
 
+  // Memoized lightweight evaluator instance. @transient so it is rebuilt 
(from the
+  // serializable fields) when InSubqueryExec is deserialized in an executor; 
constructed
+  // only for the multi-column path after prepareResult() has populated result.
+  @transient private lazy val multiColEvaluator: 
MultiColumnInSubqueryEvaluator = {
+    val rows = result.map(_.asInstanceOf[InternalRow])
+    val fTypes = plan.output.map(_.dataType).toArray
+    new MultiColumnInSubqueryEvaluator(child, fTypes, rows, 
legacyNullInEmptyBehavior)
+  }
+
   override def eval(input: InternalRow): Any = {
     prepareResult()
-    if (isResultUnavailable) true else inSet.eval(input)
+    if (isResultUnavailable) {
+      true
+    } else if (plan.output.length > 1) {
+      multiColEvaluator.eval(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 the lightweight 
evaluator.
+      // Register the evaluator (not the full InSubqueryExec) so that task 
closures do not

Review Comment:
   **Non-blocking (P2):** The evaluator still captures the entire child 
expression. For a multi-column LHS containing a scalar subquery, PlanSubqueries 
embeds an execution ScalarSubquery whose non-transient BaseSubqueryExec remains 
reachable through this reference, so driver plan state is still serialized with 
generated tasks. Please remove prepared scalar-subquery plans from the captured 
child and add a regression that proves no BaseSubqueryExec remains reachable.
   
   **Recommended change:** Sanitize the evaluator child by literalizing 
prepared execution ScalarSubquery nodes before it enters 
CodegenContext.references.
   
   **Why this works:** Transform the prepared child to replace each execution 
ScalarSubquery with its toLiteral result, then discover and register 
nondeterministic descendants from the transformed child.
   
   **Scope:** InSubqueryExec generated evaluation and SubquerySuite 
serialization coverage for multi-column left-hand expressions containing scalar 
subqueries.
   
   **Compatibility:** Preserve SQL semantics, legacy branches, and the 
single-column InSet path; only remove already-evaluated scalar-subquery plan 
objects from the serialized evaluator graph.
   
   **Risks:** Nondeterministic expression identity must remain consistent with 
partition initialization. Literal replacement must occur only after scalar 
subquery results are available.
   
   **Constraints:** No BaseSubqueryExec may remain reachable through 
CodegenContext.references. Existing nullability and empty-result compatibility 
behavior must remain unchanged.
   
   **Success:** A reachable multi-column scalar-subquery LHS serializes and 
executes with no BaseSubqueryExec reachable from task references, while the 
existing evaluator behavior remains intact.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/subquery.scala:
##########
@@ -124,8 +128,36 @@ case class InSubqueryExec(
   extends ExecSubqueryExpression with UnaryLike[Expression] with Predicate {
 
   @transient private lazy val inSet = InSet(child, result.toSet)
+  // Captured at construction time to mirror InSet's empty-result 
short-circuit (SPARK-44550).
+  private val legacyNullInEmptyBehavior = SQLConf.get.legacyNullInEmptyBehavior
 
-  override def nullable: Boolean = child.nullable
+  // Mirror the logical InSubquery.nullable: nullable when any output column 
is nullable
+  // (a nullable RHS field can produce UNKNOWN on a miss) or when any LHS 
field is nullable.
+  // For multi-column IN the LHS is a CreateNamedStruct whose top-level 
nullable is always false
+  // even when individual field expressions are nullable (SPARK-58481). Both 
PlanSubqueries and
+  // PlanAdaptiveSubqueries wrap multi-column LHS values in CreateNamedStruct, 
so matching on it
+  // here is precise for the current producers. The fallback to child.nullable 
is safe
+  // for the single-column case where child is the bare LHS expression.
+  // LEGACY_IN_SUBQUERY_NULLABILITY suppresses only RHS-derived nullability; 
LHS field nullability
+  // is preserved in both modes so a nullable LHS field can contribute UNKNOWN
+  // when no field differs.
+  // For the multi-column case (plan.output.length > 1) the LHS is always a 
generated
+  // CreateNamedStruct; a user-written single struct-valued IN also produces a 
CreateNamedStruct,

Review Comment:
   **Nit (P3):** A single struct-valued operand does not necessarily produce 
CreateNamedStruct: both planners use values.head unchanged when values.length 
== 1, so a struct attribute remains an AttributeReference. Please say that only 
multi-value IN is wrapped and that the single-value case explains the 
plan.output.length guard.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/subquery.scala:
##########
@@ -165,14 +197,61 @@ case class InSubqueryExec(
     }
   }
 
+  // Memoized lightweight evaluator instance. @transient so it is rebuilt 
(from the
+  // serializable fields) when InSubqueryExec is deserialized in an executor; 
constructed
+  // only for the multi-column path after prepareResult() has populated result.
+  @transient private lazy val multiColEvaluator: 
MultiColumnInSubqueryEvaluator = {
+    val rows = result.map(_.asInstanceOf[InternalRow])

Review Comment:
   **Nit (P3):** updateResult's normal paths already retain a runtime 
Array[InternalRow], so this map allocates and copies an extra n-element 
reference array before the index builders traverse it again. Please reuse 
result.asInstanceOf[Array[InternalRow]] on that path, with a conversion only 
for any supported Array[Any] fallback caller.



##########
sql/core/src/test/scala/org/apache/spark/sql/SubquerySuite.scala:
##########
@@ -2678,4 +2682,421 @@ 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 is fully non-null (stored in nonNullSet). LHS (NULL, 2) 
vs RHS (99, 2):
+      // second fields equal, first field null => UNKNOWN; exercises the 
nonNullSet scan with
+      // a nullable LHS. Fixture: lhs.a is nullable; rhs columns are NOT NULL 
(VALUES-derived
+      // to avoid Parquet dataSchema.asNullable widening that would defeat the 
RHS 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: InSubqueryExec.doGenCode registers Nondeterministic 
children " +
+      "for partition-level initialization") {
+    // A join-condition IN subquery with a nondeterministic LHS is rejected by 
CheckAnalysis,
+    // so this direct-construction test provides a focused unit check: it 
constructs
+    // InSubqueryExec directly (bypassing analysis) to verify that doGenCode 
correctly
+    // registers each Nondeterministic descendant of the LHS child for 
partition-level
+    // initialization. Without that registration, a Rand node's eval() would 
throw
+    // IllegalArgumentException (via require(initialized, ...)) because 
initialize() was
+    // never called.
+    //
+    // Two output columns force plan.output.length > 1, taking the 
multi-column fallback
+    // branch (with the explicit foreach loop), not the single-column InSet 
path where Rand
+    // would register itself through ordinary child codegen traversal.
+    //
+    // The test compiles the expression to bytecode via GeneratePredicate and 
executes it,
+    // so a missing initialize() call causes an actual failure rather than 
just a missing
+    // entry in partitionInitializationStatements.
+    val subplanOutput = spark
+      .range(0)
+      .select($"id".cast(IntegerType).as("a"), $"id".cast(IntegerType).as("b"))
+      .queryExecution.executedPlan.output
+    val subplan = SubqueryExec(
+      "test_subquery",
+      LocalTableScanExec(subplanOutput, Nil, None))
+    // LHS: struct of (Literal(1), Rand(42)>=0.0). Rand is always non-negative 
so the struct
+    // always equals the single RHS row (1, true), making IN = true.
+    val lhs = CreateNamedStruct(Seq(
+      Literal("a"), Literal(1),
+      Literal("b"), GreaterThanOrEqual(Rand(42L), Literal(0.0))))
+    val expr = InSubqueryExec(
+      lhs, subplan, ExprId(99),
+      isDynamicPruning = false,
+      result = Array(InternalRow(1, true)))
+    // Compile to bytecode, initialize (seeds the Rand RNG), and evaluate.
+    // If initialize() were skipped the Rand eval() would throw 
IllegalArgumentException.
+    val pred = GeneratePredicate.generate(expr, useSubexprElimination = false)
+    pred.initialize(0)
+    assert(pred.eval(InternalRow.empty) === true,
+      "Expected IN to be TRUE: (1, rand()>=0.0) IN ((1, true))")
+  }
+
+  test("SPARK-58481: MultiColumnInSubqueryEvaluator preserves shared 
Nondeterministic " +
+      "identity after Java serialization round-trip") {
+    // In doGenCode, ctx.references receives both the 
MultiColumnInSubqueryEvaluator (slot 0)
+    // and the bare Nondeterministic node (slot 1) for partition 
initialization. On the driver
+    // both slots point to the same Rand instance (evaluator.child embeds it). 
The risk is that
+    // Java serialization could write two independent copies, causing 
initialize() called on
+    // references[1] to seed a different instance than eval() uses inside the 
evaluator.
+    //
+    // Java's ObjectOutputStream tracks shared references within a single 
graph (via its
+    // internal handle table), so a single serialize/deserialize round-trip of 
the references
+    // array as one object produces back-references rather than copies. This 
test verifies
+    // that invariant explicitly using the same JavaSerializer Spark uses for 
task closures.
+    val rand = Rand(42L)
+    val lhs = CreateNamedStruct(Seq(
+      Literal("a"), Literal(1),
+      Literal("b"), GreaterThanOrEqual(rand, Literal(0.0))))
+    val evaluator = new MultiColumnInSubqueryEvaluator(
+      lhs,
+      Array(IntegerType, IntegerType),
+      Array(InternalRow(1, true)),
+      legacyNullInEmptyBehavior = false)
+    // Simulate what doGenCode does: put the evaluator and the bare 
Nondeterministic into
+    // the references array as separate slots (as 
WholeStageCodegenEvaluatorFactory would).
+    val references: Array[Any] = Array(evaluator, rand)
+    // Round-trip through JavaSerializer -- the same serializer Spark uses for 
task closures
+    // (SparkEnv always sets closureSerializer = new JavaSerializer(...)).
+    // Serializing the array directly is equivalent to serializing it as a 
field of
+    // WholeStageCodegenEvaluatorFactory: Java's ObjectOutputStream maintains 
a handle table
+    // across the entire stream, so shared-reference tracking is independent 
of nesting depth.
+    val ser = new JavaSerializer(new SparkConf()).newInstance()
+    val rt = ser.deserialize[Array[Any]](ser.serialize(references))
+    val rtEvaluator = rt(0).asInstanceOf[MultiColumnInSubqueryEvaluator]
+    val rtRand = rt(1).asInstanceOf[Rand]
+    // Find the Rand inside the deserialized evaluator's child expression tree.
+    var embeddedRand: Rand = null
+    rtEvaluator.child.foreach {
+      case n: Rand => embeddedRand = n
+      case _ =>
+    }
+    assert(embeddedRand ne null,
+      "Expected a Rand node inside the evaluator's child after round-trip")
+    // The critical identity check: the Rand node that would receive 
initialize(partitionIndex)
+    // must be the same instance eval() uses inside the evaluator -- if Java 
serialization wrote
+    // two independent copies, initialize() would seed a different instance 
than eval() uses.
+    assert(embeddedRand eq rtRand,
+      "Rand inside evaluator and separately-registered Rand must be the same 
instance " +
+      "after serialization round-trip")
+  }
+
+  test("SPARK-58481: multi-column IN with empty subquery short-circuits before 
evaluating LHS") {
+    // Under ANSI semantics, IN (empty subquery) is always FALSE regardless of 
the LHS value.
+    // evalMultiColumn must return false before calling child.eval so that a 
LHS expression

Review Comment:
   **Nit (P3):** evalMultiColumn no longer exists after the evaluator 
extraction. Please update both occurrences in this test area to 
MultiColumnInSubqueryEvaluator.eval, or simply call it the multi-column 
evaluator.



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