[
https://issues.apache.org/jira/browse/SPARK-58485?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Josh Rosen updated SPARK-58485:
-------------------------------
Summary: Broadcast hash join fails with "Subquery has not finished" when
the join key contains an uncorrelated scalar subquery (was: Broadcast hash
join fails with "Subquery has not finished" when the join key > contains an
uncorrelated scalar subquery)
> Broadcast hash join fails with "Subquery has not finished" when the join key
> contains an uncorrelated scalar subquery
> ---------------------------------------------------------------------------------------------------------------------
>
> Key: SPARK-58485
> URL: https://issues.apache.org/jira/browse/SPARK-58485
> Project: Spark
> Issue Type: Bug
> Components: SQL
> Affects Versions: 4.0.0
> Reporter: Josh Rosen
> Priority: Major
>
> The following query deterministically crashes:
> {code:java}
> CREATE TABLE t1(a INT) USING PARQUET; INSERT INTO t1 VALUES (1), (2);
> CREATE TABLE t2(b INT) USING PARQUET; INSERT INTO t2 VALUES (1), (NULL);
> CREATE TABLE t3(c INT) USING PARQUET; INSERT INTO t3 VALUES (1);
> SELECT * FROM t1 JOIN t2 ON t1.a = coalesce(t2.b, (SELECT min(c) FROM t3));
> {code}
> with error
> {code:java}
> org.apache.spark.sql.execution.QueryExecutionException:
> java.lang.IllegalArgumentException: requirement failed: Subquery
> subquery#186, [id=#527] has not finished
> at scala.Predef$.require(Predef.scala:337)
> at
> org.apache.spark.sql.execution.ScalarSubquery.eval(subquery.scala:99)
> at
> org.apache.spark.sql.catalyst.expressions.Coalesce.eval(nullExpressions.scala:90)
> at
> org.apache.spark.sql.catalyst.expressions.UnaryExpression.eval(Expression.scala:606)
> at
> org.apache.spark.sql.catalyst.expressions.InterpretedUnsafeProjection.apply(InterpretedUnsafeProjection.scala:84)
> at
> org.apache.spark.sql.execution.joins.LongHashedRelation$.apply(HashedRelation.scala:1084)
> at
> org.apache.spark.sql.execution.joins.HashedRelation$.apply(HashedRelation.scala:158)
> at
> org.apache.spark.sql.execution.joins.HashedRelationBroadcastMode.transform(HashedRelation.scala:1164)
> at
> org.apache.spark.sql.execution.joins.HashedRelationBroadcastMode.transform(HashedRelation.scala:1152)
> at
> org.apache.spark.sql.execution.exchange.BroadcastExchangeExec.$anonfun$relationFuture$1(BroadcastExchangeExec.scala:197)
> at
> org.apache.spark.sql.execution.SQLExecutionThreadLocalCaptured.$anonfun$runWith$2(SQLExecution.scala:63)
> [...] {code}
> Root cause:
> * {{ExtractEquiJoinKeys}} accepts a join key containing an uncorrelated
> scalar subquery (it has no references, so {{canEvaluate}} passes for either
> side), and the key lands in {{{}HashedRelationBroadcastMode(key:
> Seq[Expression]){}}}.
> * Scalar subqueries are normally started and awaited by
> {{SparkPlan.prepareSubqueries}} and {{{}waitForSubqueries{}}}, which discover
> them through {{{}QueryPlan.expressions{}}}.
> * That walk only surfaces Expression, Option and Iterable constructor
> arguments (QueryPlan.scala:327-341).
> * {{BroadcastMode}} is an opaque case class, so the subquery inside
> {{mode.key}} is never registered or awaited on the exchange node.
> * The broadcast thread then evaluates/code-generates the un-updated
> {{execution.ScalarSubquery}} and hits the {{require(updated)}} guard
> (subquery.scala:99/108). The join node's own {{waitForSubqueries}} does see
> its {{{}leftKeys{}}}/{{{}rightKeys{}}}, but has no happens-before
> relationship with the broadcast thread, which {{doPrepare}} launches eagerly.
> This is the same bug class as SPARK-35874 (AQE shuffle materialization not
> waiting for its subqueries), SPARK-42937 (InSubqueryExec broadcast), and
> SPARK-45584 (TakeOrderedAndProjectExec.executeCollect bypassing
> executeQuery): an execution entry point that evaluates operator expressions
> outside of executeQuery
> Possible fixes: have {{BroadcastExchangeExec}} surface {{mode.key}} to the
> subquery (register in prepareSubqueries and wait before {{mode.transform}} in
> {{{}relationFuture{}}}, mirroring SPARK-35874), or ensure planning never
> places unexecuted subquery expressions into a {{BroadcastMode}} (e.g. reject
> them in {{ExtractEquiJoinKeys}} or substitute executed results).
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]