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