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]