Skip to content

[SPARK-59884][SQL] Mark SparkPlan's columnar execution cache transient - #59146

Open
sunchao wants to merge 3 commits into
apache:masterfrom
sunchao:codex/columnar-rdd-transient-upstream
Open

sunchao wants to merge 3 commits into
apache:masterfrom
sunchao:codex/columnar-rdd-transient-upstream

Conversation

@sunchao

@sunchao sunchao commented Sep 30, 2026 •

Copy link
Copy Markdown
Member

What changes were proposed in this pull request?

Mark SparkPlan.executeColumnarRDD transient, matching the existing row execution cache.

Add two SPARK-59884 regressions:

  • A direct plan serialization round trip with an initialized columnar cache containing a nonserializable sentinel. The test also checks driver-side cache reuse.
  • A real cached-DataFrame aggregation across a shuffle. After checking the query results, the test serializes the complete downstream result-stage RDD/function tuple, as DAGScheduler does, and checks that the upstream InMemoryTableScanExec columnar RDD is absent from that object graph.

JIRA: SPARK-59884

Why are the changes needed?

Plans retained by downstream operators can reach upstream plans across a shuffle boundary. Their initialized SparkPlan.executeColumnarRDD caches then retain RDDs that are no longer part of the downstream stage's input lineage. Excluding this driver-side cache removes that serialization path while preserving driver-side reuse.

The cache was introduced by SPARK-48195 / #48037 in d2e8c1c, first included in Spark 4.0.0. The same nontransient field remains in branches 4.0 through 4.3, so those branches are candidates for backporting.

This change is limited to the base SparkPlan columnar cache. Scan-owned inputRDD/maybeCoalesceInputRDD, AQEShuffleReadExec.shuffleRDD, the broadcast execution cache, and whole-plan capture remain separate follow-up work. It does not remove every route by which a plan can retain an RDD.

Does this PR introduce any user-facing change?

Yes. It avoids serializing an upstream columnar execution RDD through the base plan cache when that RDD is not needed by a downstream stage. Query results and public APIs are unchanged. It does not make a nonserializable RDD usable as a stage's actual input.

How was this patch tested?

  • Compiled the current SparkPlan.scala and revised SparkPlanSuite.scala against the Spark 5.0.0-SNAPSHOT artifacts built from this PR's unchanged production sources in CI run 36649672228, using Scala 2.13.18 and ScalaTest 3.2.19. The complete SparkPlanSuite passed: 14 tests, 0 failures.
  • Recompiled with only the executeColumnarRDD annotation removed. Both new regressions failed as expected: the direct round trip reached the nonserializable sentinel, and the cross-shuffle test detected the upstream columnar RDD in the downstream object graph.
  • Scoped Scalastyle passed for both modified files with 0 errors and 0 warnings. git diff --check passed.
  • Full fork CI passed on d9fe008b13f: 28 successful checks and 2 skipped checks.
  • The latest follow-up only clarifies the shuffle-lineage comment. Executable code is unchanged; whitespace, ASCII, and line-length checks passed. The Spark suite was not rerun locally for this comment-only edit; CI will rerun on the new commit.

Was this patch authored or co-authored using generative AI tooling?

Generated-by: Codex (GPT-6).

@sunchao sunchao changed the title [SQL][WIP] Avoid serializing cached columnar RDDs in whole-stage closures [SPARK-59884][SQL] Avoid serializing cached columnar RDDs in whole-stage closures Sep 30, 2026
@sunchao
sunchao marked this pull request as ready for review September 30, 2026 00:43
@sunchao

sunchao commented Sep 30, 2026

Copy link
Copy Markdown
Member Author

cc @viirya @cloud-fan @peter-toth @dongjoon-hyun

@dongjoon-hyun dongjoon-hyun left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thank you for making a PR, @sunchao.

The one-line change itself looks correct and safe to me. It is consistent with executeRDD, and a deserialized plan cannot reach executeColumnarRDD.get anyway because executeQuery evaluates sparkContext from the @transient session first. I left some comments. Here is the summary by priority.

P1. The fix is narrower than the PR title and description suggest

  • Spark's built-in columnar scans still hold the same RDD chain in non-transient fields (FileSourceScanExec.inputRDD/maybeCoalesceInputRDD, BatchScanExec.inputRDD, and the streaming DSv2 scans), so this PR only removes the thin wrapper RDD for them.
  • AQEShuffleReadExec.shuffleRDD is the very object returned by doExecuteColumnar(), so this change has no effect there.
  • The new test covers a same-stage plan shape where the cached RDD is part of the task lineage anyway. It doesn't cover the cross-stage capture where this fix actually matters (e.g., a cached InMemoryTableScanExec behind a shuffle).

P2. Consistency, root cause, and backporting

  • executeBroadcastBcast is now the only non-transient LazyTry cache in SparkPlan.
  • The root cause is ctx.addReferenceObj("plan", this) in SortMergeJoinExec, SortExec, HashAggregateExec, and ShuffledHashJoinExec, which serializes whole subtrees, including upstream stages, into every task. It may deserve a follow-up JIRA.
  • This is a regression from SPARK-48195 in Apache Spark 4.0.0. Please mention it in the PR description so that we can consider backporting to branch-4.0 through branch-4.3, and please update the "Draft pending full source-build validation" note.

P3. Test cleanups (nits)

  • SparkPlanSuite lines 61-62 (the readback round trip and cleanupResources()) are redundant and cannot fail.
  • A direct plan round trip is enough; the SMJ/WSCG scaffolding couples the test to SMJ codegen internals.
  • The durationMs argument receives numOutputRows.
  • EmptyRDD can replace the anonymous RDD.
  • The empty-parameter case class cannot be canonicalized.
  • Please add the SPARK-59884: prefix to the test name.


@transient
private val executeColumnarRDD = LazyTry {
doExecuteColumnar()

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

AQEShuffleReadExec is a case where this change has no effect at all. Its doExecuteColumnar() (and doExecute()) returns shuffleRDD, which is a non-transient private lazy val:

private lazy val shuffleRDD: RDD[_] = {
shuffleStage match {
case Some(stage) =>
sendDriverMetrics()
stage.shuffle.getShuffleRDD(partitionSpecs.toArray)
case _ =>
throw SparkException.internalError("operating on canonicalized plan")
}
}
override protected def doExecute(): RDD[InternalRow] = {
shuffleRDD.asInstanceOf[RDD[InternalRow]]
}
override protected def doExecuteColumnar(): RDD[ColumnarBatch] = {
shuffleRDD.asInstanceOf[RDD[ColumnarBatch]]
}

So the very object cached in executeColumnarRDD is still serialized through shuffleRDD when a downstream stage's SortExec, HashAggregateExec, SortMergeJoinExec, or ShuffledHashJoinExec captures a subtree that reaches an earlier stage's AQEShuffleReadExec through the non-transient ShuffleQueryStageExec.plan. For the columnar path this only happens with plugin ShuffleExchangeLike implementations, but the row path has bypassed the @transient executeRDD in the same way since SPARK-48195. @transient private lazy val shuffleRDD (like partitionDataSizes and metrics in the same class) would close it.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agreed. AQEShuffleReadExec.shuffleRDD remains an independent reference to the same RDD, so this annotation does not remove that retention path. I documented this limitation and kept that cache change for separate follow-up coverage.

Comment thread sql/core/src/test/scala/org/apache/spark/sql/execution/SparkPlanSuite.scala Outdated
@@ -218,6 +218,7 @@ abstract class SparkPlan extends QueryPlan[SparkPlan] with Logging with Serializ
executeBroadcastBcast.get.asInstanceOf[broadcast.Broadcast[T]]

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

SPARK-48195 introduced three LazyTry caches, and with this PR executeBroadcastBcast becomes the only non-transient one. So every plan node serialized through codegen plan references still carries an unused LazyTry with its initializer lambda, and materialized broadcast nodes carry Success(TorrentBroadcast). This is harmless in practice because only the broadcast handle is serialized. However, a deserialized plan can never reach it either, since executeQuery dereferences the @transient session first. Adding @transient to executeBroadcastBcast would make the three caches consistent. WDYT?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The consistency change makes sense as a follow-up. I kept this PR focused on the columnar RDD retention regression and its distinguishing tests; executeBroadcastBcast remains unchanged. The description now makes that scope explicit.

Seq(Literal(1)), Seq(Literal(1)), Inner, None,
ColumnarToRowExec(source), LocalTableScanExec(Nil, Nil, None))
val (ctx, code) = WholeStageCodegenExec(join)(0).doCodeGen()
assert(ctx.references.exists(_.asInstanceOf[AnyRef] eq join))

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This assertion points at the root cause, which is one level above the patched field. SortMergeJoinExec, SortExec, HashAggregateExec, and ShuffledHashJoinExec put this into the codegen references (ctx.addReferenceObj("plan", this)). As a result, the whole subtree is serialized into each task, and it crosses stage boundaries through the non-transient QueryStageExec.plan, QueryStageExec._canonicalized, and ShuffleExchangeExec.child. Downstream tasks therefore also deserialize all upstream stage plans and register their SQLMetrics. The non-codegen paths (the SortExec.doExecute closure and SortMergeJoinEvaluatorFactory(left, right)) capture the subtree too.

A field-level @transient fixes one symptom. It may be worth a follow-up JIRA to (1) make all execution caches in SparkPlan and its subclasses transient consistently, and (2) avoid serializing plans across stage boundaries, e.g., by keeping only output.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agreed that reducing whole-plan capture deserves follow-up work. This PR now explicitly limits itself to the base columnar cache. A broader change will also need to preserve the operator methods and cleanup behavior used by the captured plans, so I kept it separate from this fix.

Comment thread sql/core/src/test/scala/org/apache/spark/sql/execution/SparkPlanSuite.scala Outdated
Comment thread sql/core/src/test/scala/org/apache/spark/sql/execution/SparkPlanSuite.scala Outdated
Comment thread sql/core/src/test/scala/org/apache/spark/sql/execution/SparkPlanSuite.scala Outdated
Comment thread sql/core/src/test/scala/org/apache/spark/sql/execution/SparkPlanSuite.scala Outdated
Comment thread sql/core/src/test/scala/org/apache/spark/sql/execution/SparkPlanSuite.scala Outdated
@sunchao sunchao changed the title [SPARK-59884][SQL] Avoid serializing cached columnar RDDs in whole-stage closures [SPARK-59884][SQL] Mark SparkPlan's columnar execution cache transient Sep 30, 2026

@viirya viirya left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the update. The revised tests address the main coverage gap identified in the earlier review: serializing the downstream RDD/function tuple checks whether the upstream columnar RDD is retained across a shuffle boundary, which the original factory-only round trip could not establish. The direct plan serialization test separately checks the cache exclusion and driver-side reuse.

The production change looks correct to me and is consistent with executeRDD. The narrowed title and description also accurately distinguish this fix from the remaining scan-owned caches and whole-plan capture. Keeping those broader changes separate seems reasonable, particularly since reducing plan capture needs to preserve executor-side operator methods and cleanup behavior.

I have no blocking concerns with this revision. One minor clarification in the test comment below.

Comment thread sql/core/src/test/scala/org/apache/spark/sql/execution/SparkPlanSuite.scala Outdated

@dongjoon-hyun dongjoon-hyun left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thank you for updating the PR, @sunchao. I left 5 inline comments on the latest commit. The production change still looks good to me. Four comments are about the new cross-shuffle test, and one is a Scaladoc nit. Here is a summary, ordered by priority.

Test validity

  1. SparkPlanSuite L75: the cached DataFrame is read row-wise (SPARK-37369), so the job never creates the scan's columnar RDD. Only this line creates it, after the job, and that is why the test fails without the fix.
  2. SparkPlanSuite L97: assert(serializedStageRDD) always holds, and nothing checks that the downstream stage captures the upstream plan.

Nits

  1. SparkPlanSuite L64: cached.count() is redundant, and checkAnswer runs the aggregation twice.
  2. SparkPlan L223: the Scaladoc of executeColumnar() refers to a nonexistent doColumnarExecute. This is not from this PR.
  3. SparkPlanSuite L80: OutputStream.nullOutputStream() can replace the unread ByteArrayOutputStream.

Item 1 is the most important one, since the PR description presents this test as a real cached-DataFrame aggregation across a shuffle.

case scan: InMemoryTableScanExec => scan
}.get
assert(scan.supportsColumnar)
val columnarRDD = scan.executeColumnar()

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The cached InMemoryTableScanExec shape that I suggested in the previous round turns out not to exercise the columnar path. Since SPARK-37369, InMemoryTableScanExec.supportsRowBased is always true, so no ColumnarToRowExec is inserted, and the query reads the cache row-wise through InputAdapter.inputRDD (child.execute()).

So the job never creates the scan's columnar RDD. This line creates it after checkAnswer, and line 98 fails without the fix only because of that. I ran the same query on Apache Spark 4.3.0, which doesn't have this fix and has the same code paths as master:

  • The executed plan has no ColumnarToRow. *(1) HashAggregate reads the InMemoryTableScan through an InputAdapter.
  • After the query runs, executeColumnarRDD is not materialized in any node.
  • Before this line, the serialized (stageRDD, taskFunc) contains only MapPartitionsRDD and ShuffledRowRDD, i.e., no upstream RDD. After this line, it contains the columnar RDD and its whole lineage down to ParallelCollectionRDD.

So the test doesn't reproduce a leak of this query, although the PR description presents it as "a real cached-DataFrame aggregation across a shuffle". assert(scan.supportsColumnar) doesn't prove that the columnar path is used either. Could you assert that the job itself created columnarRDD, like the "or the test exercises nothing" guards in DataFrameSetOperationsSuite?

        def lineage(rdd: RDD[_]): Seq[RDD[_]] = rdd +: rdd.dependencies.flatMap(d => lineage(d.rdd))
        assert(lineage(exchanges.head.inputRDD).exists(_ eq columnarRDD),
          "the job must create the columnar RDD, or the test exercises nothing")

This assertion fails with the cached DataFrame. A vectorized Parquet scan passes it because ColumnarToRowExec.inputRDDs() materializes FileSourceScanExec.executeColumnarRDD during the job. For Parquet, note that FileSourceScanExec.inputRDD is still serialized after this fix, as discussed in the previous round, so the identity check covers only the wrapper RDD. Otherwise, please update the test comments and the PR description to say that this line stands in for a columnar consumer, e.g., a plugin operator.

} finally {
out.close()
}
assert(serializedStageRDD)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

serializedStageRDD is always true. stageRDD is _1 of the tuple passed to writeObject, so replaceObject always sees it while the hook is enabled, and this line only checks that the hook works.

Nothing checks that the downstream stage captures the upstream plan, which is the precondition of line 98. Today the only path is ctx.addReferenceObj("plan", this) of the final HashAggregateExec:

If it goes away, e.g., when whole-plan capture is reduced as discussed in the previous round, line 98 passes even without this fix. I simulated it on Apache Spark 4.3.0 by writing the final HashAggregateExec as null in the same stream. Neither scan nor columnarRDD was serialized, so the assertions of this test would pass without the fix. How about tracking the scan node instead of stageRDD?

        var serializedScan = false
        ...
            serializedScan ||= obj eq scan
        ...
        assert(serializedScan, "the downstream stage no longer captures the upstream plan")

SQLConf.SHUFFLE_PARTITIONS.key -> "2") {
val cached = spark.range(0, 8, 1, 2).selectExpr("id % 2 AS key").cache()
try {
cached.count()

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: cached.count() isn't needed because cache() already registers the relation and the query builds the buffers on its first run. In addition, checkAnswer runs the aggregation twice: materializedRdd.count() on a separate rddQueryExecution plan, and then collect() on this plan. Only the latter matters for the checks below, so QueryTest.checkAnswer(query, Seq(Row(0L, 4L), Row(1L, 4L)), checkToRDD = false) would be enough.


@transient
private val executeColumnarRDD = LazyTry {
doExecuteColumnar()

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: This is not from this PR, but the Scaladoc of executeColumnar() right below refers to doColumnarExecute twice (lines 227 and 230), which doesn't exist. It should be this doExecuteColumnar.


var serializedStageRDD = false
var serializedColumnarRDD = false
val out = new ObjectOutputStream(new ByteArrayOutputStream()) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: The serialized bytes are never read. OutputStream.nullOutputStream() avoids buffering the whole task binary and makes the intent clearer, like the NullOutputStream of SerializationDebugger.ListObjectOutputStream.

import java.io.{ObjectOutputStream, OutputStream}
...
        val out = new ObjectOutputStream(OutputStream.nullOutputStream()) {

HyukjinKwon

This comment was marked as outdated.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants