[ 
https://issues.apache.org/jira/browse/SPARK-59884?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Dongjoon Hyun reassigned SPARK-59884:
-------------------------------------

    Assignee: Chao Sun

> Avoid serializing cached columnar RDDs in whole-stage closures
> --------------------------------------------------------------
>
>                 Key: SPARK-59884
>                 URL: https://issues.apache.org/jira/browse/SPARK-59884
>             Project: Spark
>          Issue Type: Improvement
>          Components: SQL
>    Affects Versions: 4.2.0
>            Reporter: Chao Sun
>            Assignee: Chao Sun
>            Priority: Major
>              Labels: pull-request-available
>
> SparkPlan.executeColumnarRDD retains the LazyTry containing a plan's columnar 
> execution RDD, but unlike executeRDD it is not transient.
> Whole-stage sort-merge join code generation adds the join plan to 
> CodegenContext.references for executor cleanup. If a descendant columnar plan 
> has initialized executeColumnarRDD, serializing the 
> WholeStageCodegenEvaluatorFactory also follows that driver-side cached RDD. 
> This unnecessarily captures the cached RDD object graph and can fail with 
> NotSerializableException when it contains driver-only state. Executable 
> inputs are already supplied to the evaluator separately.
> A bounded reproducer builds a SortMergeJoinExec with a ColumnarToRowExec 
> child whose cached RDD has a nonserializable sentinel, obtains actual 
> whole-stage codegen references, initializes the columnar cache, and 
> serializes the evaluator with Spark's Java closure serializer. Upstream fails 
> through:
> WholeStageCodegenEvaluatorFactory.references -> SortMergeJoinExec -> 
> ColumnarToRowExec -> SparkPlan.executeColumnarRDD -> LazyTry.tryT -> cached 
> RDD -> nonserializable sentinel
> Marking executeColumnarRDD @transient fixes this capture path and matches 
> executeRDD. Driver-side cache reuse is unchanged; serialization, 
> deserialization, and cleanup succeed with the fix.
> The issue was reproduced using Apache Spark branch-4.2 source at 
> fc01a9e3165f9c371074390f4b3e6f7aeabef04b against public Spark 4.2.0 
> dependencies. The same nontransient field exists on current master and 
> branches 4.x, 4.2, 4.1, and 4.0. Branch 3.5 directly invokes 
> doExecuteColumnar and has no such cache field.
> This report concerns the demonstrated cached-RDD capture path. It does not 
> claim to eliminate every source of large task closures.
> Proposed fix and regression test: https://github.com/apache/spark/pull/59146
> Current-master validation at a026a0be4c88e7ace59d1da7c18ab73b0d81edde: the 
> focused regression passes with the fix and fails with 
> NotSerializableException when only the annotation is removed. The exact 
> changed sources were compiled against public Spark 4.2.0 dependencies; this 
> is not a full current-master build. Scoped Scalastyle passes with zero errors 
> and warnings.



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