[
https://issues.apache.org/jira/browse/SPARK-60047?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Max Gekk updated SPARK-60047:
-----------------------------
Description:
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, AtLeastNNonNulls, Greatest, Least,
In, Concat, ConcatWs, Elt, FormatString, ArraysZip, MapConcat, CreateArray,
CreateNamedStruct, Stack, the hash expressions (HashExpression, which
Murmur3Hash and XxHash64 extend, and HiveHash), and the expressions in
objects.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.
was:
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.
> 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
> Priority: Major
>
> 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, AtLeastNNonNulls,
> Greatest, Least, In, Concat, ConcatWs, Elt, FormatString, ArraysZip,
> MapConcat, CreateArray, CreateNamedStruct, Stack, the hash expressions
> (HashExpression, which Murmur3Hash and XxHash64 extend, and HiveHash), and
> the expressions in objects.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]