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]

Reply via email to