Repository navigation
Conversation
…trings 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.
…on 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.
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.
"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.
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.
…enchmark 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.
…ection-projection
…ection-projection
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.
…ample 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.
…ection' into feat/cache-enabled-by-default # Conflicts: # docs/source/user-guide/latest/in-memory-cache.md
…-default # Conflicts: # spark/src/main/scala/org/apache/spark/sql/comet/execution/arrow/ArrowCachedBatchSerializer.scala
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.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Comet’s Arrow cache was opt-in. Enabling it exposes cached queries to native execution, including the previously reported union-coalescing problem.
- Design approach: Change the default, adapt Spark’s cache tests, document the upgrade, add adaptive benchmarks, and include the stacked coalescing fix from #6459.
- Correctness / compatibility analysis: Compared relevant Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. Earlier Kryo and binary-reverse concerns are addressed. The prerequisite’s distribution and ancestor-join regressions are also fixed. One existing coalescing concern remains, described below.
- Key design decisions: Reusing Spark’s coalescer avoids duplicating version-specific sizing logic. The
UnknownPartitioningguard protects downstream distribution assumptions. The added rule performs driver-side plan traversal, without adding batch processing. Its whole-stage early return leaves some independent groups unhandled. - Implementation sketch: Reviewed all 15 base-relative files, including the prerequisite. Reconstructed all 71 affected Spark source files and verified upstream and patched Git blob hashes. Test adaptations retain answer checks and distinguish Spark-specific plan and statistics assertions.
- Behavioral changes worth calling out: Compared the affected paths with
branch-1.1at992c806a7e38c2e88bd018aa5774164b0850e1fa. Eligible applications intentionally receive Arrow caches automatically. Runtime disabling retains that format. Published benchmarks show improvements for several shapes, approximately 10% slower three-long-column reads, and slower Spark-scan fallback reads. These measurements were not independently rerun. - Suggested improvements: Complete the existing #6454 fix for stages containing both an already-coalesced branch and an untouched cached-union branch before enabling the default.
No additional introduced P1/P2 issues found within this review beyond the existing concern below.
The existing P2 concern in #6454 remains partially unresolved. For CartesianProductExec(CometUnionExec(shuffle, cacheScan), shuffle), Spark can coalesce the independent right branch first. At spark/src/main/scala/org/apache/comet/rules/CometCoalesceShufflePartitions.scala:83, finding that AQEShuffleReadExec makes the new rule return the entire stage unchanged. The cached union’s shuffled branch consequently retains 200 partitions where Spark uses one. The new cache default exposes this shape by allowing the cached branch and union to become Comet operators.
A freshly compiled Spark 4.1.3 probe reproduced MIXED_READ_SPARK=Vector(1, 1); COMET=Vector(200, 1) using two 200-partition shuffles with 1 KiB per-partition statistics and a real Spark cache scan as the non-shuffle leaf. This establishes unnecessary reduce partitions without claiming measured wall-clock overhead. Reproduce with python3 /tmp/5634-current-rule-prc_9oi_/run.py. Handle untouched independent groups while preserving existing AQE reads, ancestor sizing context and distribution guards. This retains the previously reported concern rather than duplicating it as a new finding.
Reviewed full head 14e162b62482f908e0e1a978f2c948326baff83a against base 8a7e6644139c55a10b8f3ee1b1c367bf7eac7d91. Confirmed non-draft status and read existing discussions, excluding Copilot. Routed skills: review-comet-pr, review-comet-shuffle-pr, review-comet-ffi-pr, review-comet-memory-pr, and review-comet-expression-pr.
Exact-head CI at 17:14 UTC: 14 successful checks, nine running, 11 queued and 50 skipped, with no reported failures. All five Spark SQL profiles and Iceberg builds are underway. There is no completed build/integration verdict yet.
Validation: Fresh planner/RDD checks confirmed the prerequisite fixes and reproduced the remaining case. Spark 3.5.9 startup probes passed 21 assertions plus Comet-payload and Spark-fallback disk write/read checks. Verified and reused 16 prior IPC sizing/lifecycle cases against byte-identical source, with Arrow allocations returning to zero.
Limits: Probes use minimal adapters and synthetic shuffle statistics. No full Comet/native build, native query execution, full Spark SQL matrix, Iceberg suite or performance benchmark was run locally. Project code remains unchanged.
… 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.
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.
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.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Comet’s Arrow cache was opt-in, limiting coverage outside dedicated cache suites.
- Design approach: Enable the existing cache by default, adapt Spark’s tests, document the upgrade, add adaptive benchmarks, and include #6459’s union-coalescing fix.
- Correctness / compatibility analysis: Compared relevant Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. Earlier Kryo and binary-reverse concerns are addressed. Fresh Spark 4.1 checks pass for mixed AQE reads, ancestor-join sizing and distribution preservation. The existing coalescing concern remains exposed on Spark 3.4.
- Key design decisions: Reusing Spark’s coalescer avoids duplicating version-specific sizing logic.
UnknownPartitioningguards and temporary leaves preserve distribution assumptions and existing AQE reads. The added work is driver-side plan traversal, without new per-batch processing. - Implementation sketch: Reviewed all 15 base-relative files, including the stacked prerequisite. Reconstructed and verified upstream and patched Git hashes for all 76 affected Spark source files. Reviewed the benchmark’s serializer and executed-plan assertions and the test adaptations.
- Behavioral changes worth calling out: Compared affected paths with
branch-1.1at7b7eec69282a7abc33a0e8e687ca594c72f140ec. Automatic Arrow caching and improved coalescing on Spark 3.5+ are intended. Runtime disabling retains the stored format. Published benchmarks document gains and roughly 10% slower three-long-column reads, plus slower Spark-scan fallback cases. These measurements were not independently rerun. - Suggested improvements: Address the remaining Spark 3.4 exposure of #6454 before enabling the cache by default there.
No additional introduced P1/P2 issues found within this review beyond the existing blocker below.
The existing P2 concern in #6454 remains unresolved on Spark 3.4. The default change at spark/src/main/scala/org/apache/comet/CometConf.scala:362 allows a union of a native shuffled input and a cached relation to become CometUnionExec. However, registration at spark/src/main/scala/org/apache/comet/CometSparkSessionExtensions.scala:114 is a no-op through the 3.4 shim. Spark’s coalescer recognizes UnionExec, but cannot independently coalesce the shuffled branch of this Comet union.
A freshly compiled Spark 3.4.3 probe using the exact-head union implementation, a real materialized 200-partition shuffle and a real cached sibling reproduced one shuffle partition under Spark’s union versus 200 under Comet’s. This establishes unnecessary reduce partitions without claiming measured elapsed-time overhead. Reproduce with python /tmp/5634-root3483-spark34-hr6c7za7/run.py. Preserve Spark’s union until coalescing on 3.4, provide equivalent supported integration, or retain the opt-in default on that version. This retains the existing concern rather than duplicating it as a new finding.
Reviewed full head 3483f5b14f706d9edbe88157733ee81f9d8c81d8 against base 74257803f1a364a4e0b177f526b5ba30f637df80. Confirmed non-draft status and read existing discussions, excluding Copilot. Routed skills: review-comet-pr, review-comet-shuffle-pr, review-comet-ffi-pr, review-comet-memory-pr, and review-comet-expression-pr.
Exact-head CI: 107 successful checks, 11 skipped, no failures. All five Spark SQL profiles and Iceberg 1.8–1.11 passed. The Linux execution job passed 1,299 tests, including the new coalescing regressions. Its merge checkout has the same tree as this head. macOS, Delta, PyArrow and benchmark checks were skipped.
Validation limits: Local planner probes use minimal Comet adapters. The 4.1 checks use synthetic shuffle statistics and the production union RDD shim. No full local Comet/native build, native query execution, Spark SQL matrix or performance benchmark was run. Project code remains unchanged.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Comet’s Arrow cache was opt-in, limiting coverage outside dedicated cache suites.
- Design approach: Enable the existing cache by default, adapt Spark’s tests, document the upgrade, add adaptive benchmarks, and include #6459’s coalescing fix.
- Correctness / compatibility analysis: Compared relevant Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. Earlier Kryo and binary-reverse concerns are addressed. Fresh Spark 4.1 probes pass for mixed AQE reads, ancestor-join sizing and distribution preservation. The existing Spark 3.4 concern remains.
- Key design decisions: Reusing Spark’s coalescer avoids duplicating version-specific sizing logic.
UnknownPartitioningguards protect distribution assumptions, and temporary leaves preserve existing AQE reads. The rule adds driver-side traversal without new per-batch processing. - Implementation sketch: Reviewed all 15 base-relative files, including the stacked prerequisite. Verified upstream and patched Git hashes for all 76 affected Spark source files and inspected every distinct test adaptation, benchmark change and documentation change.
- Behavioral changes worth calling out: Compared affected paths with
branch-1.1at7b7eec69282a7abc33a0e8e687ca594c72f140ec. Automatic Arrow caching and improved coalescing on Spark 3.5+ are intended. Runtime disabling retains the stored format. Published benchmarks document gains, roughly 10% slower three-long-column reads, and slower Spark-scan fallback cases. Those measurements were not independently rerun. - Suggested improvements: Address the existing Spark 3.4 exposure of #6454 before enabling the cache by default there.
No additional, non-duplicate introduced P1/P2 issues found within this review.
The existing P2 blocker in #6454, already reported in the exact-head review, remains substantiated. The default change at spark/src/main/scala/org/apache/comet/CometConf.scala:362 allows a native shuffled input unioned with a cached relation to become CometUnionExec. Registration at spark/src/main/scala/org/apache/comet/CometSparkSessionExtensions.scala:114 is a no-op through the Spark 3.4 shim, so the new rule cannot restore independent coalescing there. A freshly compiled Spark 3.4.3 probe using the exact-head union implementation, a real materialized shuffle and a real cached sibling reproduced one shuffle partition under Spark’s union versus 200 under Comet’s. This establishes unnecessary reduce partitions without claiming measured elapsed-time overhead. Preserve Spark’s union until coalescing, provide equivalent integration on 3.4, or retain opt-in behavior there. Reproduction: python3 /tmp/5634-review-validation-p4os90pg/spark34/run.py.
Reviewed full SHA 3483f5b14f706d9edbe88157733ee81f9d8c81d8 against base 74257803f1a364a4e0b177f526b5ba30f637df80. Confirmed non-draft status and read existing reviews, issue comments, inline comments and threads, excluding Copilot. Routed skills: review-comet-pr, review-comet-shuffle-pr, review-comet-ffi-pr, review-comet-memory-pr, and review-comet-expression-pr.
Exact-head CI: 107 successful checks, 11 skipped, no failures. All five Spark SQL profiles and Iceberg 1.8–1.11 passed. The Linux execution job passed 1,299 tests, including all seven added regressions. Its merge checkout has the same tree as the reviewed head. macOS, Delta, PyArrow and benchmark checks were skipped.
Validation limits: Local probes use minimal Comet adapters. Spark 4.1 checks use synthetic shuffle statistics and the production union RDD shim. No full local Comet/native build, native query execution, Spark SQL matrix or performance benchmark was run. Project code remains unchanged.
|
The mixed-stage case of #6454 is fixed in #6459 and merged here. One difference from Spark remains. A group that already has a read no longer counts when Spark's rule divides the minimum partition count over the stage's groups, so the groups it coalesces can keep more partitions than on Spark, never fewer. The CI run at 14e162b failed one test on every Spark version: Main is merged in again, including #5859. The cache guide keeps its description of the fused reader and its measurements. |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Comet’s Arrow cache was opt-in, limiting coverage outside dedicated cache suites.
- Design approach: Enable the existing cache by default, adapt Spark’s tests, document the upgrade, add adaptive benchmarks, and include #6459’s shuffle-coalescing fix.
- Correctness / compatibility analysis: Compared relevant Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. Earlier Kryo and binary-reverse concerns are addressed. Fresh Spark 4.1 probes pass mixed-stage coalescing, ancestor sizing and distribution-preservation checks. The existing Spark 3.4 blocker remains.
- Key design decisions: Reusing Spark’s coalescer avoids duplicating version-specific sizing logic. The
UnknownPartitioningguard protects downstream distribution assumptions, and temporary leaves preserve existing AQE reads. The rule adds driver-side plan traversal without introducing per-batch processing. - Implementation sketch: Reviewed all 15 base-relative changed files, including the stacked prerequisite. Verified reconstructed base/head source hashes for all 76 affected Spark files and inspected every distinct test adaptation, benchmark change and documentation change.
- Behavioral changes worth calling out: Compared affected paths with
branch-1.1at7b7eec69282a7abc33a0e8e687ca594c72f140ec. Automatic Arrow caching and improved coalescing on Spark 3.5+ are intended. Runtime disabling retains the stored format. Published benchmarks document gains, approximately 10% slower three-long-column reads, and slower Spark-scan fallback cases. These measurements were not independently rerun. - Suggested improvements: Address the existing Spark 3.4 exposure of #6454 before enabling the cache by default there.
No additional, non-duplicate introduced P1/P2 issues found within this review.
The existing P2 blocker in #6454, already reported in the exact-head reviews, remains substantiated. The default change at spark/src/main/scala/org/apache/comet/CometConf.scala:362 allows a native shuffled input unioned with a cached relation to become CometUnionExec. However, registration at spark/src/main/scala/org/apache/comet/CometSparkSessionExtensions.scala:114 is a no-op through the Spark 3.4 shim. The new rule therefore cannot restore independent coalescing there.
A freshly compiled Spark 3.4.3 probe using the verified exact-head union implementation, a real materialized shuffle and a real cached sibling reproduced one shuffle partition under Spark’s union versus 200 under Comet’s. This demonstrates unnecessary reduce partitions without claiming measured elapsed-time overhead. Preserve Spark’s union until coalescing, provide equivalent integration on 3.4, or retain opt-in behavior there. Reproduction: python /tmp/5634-final-validation-4cj990fw/spark34/run.py. This retains the existing concern rather than duplicating it as a new finding.
Reviewed full SHA 3483f5b14f706d9edbe88157733ee81f9d8c81d8 against base 74257803f1a364a4e0b177f526b5ba30f637df80. Confirmed non-draft status and read existing reviews, issue comments, inline comments and threads, excluding Copilot. Routed skills: review-comet-pr, review-comet-shuffle-pr, review-comet-ffi-pr, review-comet-memory-pr, and review-comet-expression-pr.
Exact-head CI: 107 successful checks, 11 skipped, no failures. All five Spark SQL profiles and Iceberg 1.8–1.11 passed. The Linux execution job passed 1,299 tests, including all seven added regressions. Verified that its merge checkout has the same tree as the reviewed head. macOS, Delta, PyArrow and benchmark checks were skipped.
Validation limits: Local probes use minimal Comet adapters. Spark 4.1 checks use synthetic shuffle statistics and the production union RDD shim. No full local Comet/native build, native query execution, Spark SQL matrix or performance benchmark was run. Project code remains unchanged.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Comet’s Arrow cache was opt-in, limiting coverage outside dedicated cache suites.
- Design approach: Enable the existing cache by default, adapt Spark’s tests, document the upgrade, add adaptive benchmarks, and include #6459’s coalescing fix.
- Correctness / compatibility analysis: Compared relevant Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. Earlier Kryo-registration and binary-reverse concerns are addressed. Fresh Spark 4.1 probes pass mixed-stage coalescing, ancestor-sizing and distribution-preservation checks. The existing Spark 3.4 blocker remains.
- Key design decisions: Reusing Spark’s coalescer avoids duplicating version-specific sizing logic. The
UnknownPartitioningguard protects downstream distribution assumptions, and temporary leaves preserve existing AQE reads. The rule adds driver-side traversal without adding per-batch processing. - Implementation sketch: Reviewed all 15 files in the full base-relative diff, including the stacked prerequisite. Reconstructed all 76 affected Spark source files, verified upstream and patched Git hashes, and inspected the test adaptations, benchmark and documentation.
- Behavioral changes worth calling out: Compared affected paths with
branch-1.1at7b7eec69282a7abc33a0e8e687ca594c72f140ec. Automatic Arrow caching and improved coalescing on Spark 3.5+ are intended. Runtime disabling retains the stored format. Published benchmarks document gains, approximately 10% slower three-long-column reads, and slower Spark-scan fallback cases. These measurements were not independently rerun. - Suggested improvements: Address the existing Spark 3.4 exposure of #6454 before enabling the cache by default there.
No additional, non-duplicate introduced P1/P2 issues found within this review.
The existing P2 blocker in #6454 remains substantiated. The default change at spark/src/main/scala/org/apache/comet/CometConf.scala:384 allows a native shuffled input unioned with a cached relation to become CometUnionExec. Registration at spark/src/main/scala/org/apache/comet/CometSparkSessionExtensions.scala:114 is a no-op through the Spark 3.4 shim, so the new rule cannot restore independent coalescing there.
A freshly compiled Spark 3.4.3 probe using the verified current union implementation, a real materialized shuffle and a real cached sibling reproduced one shuffle partition under Spark’s union versus 200 under Comet’s. This demonstrates unnecessary reduce partitions without claiming measured elapsed-time overhead. Preserve Spark’s union until coalescing, provide equivalent integration on 3.4, or retain opt-in behavior there. Reproduction: python3 /tmp/5634-dbca-review-pmkjsd8s/spark34/run.py. This retains the existing concern without duplicating it as a new finding.
Reviewed full SHA dbca563e894b34d53be95e64a6b9125f0092f0d8 against base 9dc8c3ca89962506078831f7064a25c75da83883. Confirmed non-draft status and read existing reviews, issue comments, inline comments and threads, excluding Copilot. Routed skills: review-comet-pr, review-comet-shuffle-pr, review-comet-ffi-pr, review-comet-memory-pr, and review-comet-expression-pr.
Exact-head CI at 11:44 UTC: seven successful checks, 20 running and five skipped, with no reported failures. All five Spark SQL profiles and Iceberg builds are underway in run 37457592767. There is no completed integration verdict. macOS, benchmark, Delta and PyArrow checks were skipped.
Validation: Fresh startup probes passed 21 assertions plus Comet-payload and Spark-fallback disk write/read checks. Verified and reused 16 passing IPC sizing/lifecycle cases against unchanged implementation code, with Arrow allocations returning to zero.
Limits: Local probes use minimal Comet adapters. Spark 4.1 checks use synthetic shuffle statistics and the production union RDD shim. No full local Comet/native build, native query execution, Spark SQL matrix, Iceberg suite or performance benchmark was run. Project code remains unchanged.
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.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Comet’s Arrow cache was opt-in. Enabling it also exposes cached unions to Spark’s incomplete recognition of Comet operators during AQE coalescing.
- Design approach: Enable the cache by default on Spark 3.5+, retain opt-in behavior on Spark 3.4, and include the stacked coalescing fix from #6459.
- Correctness / compatibility analysis: Compared relevant Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. Focused checks confirm preserved union partitioning, ancestor-join sizing and existing AQE reads. Earlier startup, Kryo, binary-reverse and coalescing concerns are addressed. Keeping Spark 3.4 opt-in resolves the remaining default-enabled exposure.
- Key design decisions: Reusing Spark’s coalescer avoids duplicating version-specific sizing logic. The
UnknownPartitioningguard protects downstream distribution assumptions. Temporary leaves preserve existing shuffle reads. The rule adds driver-side traversal without adding per-batch processing. - Implementation sketch: Reviewed all 14 files in the full base-relative diff, including the prerequisite, tests, benchmark and documentation. Reconstructed all 63 affected Spark source files and verified upstream and patched Git blob hashes. Test adaptations distinguish Spark-specific assertions from result checks.
- Behavioral changes worth calling out: Compared affected paths with
branch-1.1at7b7eec69282a7abc33a0e8e687ca594c72f140ec. Automatic Arrow caching on Spark 3.5+ and improved coalescing are intended changes. Runtime disabling retains the stored format. Published benchmarks document gains, approximately 10% slower three-long-column reads, and slower Spark-scan fallback cases. These measurements were not independently rerun. - Suggested improvements: None at P1/P2 priority. No introduced P1/P2 issues found within this review.
Reviewed full SHA 88258a7367a11772cf0d66a7e069d05b82fe768a against base 9dc8c3ca89962506078831f7064a25c75da83883. Confirmed non-draft status and read existing reviews, issue comments, inline comments and threads, excluding Copilot. Routed skills: review-comet-pr, review-comet-shuffle-pr, review-comet-ffi-pr, review-comet-memory-pr and review-comet-expression-pr.
Exact-head CI at 12:48 UTC: 21 successful checks, 36 running and five skipped, with no failures reported. Spark SQL and Iceberg validation remain underway in run 37464180395. There is no completed integration verdict. macOS, benchmark, Delta and PyArrow checks were skipped.
Validation: Fresh Spark 4.1.3 planner/RDD probes passed. Spark 3.4.3 and 3.5.9 startup probes each passed 21 assertions plus Comet-payload and Spark-fallback disk write/read checks. Verified and reused 16 passing IPC sizing/lifecycle cases against byte-identical implementation code, with Arrow allocations returning to zero.
Limits: Local probes use minimal adapters and synthetic shuffle statistics. No full local Comet/native build, native query execution, Spark SQL matrix, Iceberg suite or performance benchmark was run. Project code remains unchanged.
|
On Spark 3.4, Comet's cache format now stays opt-in (88258a7). The Spark 3.4 diff goes back to the one on
|
|
|
||
| - test("SPARK-35332: Make cache plan disable configs configurable - check AQE") { | ||
| + test("SPARK-35332: Make cache plan disable configs configurable - check AQE", | ||
| + IgnoreComet("Spark's SQL UI shows a cached plan only under Spark's own cache scan")) { |
Which issue does this PR close?
Part of #5487.
#6459, which fixes #6454, has landed, as have #6537, #6538, #6412 and #6421, which this branch was stacked on before.
Rationale for this change
spark.comet.exec.inMemoryCache.enabledhas been off by default since #5051, so the native cache path only ever runs underCometInMemoryCacheSuiteandCometInMemoryCacheKryoSuite, which exercise it deliberately. Nothing tells us how it behaves under the rest of the suite: the Spark SQL test diffs, the fuzz suites, the Iceberg and Delta jobs, and any test that callscache()/persist()incidentally.This PR flips the default so a full CI run exercises the cache format everywhere caching happens. It was opened as a draft to collect that signal, and now proposes the default for 1.2.0. The review asked for four things before that, which are all on
mainor merged into this branch:What changes are included in this PR?
spark.comet.exec.inMemoryCache.enableddefaults totruefrom Spark 3.5 and staysfalseon Spark 3.4. Spark 3.4 has no hook for query-stage optimizer rules, so fix: coalesce shuffle partitions under Comet unions the way Spark does #6459's rule cannot run there, and AQE could not coalesce the shuffled branch of a Comet union whose other branch reads a cache in Comet's format.CometDriverPluginalready follows the entry's default since fix: install Comet's cache serializer only when Comet and native execution are enabled #6360, so the plugin needs no change here. A new test inCometInMemoryCacheSuitepins the default for each version.The in-memory cache user guide records the new default, says why Spark 3.4 keeps the old one, and shows how to turn the feature off or on. Its Limitations section keeps perf: fuse Comet cache vector reads into Spark codegen #5859's account of how Spark operators read Comet's format, and no longer gives the slower reads as the reason the feature is off by default. The 1.2.0 section of the upgrade guide describes the change, how to keep Spark's format, and when Comet keeps it without being asked. The operator and compatibility pages say the feature is enabled by default from Spark 3.5.
CometInMemoryCacheBenchmarkgains cases that compare the two formats as an application gets them: Comet and AQE on, Comet's other settings at their defaults, with Comet operators or a Spark operator above the cache scan. The results are in the cache guide.The Spark SQL diffs for Spark 3.5 and later install Comet's cache serializer. They turn Comet on through
SharedSparkSessionandTestHiverather than by loadingCometPlugin, and the plugin is what installs the serializer, so until now every Spark SQL suite cached in Spark's own format whatever the default was. Both now setspark.sql.cache.serializerwhen Comet is enabled, as the plugin would. The tests that assume Spark's cache scan are adapted:CometInMemoryTableScanExec;PartitionBatchPruningSuitestill checks every answer, but reads its test-only accumulators only from Spark's scan, which Comet's lacks;CacheTableInKryoSuiteregisters Comet's classes withCometKryoRegistrator, as Comet asks of any application that setsspark.kryo.registrationRequired;CachedBatchSerializerNoUnwrapSuiteclears Spark's JVM-wide cache serializer before and after it runs, asCachedBatchSerializerSuitealready does, so it tests the serializer it sets rather than Comet's;BroadcastJoinSuite'sSPARK-23214andSPARK-37742tests also countCometInMemoryTableScanExecandCometBroadcastHashJoinExec. Since test: run Comet in Spark SQL and Iceberg suites that build their own sessions #6415, suites that build their own session run Comet too, and Spark keeps one cache serializer for the whole JVM, so this suite caches in Comet's format after the suites before it;DynamicPartitionPruningSuite's check that every subquery of an adaptive plan is adaptive skips Spark'sInMemoryTableScanExec. Since fix: draw the cached plan below CometInMemoryTableScan in the SQL tab and event log #6577, Comet's cache scan exposes Spark's scan as its one subquery so that the SQL UI can draw the cached plan below it, but the query never runs it as one;IgnoreComet: exact cache size estimates, the union's columnar support, and subquery reuse through the scan's predicates.Spark's
SPARK-42101: Coalesce shuffle partition with union even if exists TableCacheQueryStageruns with Comet again: it failed because of AQE does not coalesce shuffle partitions under a CometUnion when a branch is not a shuffle #6454, which fix: coalesce shuffle partitions under Comet unions the way Spark does #6459 fixes.The Spark 3.4 diff is the same as on
main. Its sessions cache in Spark's format, as the plugin now does by default on that version, so none of these adaptations apply there.Each diff was regenerated from a clone of its Spark tag with the existing diff applied, after checking that it round-tripped unchanged. Where
mainchanged the same diffs in the meantime, no Spark file was changed on both sides, so each merged diff is the union of the two sides file by file, and applies to its tag.The feature is still described as experimental.
How are these changes tested?
The point of the PR is the CI run itself. On Spark 3.5 and later, every job now builds cached tables in Comet's Arrow format wherever a test caches anything, rather than only in the two suites that opt in, and that now includes the Spark SQL suites.
The test changes come from running Spark's cache suites and 54 other
sql/coresuites that cache, from each version's test jar, once with Spark's cache format and once with Comet's, on Spark 3.4, 3.5, 4.0, 4.1 and 4.2, and from the first Spark 4.1 CI run of this PR. Every test that fails only with Comet's format is adapted or tagged above, except for real bugs. Two of them were fixed by PRs that have since landed:DataFrameCallbackSuite'sSPARK-35695: get observable metrics with persist by callback, fixed by fix: keep Spark's cache scan for a relation whose cached plan records observed metrics #6421.AdaptiveQueryExecSuite'sSPARK-37742: AQE reads invalid InMemoryRelation stats and mistakenly plans BHJ. The cached relation reported its compressed size, 162,440 bytes, under the test's 1,048,584-byte broadcast threshold, so AQE broadcast it. fix: report a cached relation's decoded size to the planner, not its compressed size #6412 makes a cached relation report its decoded size, 1,447,524 bytes here.Two more failed until they were fixed:
DataFrameFunctionsSuite'sreverse function - binary, fixed by fix: run reverse of a binary value through the codegen dispatcher #6461. The bug is not specific to the cache, which only makes this test's projection native.AdaptiveQueryExecSuite'sSPARK-42101: Coalesce shuffle partition with union even if exists TableCacheQueryStage, fixed by fix: coalesce shuffle partitions under Comet unions the way Spark does #6459. AQE did not coalesce the shuffled branch of a union that Comet ran, when another branch was a scan or a table-cache stage (AQE does not coalesce shuffle partitions under a CometUnion when a branch is not a shuffle #6454).#6453 and #6577 fixed two more problems those runs found, in how
EXPLAINand the SQL UI show Comet's cache scan.DynamicPartitionPruning*SuiteAEOn'sfiltering ratio policy fallbackfailed on every Spark version, insql_core-1and the Hive jobs, in the CI run of an earlier head, once #6577 was merged in. Its cached-dimension case planned the pruning filter as an adaptive subquery and returned the right rows, but the suite's subquery check also met the Spark scan that Comet's cache scan exposes for the UI. A scratch suite that runs that case on Spark 4.1 reproduces the failure from the check alone, and passes with Spark's scan skipped.BroadcastJoinSuite'sSPARK-23214andSPARK-37742failed in the Spark 4.1 and 4.2 CI run of an earlier head, after the merge that brought in #6415. With Comet's cache format, both queries plan oneCometInMemoryTableScanExecand no broadcast join, with and without AQE, which the adapted assertions now accept.The existing cache suites are unaffected: they set the config explicitly, including the two cases that set it to
false.CometInMemoryCacheSuite,CometInMemoryCacheKryoSuite,CometInMemoryCacheKryoUnregisteredSuite,CometInMemoryCacheKryoClassesToRegisterSuite,CometInMemoryCachePruningSuite,CachedBatchRowIteratorSuite,CometPluginsSuite,CometPluginsDefaultSuite,CometExecSuite, with the union coalescing tests that #6459 adds, andCometJoinSuitepass locally on this branch with Spark 4.1, and the branch passes the Spark 3.5 strict-warnings and scalafix checks. The cache, Kryo, pruning and plugin suites also pass on Spark 3.4, where the default isfalse.The new adaptive cases in
CometInMemoryCacheBenchmarkran on an AMD Ryzen 9 7950X3D, with Comet and AQE on. Comet's format came out 1.2x to 1.5x faster than Spark's with Comet operators above the cache scan, and 1.0x to 2.4x faster with a Spark operator above it. The exception is a read of three long columns, which is about 10% slower in both shapes because ofzstddecompression; with thenonecodec it is faster too. The tables are in the cache guide.