Repository navigation
fix: draw the cached plan below CometInMemoryTableScan in the SQL tab and event log - #6577
Conversation
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Comet cache scans omitted the cached plan from EXPLAIN, the SQL graph and event-log plans.
- Design approach: Expose the relation through
innerChildrenand the cached physical plan throughsubqueries, while excluding both from Comet’s execution reporting. - Correctness / compatibility analysis: Checked Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. The graph construction works, but the generic
subqueriesoverride introduces the Spark 4.2 metric regression below. - Key design decisions: Executable children remain unchanged, and the lazy override accommodates Spark 3.x. The implementation is small, but using an execution-related traversal API for display data has observable consequences beyond the UI. Additional traversal is identifiable, but no performance regression was measured.
- Implementation sketch: Two scan overrides, one reporting exclusion, a test helper and AQE-on/off regression tests.
- Behavioral changes worth calling out: Compared with
branch-1.1, displaying the cached subtree is intended. Codegen inspection also reaches that subtree, and Spark 3.x can emit another final plan update. The last-attempt metric failure is unintended. - Suggested improvements: Address the P2 finding and cover a metric recorded outside an already materialized cache containing a Comet shuffle.
Reviewed the entire four-file diff from f980d6fb59f51ae20c0a6916eec9771ace8e50d1 to a422a549d327933010fc04c59bc29ec68bef2863, including prerequisite commit 81736558d43e00f55499192761c6cd9303c31444. Verified its reusable source evidence. The PR remains non-draft. Snapshot and live discussion checks contained no existing review concerns. Routed skill: .ai/skills/review-comet-pr/SKILL.md. No sibling area skill applies.
Exact-head CI at 2026-10-03 19:19 UTC: 38 successful, 7 running and 35 skipped checks, with no failures. Spark 4.1 and 3.4 execution jobs passed. The Spark 4.1 log confirms both new tests passed within 1,229 successful tests. Execution jobs for 3.5, 4.0 and 4.2 remained pending. Upstream Spark SQL suites were skipped.
Validation: git diff --check passed. A Spark 3.5.9 API probe passed eight cases covering AQE, nested caches and materialization. A bounded replay of Spark 4.2’s scope-extraction method confirmed the finding. These probes did not execute Comet native code. The Spark 4.1 probe compiled but could not start because the available runtime lacked KVStore. No local full Comet build ran. Build artifacts and dependencies were absent, and Maven Central access was blocked.
Review state: Request changes.
| // Spark only walks this list: subqueries run from a plan's expressions. Its other walkers, such | ||
| // as collectWithSubqueries, follow it into the cached plan too. A lazy val, because Spark 3.x | ||
| // declares subqueries as one, and a lazy val overrides Spark 4's def as well. | ||
| @transient override lazy val subqueries: Seq[SparkPlan] = Seq(originalPlan.relation.cachedPlan) |
There was a problem hiding this comment.
[P2] Keep cached shuffles out of last-attempt metric traversal. On Spark 4.2 with Comet’s cache enabled, materialize spark.range(0, 100, 1, 2).repartition(2).cache(), then increment a SQLLastAttemptMetrics.createMetric accumulator in a subsequent .map over that cache. The consuming query has no shuffle, and the metric is entirely outside the cache, so lastAttemptValueForDataset should return Some(100). This override makes SQLLastAttemptAccumulator.extractStageRDDScopes enter the cached plan, encounter CometShuffleExchangeExec, and return Left(Unsupported ShuffleExchangeLike: ...). The metric accessor consequently returns None. Previously the cached shuffle was outside this traversal. Spark’s documented undefined behavior applies when the metric itself was used inside the cached plan, which does not cover this case. Please isolate display-only cached plans from this walker, or provide compatible handling, and add a regression for an outside-cache metric.
Evidence: A bounded JVM probe replayed Spark v4.2.0’s extractStageRDDScopes method byte-for-byte using Spark 4.1.3 plan classes, an equivalent helper companion and stable substitute scope IDs. Its cached subtree contained a plugin exchange implementing ShuffleExchangeLike, matching Comet’s inheritance, beneath a consuming stage. Changing only whether the cache exposed that subtree through subqueries changed the result from Right(List(3)) to Left(Unsupported ShuffleExchangeLike: org.apache.spark.sql.execution.metric.PluginShuffle). Probe: /tmp/pr6577-review-a422a549/ScopeRegressionProbe.scala, output: /tmp/pr6577-review-a422a549/scope-probe.log. Spark v4.2.0 SQLLastAttemptAccumulator.scala lines 343–350 reject non-Spark shuffle implementations, lines 424–426 traverse these subqueries, and lines 257–263 convert the failure to None. This validates the traversal regression, not an end-to-end Comet execution.
There was a problem hiding this comment.
Thanks, confirmed and fixed in c190f6e (rebased onto main after #6574 merged). subqueries now returns Spark's own InMemoryTableScanExec (originalPlan) rather than the cached plan. SparkPlanInfo special-cases that class, so the SQL tab and event log still draw the cached plan, one level down, while SQLLastAttemptAccumulator and the other subquery walkers stop at that scan, as they do in Spark's own plans.
Your outside-cache scenario is now a 4.2-only CometInMemoryCacheLastAttemptMetricSuite: it checks that the cached plan has a Comet shuffle, then expects Some(100). A check on every version also asserts that collectWithSubqueries no longer reaches the cached plan's shuffle. The 4.2 suite could not run locally, so the run-all-spark-profiles run gives its first result.
a422a54 to
5d5b520
Compare
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Comet cache scans omitted the cached subtree from EXPLAIN, the SQL graph and event-log plans.
- Design approach: Expose the relation through
innerChildrenand its cached physical plan throughsubqueries, while excluding it from Comet’s coverage and fallback reporting. - Correctness / compatibility analysis: Checked Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. The display behavior follows Spark’s implementation. The existing Spark 4.2 metric regression remains unresolved: a metric recorded outside a materialized cache can return
Noneinstead ofSome(100)because traversal now encounters a cached Comet shuffle. - Key design decisions: Executable children remain unchanged, and the lazy override supports both Spark 3.x and 4.x. The implementation is small, but using
subqueriesfor display couples it to other plan walkers. No additional reproducible performance regression was identified. - Implementation sketch: Two scan overrides, one reporting exclusion, a test helper and AQE-on/off regression tests.
- Behavioral changes worth calling out: Compared with
branch-1.1, displaying the cached subtree is intended. Codegen inspection also reaches it, and Spark 3.x can emit an additional final plan update. The last-attempt metric regression is unintended. - Suggested improvements: Address the existing metric-traversal concern by isolating the display subtree or providing compatible handling, with a regression covering a metric outside an already materialized cache.
Reviewed full SHA 5d5b52045d069dae6e23be1ab2106f572f67680d, including all four files and both prerequisite commits from #6574. The requested base is 569eaa59d032f758669777964aee2d24eb55ebae; the PR merge-base is f980d6fb59f51ae20c0a6916eec9771ace8e50d1. Read AGENTS.md and routed through .ai/skills/review-comet-pr/SKILL.md. No sibling area skill applies. The PR remains non-draft.
Read the snapshot and live discussions. No additional introduced P1/P2 issues found within this review. The existing P2 concern remains at spark/src/main/scala/org/apache/spark/sql/comet/CometInMemoryTableScanExec.scala:95. Production behavior is unchanged from the previously reviewed head, so this is not duplicated as a new finding.
Validation: git diff --check passed. Reran the Spark 3.5.9 API probe successfully across eight AQE, nested-cache and materialization cases. Verified the Spark 4.2 scope-extraction replay byte-for-byte against upstream source and reproduced the change from Right(List(3)) to Left(Unsupported ShuffleExchangeLike: ...). These probes validate Spark traversal behavior, not end-to-end Comet execution. No local Comet/native build or suite ran because this checkout lacks built artifacts and Maven Spark dependencies.
Exact-head CI at 2026-10-03 20:45 UTC: 6 successful, 2 running and 14 skipped checks, with no failures reported. Linux lint remained running, and no exact-head Comet test verdict was available. run-all-spark-profiles is applied. Upstream Spark SQL suites were skipped.
Review state: Request changes remains warranted for the existing P2 concern.
… and event log Spark builds the SQL tab's graph and the plans in the event log with SparkPlanInfo.fromSparkPlan. It gives its own InMemoryTableScanExec the cached plan as a child, but recognizes that scan by its class, and for any other node takes the children and the subqueries. Expose the cached plan as the one subquery of CometInMemoryTableScanExec. Spark runs subqueries from a plan's expressions and only walks this list, so the cached plan does not become part of the query that reads the cache. Closes apache#6463.
Spark 4.2's SQLLastAttemptAccumulator walks a plan's subqueries to find a metric's stages and gives up on a shuffle it does not know. With the cached plan as the subquery it reached the cached plan's Comet shuffle, so a last-attempt metric used outside the cache returned None. SparkPlanInfo draws the cached plan below Spark's own InMemoryTableScanExec, matched by class, so expose that scan instead. Every other walker of subqueries stops at it, as in Spark's own plans. Adds a check on every version that collectWithSubqueries does not reach the cached plan's shuffle, and a Spark 4.2 suite for the outside-cache metric.
1a36ec9 to
c190f6e
Compare
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Comet cache scans omitted the cached plan and its metrics from the SQL graph and event-log plans.
- Design approach: Expose the original Spark cache scan through
subqueries. Spark’s existingSparkPlanInfohandling then displays the cached plan underneath it. - Correctness / compatibility analysis: Checked Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. Preparation and transformations follow expressions, so the added scan does not execute. The current override addresses the existing Spark 4.2 metric concern by keeping the cached Comet shuffle outside scope extraction.
- Key design decisions: A single lazy override supports Spark 3.x and 4.x without another production abstraction. The existing
innerChildrenoverride preserves EXPLAIN and Comet reporting. Graph construction and observed-metric collection gain traversal work, but no reproducible P1/P2 performance regression was identified. - Implementation sketch: One production override, a test helper, AQE-on/off graph and traversal assertions, and a Spark 4.2 metric regression suite registered in both workflows.
- Behavioral changes worth calling out: Compared with
branch-1.1at992c806a7e38c2e88bd018aa5774164b0850e1fa, restoring the cached subtree in the SQL graph and event log is intended. The EXPLAIN and observed-metric fallback changes are already in the supplied base. Spark 3.x can emit an additional final plan update. Observed-metric collection now reaches Spark’s scan, while the existing fallback for observed caches remains intact. - Suggested improvements: No further P1/P2 changes requested.
Reviewed all six files and all three commits in the full diff from d98fd2494c9a6cc8552efe6aa0a73d5714486b23 to c190f6e3ac0252987d3d77dc7bdf15313ce7147c. The PR remains non-draft. Read AGENTS.md, the snapshot and live discussions. Routed skill: .ai/skills/review-comet-pr/SKILL.md. No sibling skill applies.
No introduced P1/P2 issues found within this review. The existing P2 concern is addressed by the current override.
Validation: git diff --check passed. Spark API probes passed eight cases each on 3.5.9 and 4.1.3, covering AQE, nested caches and materialization. A byte-for-byte replay of Spark 4.2’s scope-extraction method on Spark 4.1 classes reproduced the old shuffle rejection and succeeded with the current override. The replay substitutes stable scope IDs and does not validate end-to-end accumulator execution. No full Comet build or native suite ran locally because this checkout lacks built artifacts and Maven dependencies.
Exact-head CI at 2026-10-04 16:52 UTC: 8 successful, 11 running and 14 skipped checks, with no failures. Comet execution results, including the new Spark 4.2 suite, remain pending. run-all-spark-profiles is applied. Upstream Spark SQL suites and macOS checks were skipped. This review does not establish a completed CI verdict.
Since apache#6577, CometInMemoryTableScanExec exposes Spark's own InMemoryTableScanExec as its one subquery, so that Spark's UI draws the cached plan below it. With the cache on by default, the DPP suite's 'filtering ratio policy fallback' caches its dimension table, and checkPartitionPruningPredicate requires every subquery of an adaptive plan to contain an AdaptiveSparkPlanExec, which Spark's scan does not. Skip it there in every diff, since the query never runs it as a subquery.
The diffs skipped CachedTableSuite's 'SPARK-35332: Make cache plan disable configs configurable - check AQE' on Spark 4.0 and later, because SparkPlanInfo drew the cached plan only below Spark's own InMemoryTableScanExec. Since apache#6577 it is drawn below CometInMemoryTableScanExec too, so the test finds the cached plan. The test's check that the cached plan was coalesced, which the diffs already adapt, looked for AQEShuffleRead only below a columnar-to-row transition, where it is when Comet's plan is cached in Spark's format. In Comet's columnar format the cached plan ends in the read itself, as in Spark, so the check accepts it there too. The test passes on Spark 4.0, 4.1 and 4.2 with either format, and with Comet off.
* perf: project cached batches by buffer selection, prune on collated strings Follow-up to apache#5051, applying items from apache#5487. Replace the per-column Arrow IPC stream layout of `CometCachedBatch` with a single encapsulated IPC record batch message per cached batch, carrying no Schema message and no end-of-stream marker. The reader rebuilds the schema from the cached relation's attributes, so a wide relation no longer repeats the same schema bytes once per cached batch. Compression moves from a whole-payload Spark codec to Arrow's per-buffer IPC compression. That is what makes projection cheap: the message metadata records every buffer's offset and length in the body, so `CachedBatchIpc.readProjected` copies out only the byte ranges of the columns a scan selected and decompresses just those. This subsumes the separate "drop the schema message" item, since there is no longer a per-column stream to frame. Dictionary-encoded columns are decoded before being stored: a payload with no schema message cannot describe a dictionary encoding. The codec defaults to zstd, and lz4 is deliberately not offered. Arrow's lz4 is commons-compress's pure-Java implementation, unrelated to the JNI-accelerated lz4-java behind `spark.io.compression.codec`. Over a 200k-row six-column relation it measured 205s to write against 347ms for zstd, while also producing larger output, so no workload prefers it. zstd also beats storing batches uncompressed on both axes (347ms and 2 MiB against 1743ms and 13 MiB), because the bytes it saves cost more to copy and store than compressing them costs. Decompression is done here rather than left to `VectorLoader`, which leaks: `VectorLoader.loadBuffers` collects a field's decompressed buffers into a local list and releases them only after the whole field loads, so a buffer that fails to decompress strands every buffer of that field decompressed before it. A string column reaches this, its offsets buffer decompressing before its data buffer throws. Also track statistics bounds for collated string columns, comparing with the collation's own ordering through a new `CometTypeShim.compareStrings`. Matching the bare `StringType` object excluded collated columns, which then got null bounds and no pruning. Benchmark over a 5M-row six-column relation, keeping the cached scan native against falling back to a Spark cache scan and converting: 1.3x on a repeated scan, 1.3x on a narrow projection and 2.3x on a full projection. * refactor: use Spark's interpreted ordering for bounds, hoist projection layout Cleanup pass over the cache format change. No behaviour change. Drop the `compareStrings` shim in favour of `TypeUtils.getInterpretedOrdering`. That method is public with the same signature on every supported Spark version, and on Spark 4 it resolves a `StringType` through `CollationFactory.fetchCollation(collationId).comparator` -- the comparison the shim was reaching for. So the collation awareness comes from Spark itself and the shim, its Spark 3.x stub and the hand-rolled per-type `compare` all go. The ordering is now resolved once per column per partition rather than being re-dispatched on the `DataType` twice per row. Build the projection's index layout once per partition instead of per batch. The node, buffer and variadic index arithmetic is a pure function of the cached schema and the selected columns, but it walks every field of the relation, so recomputing it per batch made the bookkeeping O(total columns) against O(selected columns) of useful work -- worst in the wide-relation, narrow-projection case the format exists for. `CachedBatchIpc.Projection` now holds that layout and the projected schema, and owns the whole decode; `ProjectedBatch` is left with ownership only. This also puts the projected schema next to the code that packs buffers in the same order, an invariant that previously spanned two files unstated. Smaller cleanups: use Arrow's `DataSizeRoundingUtil.roundUpTo8Multiple` rather than open-coding IPC body alignment; size the serialization buffer from the record batch's known body length instead of growing from 32 bytes; resolve decompressors once instead of per batch; share the dictionary lookup guard between `Utils.combineDictionaryProviders` and the cache writer; read the codec config through one helper carrying the driver-vs-executor rationale; and collapse the duplicated compressed-buffer predicate and scramble loop in the test helper. Corrects two `Utils` scaladocs that still described the per-column stream format this change replaced. Benchmark and codec figures in the docs re-measured against the current code. * fix: relocate the arrow-compression service file when shading arrow-compression ships META-INF/services/org.apache.arrow.vector.compression.CompressionCodec$Factory. The shade plugin copies it verbatim without a ServicesResourceTransformer, so the jar declared a provider for Spark's own unshaded Arrow interface while naming a class that exists here only under the relocated package. Every ServiceLoader lookup Spark's Arrow made then failed with a ServiceConfigurationError, which took CompressionCodec.Factory's static initializer down with it and broke unrelated Arrow IPC reads, including mapInArrow. Add ServicesResourceTransformer so the service file name and its contents are both relocated. arrow-compression is the only bundled artifact that ships one. Also drop an unused NonFatal import that scalafix flagged. * test: drop the cache leak test that depends on zstd corruption detection "releases its vectors when a column fails part way through" zeroed the last 16 bytes of a compressed buffer and required the read to fail. Whether that fails is a property of the zstd runtime, not of Comet: the cached payload is byte-identical across Spark versions, but Comet takes zstd-jni from Spark rather than from arrow-compression, and 1.5.5 (Spark 3.4, 3.5) decodes that frame while 1.5.7 (Spark 4.x) reports it corrupt. So the test passed on 4.x and failed on 3.4 and 3.5. The scenario it claimed to cover is also unreachable: CachedBatchIpc decompresses every selected buffer before VectorLoader runs, so no content corruption can fail part way through the load. The two remaining leak tests corrupt a frame from its header onwards, which every zstd release rejects, and already cover a failure at a column's first buffer and a failure after an earlier buffer of the same column decoded. Records the constraint on scramble so a future test does not reach for a tail-only corruption again, and drops the now unused truncateColumn helper and the dictionary fixture's payload argument. * feat: enable Comet's in-memory cache by default Flip spark.comet.exec.inMemoryCache.enabled to true so cached tables are stored and scanned in Comet's Arrow format without an opt-in. CometDriverPlugin.maybeSetCacheSerializer read the config out of SparkConf with a hardcoded false default, so flipping the ConfigEntry alone would have left the serializer uninstalled unless the user set the key explicitly. It now falls back to the entry's own default, matching how the plugin reads spark.comet.metrics.enabled. Stacked on apache#5543. * test: cover nested columns in the cached-batch projection tests and benchmark Addresses review feedback asking whether nested data should be tested and benchmarked. Nested columns were already round-tripped, but only under a full projection, which cannot see the part of the format that is nontrivial for them. A flat column always owns one field node and two or three buffers; a nested one owns a run as long as its subtree, and selecting every column covers the whole sequence however it is partitioned. So the buffer-span arithmetic was only exercised in the one shape where getting it wrong does not show. Adds two tests over a six-column relation whose middle four columns are a struct, an array, a map and a struct wrapping an array: - Each column takes its turn as the sole projection while the other five are corrupted, so a run computed short or long is caught by reaching into a corrupted neighbour. - Values are compared against the uncached query across single-column, paired and out-of-order projections. Row counts cannot catch a window that is misaligned but still decompresses, and out-of-order is the case a full projection cannot stand in for. The per-column statistics test now runs over the nested relation too, since a nested column's recorded size is the sum of its whole subtree. Both new tests fail if fieldNodeCount stops recursing into children. In the benchmark, adds the three projection widths over a relation of struct columns, and asserts the width each case claims. That assertion caught the existing "full projection (6 of 6 columns)" case reading three: count() over a non-nullable column is rewritten to count(1) by NullPropagation, which prunes the column out of the scan, and only k, s1 and s2 were nullable -- and those only incidentally, because Remainder can divide by zero. Every column of both relations is now nullable so count(c) genuinely reads c, and the documented numbers are regenerated. Array and map columns are left out of the benchmark deliberately: the baseline arm needs Spark's cache scan to bridge into Comet operators, and CometSparkToColumnarExec declines ArrayType and MapType, so for those the arm does not exist and the two cases stop measuring the same boundary. The docs say so rather than leaving it to be rediscovered. * fix: drop a redundant string interpolator flagged by scalafix RedundantSyntax * review: check the cached layout, and address the rest of the review Reader-side: a cached payload carries no schema, so `Projection` derived every node and buffer window from `Utils.toArrowSchema(cacheAttributes)` with nothing checking the writer had produced that layout. `load` now compares `nodesLength()`/`buffersLength()` against the totals `selectedRange` already computes, before any unchecked `batch.buffers(j)`. Writer-side: `isArrowBacked` accepts a `FixedSizeBinaryVector` for a `BinaryType` column, which is two buffers where the reader rebuilds three, and it answers for the top-level vector only -- so a struct of large strings passes it and is stored with 64-bit offsets. `matchesReaderLayout` compares the batch's Arrow types against the reader's recursively, and a batch that disagrees takes the conversion path instead. A dictionary column's field carries the index type, so the dictionary's field is what is compared. Also: an unrecognized body-compression byte is rejected rather than read as plain bytes, `fieldVariadicCount` and the variadic plumbing are gone (the length check covers view vectors, which the counts would not have), `columnSizes` no longer re-walks each column's subtree, the write codec is a case class rather than a bare tuple, the per-partition `Projection` is lazy so a row-count-only read never builds it, `hydrateDictionaries` is `decodeDictionaries`, `Projection` takes an `IndexedSeq`, and the stale `readProjected` links and some over-long comments are fixed. Tests: the two projection tests become one parameterized over both relations, caching once and restoring the payload between columns instead of re-caching; the two leak tests become one with two corruption points. New tests cover the reader's layout check and the writer declining a fixed-size-binary batch. * review: own the write-side compression buffers, fix the activation example Compressing through VectorUnloader leaks on the failure path: appendNodes retains each input buffer and accumulates the compressed ones into a list local to getRecordBatch, so a buffer that fails to compress strands that retain and leaves every buffer compressed before it reachable from nothing. Closing the input batch afterwards undoes neither. Unload plain and compress in CachedBatchIpc.compressed instead, mirroring what decompressed already does on the read side, so every allocation stays reachable from an error path that owns it. The docs enabled the cache with spark.conf.set, which cannot work: the driver plugin picks spark.sql.cache.serializer while the SparkContext is initializing. Show it as a startup --conf. Also drops a redundant s interpolator that the scalafix lint rejected. * fix: write cached batches to the schema width, not the batch width ArrowWriter.writeColumns drove its loop from the input ColumnarBatch's width while indexing the writer's fields, which come from the schema the batch is written under. That assumed every producer hands over a batch exactly as wide as the schema. Iceberg's vectorized reader does not. BatchDeleteFilter.filterBatch reads with the delete filter's requiredSchema, which carries _pos after the projected columns when a data file has position deletes, and trims the extras back only when the file also has equality deletes. A merge-on-read UPDATE writes position deletes and no equality deletes, so the extra column survives into the batch, and caching such a relation failed with ArrayIndexOutOfBoundsException inside the write loop. Drive the loop from the writer's fields instead, which writes exactly the columns the schema describes: the extras are trailing, the same prefix Iceberg keeps when it does trim. A batch narrower than the schema is a genuine contract violation and is now refused with a message naming both widths. Closes apache#6087. * fix: report a cached relation's decoded size to the planner, not its compressed size Comet's cache format reported each batch's sizeInBytes as its stored, compressed payload. Spark's planner reads the sum of those as the size of a materialized cached relation, for the broadcast threshold and the shuffled hash join build side among others, and both of Spark's own cache formats report decoded sizes there. With the default zstd codec a cached relation could look several times smaller than in Spark's format and be broadcast where Spark would shuffle it. Record each column's decoded Arrow size, measured before compression as Spark's ArrowCachedBatchSerializer does, and let CometCachedBatch inherit sizeInBytes from SimpleMetricsCachedBatch as the sum of those. * fix: keep Spark's cache scan for a relation whose cached plan records observed metrics Spark collects Dataset.observe metrics after a query with CollectMetricsExec.collect, which reaches the metrics recorded inside a cached plan only through an InMemoryTableScanExec over it. With Comet's native cache scan in its place those metrics were lost: QueryExecution.observedMetrics came back empty, Observation.get returned an empty map on Spark 3.5 and later, and on Spark 3.4 it never returned. Keep Spark's scan for such a relation, with a fallback reason. The data stays in Comet's format and is read through the existing fallback path. The check walks the cached plan the way CollectMetricsExec.collect does, through subqueries, adaptive plans and nested caches. * test: install Comet's cache serializer in the Spark SQL test sessions The Spark SQL diffs turn Comet on through SharedSparkSession and TestHive rather than by loading CometPlugin, and the plugin is what installs Comet's cache serializer. So even with spark.comet.exec.inMemoryCache.enabled on by default, every Spark SQL suite cached in Spark's own format and none of them exercised Comet's. Set spark.sql.cache.serializer in both places when Comet is enabled, as the plugin would, and adapt the tests that assume Spark's cache scan: - plan checks that only need a cache scan to be present also accept CometInMemoryTableScanExec, reading its originalPlan where the test uses the scan's relation; - PartitionBatchPruningSuite checks its answers as before, but reads the test-only accumulators only from Spark's scan, which Comet's lacks; - CacheTableInKryoSuite registers Comet's classes with CometKryoRegistrator, as Comet asks of any application that sets spark.kryo.registrationRequired; - tests of Spark internals that Comet's scan replaces are tagged IgnoreComet: exact cache size estimates, the union's columnar support, subquery reuse through the scan's predicates, and AQE coalescing under a union (apache#6454). Each diff was regenerated from a clone of its Spark tag with the existing diff applied, after checking that it round-tripped unchanged. * fix: coalesce shuffle partitions under Comet unions the way Spark does Spark's CoalesceShufflePartitions coalesces each child of a UnionExec as its own group, and from Spark 4.0 each child of a CartesianProductExec, BroadcastHashJoinExec or BroadcastNestedLoopJoinExec too. It matches those classes, so a CometUnionExec fell through to the case that coalesces only when every leaf below it is an exchange stage. A union with a scan or a table-cache stage in one branch kept every partition of the shuffles in the others. Add CometCoalesceShufflePartitions, a query-stage optimizer rule that runs after Spark's. For a Comet operator whose Spark original is one of those, and whose shuffle stages no AQE rule has touched, it rebuilds the Spark operators over the Comet children, runs Spark's own rule over that, and swaps the Comet operators back in. It extends AQEShuffleReadRule, so AQE gates and validates it the way it does Spark's rule. Spark 3.4 has no hook for query-stage optimizer rules, so the fix applies from Spark 3.5. Closes apache#6454. * test: fix three Spark SQL test adaptations for Comet's cache format - CachedBatchSerializerNoUnwrapSuite sets its own cache serializer, but Spark keeps the first one it loads for the rest of the JVM, so the suite ran with Comet's. Clear it before and after the suite, as CachedBatchSerializerSuite does (4.1.3, 4.2.0). - The SPARK-19993 helper counted the subqueries in a Comet cache scan's original plan, which subquery reuse never reaches, on top of the same subqueries in the filter above the scan. - SPARK-35332's AQE check reads the cached plan from the SQL UI's plan tree, which shows it only under Spark's own cache scan (apache#6463). Skip it under Comet (4.0.4 and later). * fix: keep Spark's cache format when Kryo would reject Comet's or Comet disables itself CometDriverPlugin installed ArrowCachedBatchSerializer whenever Comet, its native execution and the cache config were enabled at startup. Two more startup settings decide whether that format can work: - With spark.kryo.registrationRequired=true and no CometKryoRegistrator, Kryo rejects CometCachedBatch, so caching failed with "Class is not registered" as soon as Spark serialized a cached block, where Spark's own format works. The plugin only warned. It now keeps Spark's format, and the warning says so. - With Comet shuffle enabled but neither of Comet's shuffle managers configured, Comet disables itself, so every cache was stored in Comet's format with only Spark operators to read it. The plugin now keeps Spark's format there too. The plugin's boolean config reads also honor a deprecated alternative key, such as spark.comet.exec.shuffle.enabled, as a session does. * docs: describe the in-memory cache default in the upgrade guide Add a 1.2.0 section to the upgrade guide for the new default of spark.comet.exec.inMemoryCache.enabled: what changes, how to keep Spark's cache format, and when Comet keeps it without being asked. The operator and compatibility pages still called the feature disabled by default, and the cache guide still said an application that started with the default kept Spark's format. * test: benchmark the in-memory cache against Spark's format with AQE on The published comparisons either read the same Comet-written cache both ways, with spark.comet.sparkToColumnar.enabled turned on, or read both formats with Comet off. Neither is what an application gets from the feature's default. Add cases with Comet and AQE on and Comet's other settings at their defaults, reading a relation cached in Spark's format and in Comet's, with Comet operators above the cache scan and with a Spark operator above it. * feat: record why Spark scans a relation cached in Comet's format when Comet is off spark.sql.cache.serializer is static, so a relation cached in Comet's format keeps it for the life of the application. A session that turns Comet or its native execution off then reads it with Spark's InMemoryTableScanExec, which is slower than reading Spark's own format (apache#5485). CometExecRule records a fallback reason for the other ways Spark ends up scanning a cache, but it returned before looking at the plan in these two cases, so nothing explained them. Record a reason on each such scan, naming the cause and the startup setting that keeps caches in Spark's format. * docs: say an application keeps the cache format it started with The sentence described a session started with the default, which keeps Spark's format only while the default is off. * test: format the in-memory cache benchmark * fix: bind the fallback reason result for the strict Scala warnings build -Ywarn-value-discard rejects the node that withFallbackReason returns, discarded as the last expression of a case. * docs: compare the in-memory cache with Spark's format as an application gets it Add the benchmark's adaptive results to the cache guide: with Comet and AQE on and Comet's other settings at their defaults, Comet's format is as fast or faster than Spark's in every shape but a read of three long columns, which zstd decompression makes about 10% slower. That includes a Spark operator above the native scan. The guide said that, without the feature, a CometSparkColumnarToColumnar converts each batch of Spark's cache scan for Comet. That takes spark.comet.sparkToColumnar.enabled, which is off by default, so by default the operators above Spark's scan run on Spark. The limitation that Spark operators read Comet's format slowly applies to Spark's own cache scan, not to Spark operators above Comet's. * fix: count Kryo registrations however the application made them The gate kept Spark's cache format whenever spark.kryo.registrator did not list CometKryoRegistrator, even where the application had registered Comet's cached batch another way, such as spark.kryo.classesToRegister. Comet's format works there, while Spark's does not before 4.1, which registers DefaultCachedBatch itself only from then on, so the gate broke caches that worked. Ask a Kryo instance built from the application's conf which of Comet's classes it has registered, and keep Spark's format only if Comet's cached batch is not among them. If Kryo cannot be built, fall back to looking for CometKryoRegistrator in spark.kryo.registrator. The startup warning uses the same check and names the missing classes. * docs: describe the Kryo condition by registration in the upgrade guide The plugin now keeps Spark's cache format when Kryo has not registered Comet's cached batch, however the application registers classes, rather than when spark.kryo.registrator does not list CometKryoRegistrator. * docs: name spark.comet.convert.inMemoryCache.enabled for converting Spark's cache scan spark.comet.sparkToColumnar.enabled is deprecated for in-memory cached tables since apache#6602, which gave the conversion its own config. * test: accept Comet's cache scan in BroadcastJoinSuite's cached-relation tests Since apache#6415, suites that build their own SparkSession run Comet, and with the cache on by default BroadcastJoinSuite's cached relations are read by CometInMemoryTableScanExec. SPARK-23214 and SPARK-37742 count Spark's cache scan and broadcast hash join, so count Comet's as well. * fix: coalesce with the whole stage in view and keep partitionings above Comet unions Address review of the CometCoalesceShufflePartitions rule: - Stand in only for Comet operators whose output partitioning is unknown. From Spark 4.1 a union of co-partitioned children reports their partitioning, and an aggregate above reads it without a shuffle. Coalescing one branch made Comet's aggregate, which states no required distribution, return each key twice. - Hand Spark's rule the whole stage rather than the subtree below each Comet operator, so it sees a Cartesian product or nested loop join above a union and uses its smaller target size there, and divides the minimum partition count over every group. Leave a stage that already has a read over a shuffle, which Spark's rule would coalesce again. - Copy the Spark original when withNewChildren hands it back, so a union that AQE plans again over its stages is coalesced too. - Drop the CartesianProductExec case, which has no Comet counterpart, and document the rule in CometSparkSessionExtensions. * test: run Spark's SPARK-42101 union coalescing test with Comet again The test failed with Comet's cache format because AQE did not coalesce the shuffled branch of a union that Comet ran (apache#6454), which apache#6459 fixes. * fix: coalesce Comet union branches in stages that already hold an AQE shuffle read CometCoalesceShufflePartitions returned a stage unchanged when it held any AQEShuffleReadExec, because Spark's rule asserts on reads that are not skew splits. Below a Cartesian product, Spark's rule coalesces the shuffle on the other side by itself first, so the read it leaves there kept the shuffled branch of a Comet union at every partition. Hide each read behind a leaf while Spark's rule runs, and put it back afterwards. Spark's rule then leaves alone the shuffles it would coalesce together with a read one, and coalesces the rest of the stage. * docs: say Spark's cache scan can read Comet's format more slowly The fused reader from apache#5859 reads Comet's format about as fast as Spark's own, and faster for the narrowest reads, so the upgrade note should not say Spark's scan is always slower. * test: leave Comet's cache scan out of the DPP suite's subquery check Since apache#6577, CometInMemoryTableScanExec exposes Spark's own InMemoryTableScanExec as its one subquery, so that Spark's UI draws the cached plan below it. With the cache on by default, the DPP suite's 'filtering ratio policy fallback' caches its dimension table, and checkPartitionPruningPredicate requires every subquery of an adaptive plan to contain an AdaptiveSparkPlanExec, which Spark's scan does not. Skip it there in every diff, since the query never runs it as a subquery. * feat: keep Comet's cache format opt-in on Spark 3.4 Spark 3.4 has no hook for query-stage optimizer rules, so CometCoalesceShufflePartitions cannot run there, and AQE cannot coalesce the shuffled branch of a Comet union whose other branch reads a cache in Comet's format (apache#6454). spark.comet.exec.inMemoryCache.enabled now defaults to true only from Spark 3.5, and the plugin, which follows the entry's default, keeps Spark's format on 3.4 unless it is set. The 3.4 Spark SQL diff goes back to main's: its sessions install Comet's serializer only as the plugin would, which on 3.4 it no longer does by default, so its cache test adaptations are not needed there. * docs: rewrap the Spark 3.4 cache default sentences * test: run Spark's SPARK-35332 AQE cache test with Comet again The diffs skipped CachedTableSuite's 'SPARK-35332: Make cache plan disable configs configurable - check AQE' on Spark 4.0 and later, because SparkPlanInfo drew the cached plan only below Spark's own InMemoryTableScanExec. Since apache#6577 it is drawn below CometInMemoryTableScanExec too, so the test finds the cached plan. The test's check that the cached plan was coalesced, which the diffs already adapt, looked for AQEShuffleRead only below a columnar-to-row transition, where it is when Comet's plan is cached in Spark's format. In Comet's columnar format the cached plan ends in the read itself, as in Spark, so the check accepts it there too. The test passes on Spark 4.0, 4.1 and 4.2 with either format, and with Comet off. * docs: stop describing the in-memory cache as experimental It is enabled by default from Spark 3.5. The operator table now marks InMemoryTableScanExec as supported, and the plan-node table no longer says the cache scan is disabled by default. * test: skip only the format-dependent checks in the Spark SQL diffs SPARK-36120, SPARK-33687, SPARK-22673 and SPARK-43376 now run with Comet's cache format and skip only the size or subquery-reuse assertions that depend on it. The union test in CoalesceShufflePartitionsSuite runs its coalescing check again on Spark 3.5 and later, where CometCoalesceShufflePartitions coalesces the shuffles on the union's join side. From Spark 3.5, CometTestSettings also makes Comet's cache serializer the default of the sql and hive test JVMs, next to the shuffle manager. Spark keeps the first cache serializer it loads for the rest of the JVM, so a suite that caches first through a session of its own left the suites after it in Spark's format. * test: time building the cache in Spark's format in the cache benchmark runBuildBenchmark caches the same relation in Spark's format and in Comet's, with each codec, from Comet operators and from Spark operators, and reports the footprint of each. The cache guide has the results. * test: skip only the columnar checks of SPARK-37371 under Comet Comet runs the union of two cached relations as a CometUnion over its own cache scans, so the checks on Spark's scans and union are skipped there. The union with a local relation stays Spark's non-columnar UnionExec, and both answers are checked, under Comet as well. --------- Co-authored-by: test <a@b.c>
Which issue does this PR close?
Closes #6463.
Rationale for this change
Spark builds the SQL tab's graph, and the plans the event log records, with
SparkPlanInfo.fromSparkPlan. It draws its ownInMemoryTableScanExecwith the cached plan below it, but recognizes that scan by its class. For any other node it takesplan.children ++ plan.subqueries, so for a relation cached in Comet's format the tree ended atCometInMemoryTableScan, and the cached plan's metrics were missing below it.The issue expected that Comet could not fix this alone, because a child would make the cached plan part of the query that reads the cache. A subquery does not: Spark runs subqueries from a plan's expressions and only walks the
subquerieslist.What changes are included in this PR?
It builds on #6574, now merged, and needs its explicit
innerChildrenoverride:QueryPlan.innerChildrendefaults tosubqueries, so without it EXPLAIN would draw this subquery as well, andExtendedExplainInfowould count it.CometInMemoryTableScanExecoverridessubqueriesto return the SparkInMemoryTableScanExecit replaces, not the cached plan.SparkPlanInfodraws the cached plan below that scan through its own special case, so the tree readsCometInMemoryTableScan, then Spark's scan of the cache, then the cached plan. It is alazy valbecause Spark 3.x declaressubqueriesas one, and alazy valoverrides Spark 4'sdefas well.CometSparkPlanInfoHelper, because theSparkPlanInfocompanion object isprivate[execution].Exposing Spark's scan rather than the cached plan leaves every other reader of
subquerieswhere it is for Spark's own scan, because they treat that scan as a leaf or special-case it. An earlier revision exposed the cached plan, and Spark 4.2'sSQLLastAttemptAccumulatorthen reached its Comet shuffle and gave up, so a last-attempt metric used outside the cache returnedNone(see the review thread). Checked against Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0, two differences remain:CollectMetricsExec.collectnow reaches Spark's scan through the subquery and, as it does for that scan in Spark's own plans, collects observed metrics from its cached plan. Nothing changes today, because fix: keep Spark's cache scan for a relation whose cached plan records observed metrics #6421 keeps Spark's scan for any relation whose cached plan records observed metrics.AdaptiveSparkPlanExec.finalPlanUpdateposts one more plan update when the final plan contains the scan outside a query stage, because it checks for any node with subqueries.How are these changes tested?
A new test in
CometInMemoryCacheSuite, with AQE off and on, checks that the scan'sSparkPlanInfonode has Spark's scan below it, with the cached plan's ownSparkPlanInfobelow that. It also checks thatcollectWithSubqueriesdoes not find the cached plan's shuffle in a query that has none of its own. The first check fails without the override, and the second fails when the cached plan itself is exposed, as in the earlier revision.CometInMemoryCacheLastAttemptMetricSuite, underspark/src/test/spark-4.2, runs the scenario from the review: a last-attempt metric in amapover a cache whose plan has a Comet shuffle must still reportSome(100). It is registered in both workflow files. It was not run locally, because the 4.2 profile cannot be built offline here, sorun-all-spark-profilesgives its first result.CometInMemoryCacheSuitepasses locally on the default Spark 4.1 profile (61 tests). Spark'sCachedTableSuitetest "SPARK-35332: Make cache plan disable configs configurable - check AQE", which #5634 skips under Comet, reads the cached plan from this tree and should now find it. Whether it passes was not checked.