Max Gekk created SPARK-60047:
--------------------------------

             Summary: Split large expressions other than CASE WHEN under 
whole-stage codegen
                 Key: SPARK-60047
                 URL: https://issues.apache.org/jira/browse/SPARK-60047
             Project: Spark
          Issue Type: Improvement
          Components: SQL
    Affects Versions: 5.0.0
            Reporter: Max Gekk


Outside whole-stage codegen, a large expression is split into methods that take 
the input row, so that no generated method passes the JVM's 64KB limit or 
HotSpot's 8000-byte JIT limit. Inside a whole-stage codegen stage two 
mechanisms never split, because the inputs are local variables of the 
operator's method (CodegenContext.currentVars != null) rather than a row:

1. CodegenContext.splitExpressionsWithCurrentInputs returns the pieces joined 
as they are when INPUT_ROW == null || currentVars != null. About 28 call sites 
in sql/catalyst use it, among them Coalesce, Greatest and Least, In, Concat, 
ConcatWs, Elt, CreateArray, CreateNamedStruct, CreateMap, the hash expressions 
(Murmur3Hash, XxHash64, HiveHash), Stack, and the expressions in objects.scala 
and collectionOperations.scala.
2. Expression.reduceCodeSize wraps the code of any one expression into a method 
when it is longer than spark.sql.codegen.methodSplitThreshold, only where 
INPUT_ROW != null && currentVars == null. It carries the comment "TODO: support 
whole stage codegen too".

So a large Coalesce, IN list, hash or struct constructor in a projection goes 
into the operator's method, which can grow past 64KB and fail to compile (the 
stage then falls back to the non-codegen path) or past 8000 bytes, where 
HotSpot never JIT-compiles it.

SPARK-33301 (https://github.com/apache/spark/pull/59225) adds what a fix needs, 
for CASE WHEN only:
- CodegenContext.collectInputs walks the expressions a piece of code was 
generated from and returns the locals it reads, or None when no method can take 
them;
- CodegenContext.splitExpressionsWithSources takes each piece of code together 
with its source expressions and splits it per block, falling back to inline 
code for a block it cannot split;
- a stage gate that compiles the code in one piece first and splits only the 
expressions in methods past the JIT limit, so that a stage the JIT compiles 
whole keeps its code (splitting always cost 5-15% on 64-300 branches).

Proposed: pass the sources from the call sites above to 
splitExpressionsWithSources, and make reduceCodeSize use collectInputs for the 
expression it wraps, both under the same gate. Each expression needs a test 
that the split method gets the inputs it reads, as the CASE WHEN tests do, and 
the gate must keep the stages that need no split byte for byte as they are.




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