Skip to content

fix: draw the cached plan below CometInMemoryTableScan in the SQL tab and event log - #6577

Merged
comphead merged 3 commits into
apache:mainfrom
comphead:cached-plan-sql-tab-6463
Oct 4, 2026
Merged

comphead merged 3 commits into
apache:mainfrom
comphead:cached-plan-sql-tab-6463

Conversation

@comphead

@comphead comphead commented Oct 3, 2026 •

Copy link
Copy Markdown
Contributor

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 own InMemoryTableScanExec with the cached plan below it, but recognizes that scan by its class. For any other node it takes plan.children ++ plan.subqueries, so for a relation cached in Comet's format the tree ended at CometInMemoryTableScan, 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 subqueries list.

What changes are included in this PR?

It builds on #6574, now merged, and needs its explicit innerChildren override: QueryPlan.innerChildren defaults to subqueries, so without it EXPLAIN would draw this subquery as well, and ExtendedExplainInfo would count it.

  • CometInMemoryTableScanExec overrides subqueries to return the Spark InMemoryTableScanExec it replaces, not the cached plan. SparkPlanInfo draws the cached plan below that scan through its own special case, so the tree reads CometInMemoryTableScan, then Spark's scan of the cache, then the cached plan. It is a lazy val because Spark 3.x declares subqueries as one, and a lazy val overrides Spark 4's def as well.
  • A test-only CometSparkPlanInfoHelper, because the SparkPlanInfo companion object is private[execution].

Exposing Spark's scan rather than the cached plan leaves every other reader of subqueries where 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's SQLLastAttemptAccumulator then reached its Comet shuffle and gave up, so a last-attempt metric used outside the cache returned None (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.collect now 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.
  • On Spark 3.4 and 3.5, AdaptiveSparkPlanExec.finalPlanUpdate posts 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's SparkPlanInfo node has Spark's scan below it, with the cached plan's own SparkPlanInfo below that. It also checks that collectWithSubqueries does 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, under spark/src/test/spark-4.2, runs the scenario from the review: a last-attempt metric in a map over a cache whose plan has a Comet shuffle must still report Some(100). It is registered in both workflow files. It was not run locally, because the 4.2 profile cannot be built offline here, so run-all-spark-profiles gives its first result.

CometInMemoryCacheSuite passes locally on the default Spark 4.1 profile (61 tests). Spark's CachedTableSuite test "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.

@github-actions github-actions Bot added bug Something isn't working area:scan Parquet scan / data reading labels Oct 3, 2026
@comphead comphead added the run-all-spark-profiles Run the Comet test suites against every Spark profile on this pull request, ahead of the merge queue label Oct 3, 2026

@sunchao sunchao 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.

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 innerChildren and the cached physical plan through subqueries, 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 subqueries override 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)

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.

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

@comphead comphead Oct 4, 2026 •

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

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.

@comphead
comphead force-pushed the cached-plan-sql-tab-6463 branch from a422a54 to 5d5b520 Compare October 3, 2026 20:41

@sunchao sunchao 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.

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 innerChildren and its cached physical plan through subqueries, 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 None instead of Some(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 subqueries for 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.
@comphead
comphead force-pushed the cached-plan-sql-tab-6463 branch from 1a36ec9 to c190f6e Compare October 4, 2026 16:44

@sunchao sunchao 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.

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 existing SparkPlanInfo handling 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 innerChildren override 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.1 at 992c806a7e38c2e88bd018aa5774164b0850e1fa, 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.

@comphead
comphead added this pull request to the merge queue Oct 4, 2026
Merged via the queue into apache:main with commit f8cc631 Oct 4, 2026
57 checks passed
andygrove added a commit to andygrove/datafusion-comet that referenced this pull request Oct 5, 2026
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.
andygrove added a commit to andygrove/datafusion-comet that referenced this pull request Oct 8, 2026
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.
parthchandra pushed a commit to parthchandra/datafusion-comet that referenced this pull request Oct 9, 2026
* 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>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:scan Parquet scan / data reading bug Something isn't working run-all-spark-profiles Run the Comet test suites against every Spark profile on this pull request, ahead of the merge queue

Projects

None yet

Development

Successfully merging this pull request may close these issues.

The SQL tab and event log lose the cached plan under CometInMemoryTableScan

2 participants