[ 
https://issues.apache.org/jira/browse/SPARK-60046?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Max Gekk updated SPARK-60046:
-----------------------------
    Affects Version/s: 2.3.0
                           (was: 5.0.0)

> Sort-merge join keys that split their code trip SPARK-22668's testing 
> assertion on the streamed row
> ---------------------------------------------------------------------------------------------------
>
>                 Key: SPARK-60046
>                 URL: https://issues.apache.org/jira/browse/SPARK-60046
>             Project: Spark
>          Issue Type: Improvement
>          Components: SQL
>    Affects Versions: 2.3.0
>            Reporter: Max Gekk
>            Priority: Major
>
> Under testing, the assertion of SPARK-22668 in 
> CodegenContext.splitExpressions ("split function argument ... cannot be a 
> global variable") fails for a sort-merge join whose key is an expression that 
> splits its code through splitExpressionsWithCurrentInputs, such as a Coalesce 
> with several arguments.
> SortMergeJoinExec.createJoinKey sets INPUT_ROW to the streamed row, which is 
> a mutable state field created with forceInline (smj_streamedRow_0). 
> splitExpressionsWithCurrentInputs then passes INPUT_ROW as the "InternalRow" 
> argument of the split function, and the assertion rejects it because the name 
> is in mutableStateNames.
> Reproduction (a test; the assertion is guarded by Utils.isTesting):
>   withSQLConf(
>       SQLConf.CODEGEN_METHOD_SPLIT_THRESHOLD.key -> "1",
>       SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1",
>       SQLConf.PREFER_SORTMERGEJOIN.key -> "true") {
>     spark.range(10).selectExpr("id AS a", "nullif(id, 3) AS b", "nullif(id, 
> 4) AS c",
>       "nullif(id, 5) AS d").createOrReplaceTempView("l")
>     spark.range(10).selectExpr("id AS x").createOrReplaceTempView("r")
>     sql("SELECT * FROM l JOIN r ON coalesce(b, c, d) = x").collect()
>   }
> fails with: java.lang.AssertionError: assertion failed: split function 
> argument smj_streamedRow_0 cannot be a global variable.
> This is not a user-facing bug: the assertion runs only under testing, the 
> split methods only read the row, so outside testing the generated code 
> compiles and returns the right rows. A low 
> spark.sql.codegen.methodSplitThreshold is needed to reach it with a short 
> key. It does make tests of the code generation with a low threshold fail on 
> such joins (found while fuzzing the split of CASE WHEN in whole-stage 
> codegen, SPARK-33301, https://github.com/apache/spark/pull/59225).
> Possible fixes: have createJoinKey hand the key expressions a local copy of 
> the row, or let the assertion accept INPUT_ROW where the split function only 
> reads it.



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