Conversation
dongjoon-hyun
left a comment
There was a problem hiding this comment.
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.shuffleRDDis the very object returned bydoExecuteColumnar(), 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
InMemoryTableScanExecbehind a shuffle).
P2. Consistency, root cause, and backporting
executeBroadcastBcastis now the only non-transientLazyTrycache inSparkPlan.- The root cause is
ctx.addReferenceObj("plan", this)inSortMergeJoinExec,SortExec,HashAggregateExec, andShuffledHashJoinExec, 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.0throughbranch-4.3, and please update the "Draft pending full source-build validation" note.
P3. Test cleanups (nits)
SparkPlanSuitelines 61-62 (thereadbackround trip andcleanupResources()) are redundant and cannot fail.- A direct plan round trip is enough; the SMJ/WSCG scaffolding couples the test to SMJ codegen internals.
- The
durationMsargument receivesnumOutputRows. EmptyRDDcan 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() |
There was a problem hiding this comment.
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:
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.
There was a problem hiding this comment.
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.
| @@ -218,6 +218,7 @@ abstract class SparkPlan extends QueryPlan[SparkPlan] with Logging with Serializ | |||
| executeBroadcastBcast.get.asInstanceOf[broadcast.Broadcast[T]] | |||
There was a problem hiding this comment.
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?
There was a problem hiding this comment.
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)) |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
viirya
left a comment
There was a problem hiding this comment.
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.
dongjoon-hyun
left a comment
There was a problem hiding this comment.
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
SparkPlanSuiteL75: 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.SparkPlanSuiteL97:assert(serializedStageRDD)always holds, and nothing checks that the downstream stage captures the upstream plan.
Nits
SparkPlanSuiteL64:cached.count()is redundant, andcheckAnswerruns the aggregation twice.SparkPlanL223: the Scaladoc ofexecuteColumnar()refers to a nonexistentdoColumnarExecute. This is not from this PR.SparkPlanSuiteL80:OutputStream.nullOutputStream()can replace the unreadByteArrayOutputStream.
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() |
There was a problem hiding this comment.
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) HashAggregatereads theInMemoryTableScanthrough anInputAdapter. - After the query runs,
executeColumnarRDDis not materialized in any node. - Before this line, the serialized
(stageRDD, taskFunc)contains onlyMapPartitionsRDDandShuffledRowRDD, i.e., no upstream RDD. After this line, it contains the columnar RDD and its whole lineage down toParallelCollectionRDD.
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) |
There was a problem hiding this comment.
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() |
There was a problem hiding this comment.
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() |
There was a problem hiding this comment.
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()) { |
There was a problem hiding this comment.
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()) {
What changes were proposed in this pull request?
Mark
SparkPlan.executeColumnarRDDtransient, matching the existing row execution cache.Add two SPARK-59884 regressions:
DAGSchedulerdoes, and checks that the upstreamInMemoryTableScanExeccolumnar 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.executeColumnarRDDcaches 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
SparkPlancolumnar cache. Scan-ownedinputRDD/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?
SparkPlan.scalaand revisedSparkPlanSuite.scalaagainst 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 completeSparkPlanSuitepassed: 14 tests, 0 failures.executeColumnarRDDannotation 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.git diff --checkpassed.d9fe008b13f: 28 successful checks and 2 skipped checks.Was this patch authored or co-authored using generative AI tooling?
Generated-by: Codex (GPT-6).