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