Max Gekk created SPARK-60046:
--------------------------------
Summary: 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: 5.0.0
Reporter: Max Gekk
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]