docs: split the PR review skill by area and correct the shuffle contributor docs - #6018
Conversation
review-comet-pr keeps the workflow that applies to every PR and gains a routing table that points at the area skill for the files a PR touches. Expression review depth moves into review-comet-expression-pr, and new sibling skills cover FFI, memory management, and shuffle. Each area skill names the contributor guide docs to read for that subsystem and lists the specific claims in those docs that a PR can make stale.
The Rust Side table pointed at native/core/src/execution/shuffle/, which no longer exists. The implementation is now its own workspace crate, datafusion-comet-shuffle, under native/shuffle/, and codec.rs is gone. Replace the table with the current file layout, including the partitioners and writers modules. The Memory Management section named PartitionBuffer and SpillFile, neither of which exists; describe what MultiPartitionShuffleRepartitioner actually holds, including the pinned-buffer deduplication that keeps sliced batches from being charged once per slice. The default compression codec is lz4, not zstd, in both shuffle docs.
spark.comet.shuffle.directRead.enabled defaults to true and makes native operators decode compressed shuffle blocks in the plan instead of importing Arrow arrays over FFI, but neither shuffle doc mentioned it and the read paths described only the JVM decode route. Add a Direct Read section to native_shuffle.md covering how the path is selected during plan serialization, why the serialized protobuf rather than the config is the source of truth for which slots are direct read, the block iterator's header format and buffer lifetime contract, and the ShuffleScanExec constraints on JNI thread affinity and dictionary unpacking. Note the branch in both read paths and add the config to both tables. Expand the CometShuffleBlockIterator aside in ffi.md, since a direct-read slot does not cross the FFI boundary at all.
comphead
left a comment
There was a problem hiding this comment.
Thanks @andygrove this is a good direction to identify a scope first and get more detailed context for the modified part
There was a problem hiding this comment.
Correctness
Reviewed e2f054991f28c49f85e53e8bb4a35985ca297be6 against 4bf725f62c0b8fe61b432730f8f960671b7394ac. This changes eight Markdown files: it turns the expression-heavy review skill into a generic entry point with four area skills, and updates shuffle/FFI contributor guidance. There are no executable-code changes, so SQL results, ANSI/error behavior, null handling, overflow, and runtime fallback decisions are unchanged by this PR.
The documentation nevertheless needs four corrections before reviewers can rely on its new checks: unsupported shuffle output types do not retry as a regular native Scan. Local native shuffle partitions share one spill file. Nested hash keys have an opt-in native path. The FFI checklist names removed selection-vector classes and an absent import hook. The inline comments link the exact implementation and existing test source for each.
For the Spark compatibility review, I compared the hash-partition contract with maintained Spark 3.5 (5947fd6e74a1) and 4.0 (03f28fc43180) sources: both require a nonnegative partition ID through Pmod, and 4.0 uses collation-aware hashing. The review guidance must preserve the corresponding type/collation and version gates. The nested-key rule currently obscures one of those distinctions. The required maintained Spark 3.4/4.1 branches were unavailable locally, so I do not claim source or runtime verification for them.
Validation
All eight files pass local Prettier 3.9.6, and local link/anchor and sibling-skill routing checks pass. CI Preflight also passed formatting on merge a3a2c4f38adb5bc850980a99035a91a765479593. Its parents are the reviewed base/head and its tree equals the head tree. At the review cutoff, seven checks succeeded and fifteen were skipped. Spark/native builds and tests were skipped for this documentation-only change. I inspected the relevant existing test source without claiming test execution. Local Sphinx rendering stopped at the unavailable sphinxcontrib.mermaid extension, so rendered documentation is not verified.
Performance
The PR adds no runtime allocations, copies, scans, or JNI work. The new material correctly calls attention to buffer ownership, reservation accounting, and the direct-read path that avoids the JVM decode/FFI round trip. These are review targets, not measured performance improvements from this PR. Correcting the shared-spill ownership and default-versus-opt-in nested hashing descriptions matters because they determine which resource and partitioning costs a future reviewer should examine. No runtime benchmark is warranted for these Markdown changes.
Design
The prior skill mixed generic workflow with detailed expression checks, leaving other subsystems without equivalent guidance. A generic router plus expression, FFI, memory, and shuffle siblings is a reasonable separation: overlapping file patterns can select multiple skills, while metadata, CI, fallback explanation, and documentation freshness remain common concerns. The contributor-guide links and direct-read explanation make the intended review path concrete. The four factual corrections are the remaining actionable improvements. The overall split does not need a different mechanism.
Abstraction & complexity
The new units are plain Markdown skills with explicit names, routes, and links back to the generic workflow. They add no code abstraction or execution dependency. I checked that all four routes resolve and that the extracted expression material retains the support-level, scalar-return-type, SQL-test directive, timezone, overflow, and benchmark guidance. Keeping ownership and fallback claims aligned with the implementation is more useful here than adding further checklist layers. The obsolete FFI selection-vector section should be replaced with current contracts rather than preserved as an additional abstraction.
| `ShuffleScan` operator in place of the usual `Scan`. If any output type fails | ||
| `supportedSinkDataType`, it records the fallback reason `Unsupported data type for shuffle direct read` | ||
| and the slot serializes as a regular `Scan` instead. |
There was a problem hiding this comment.
Correctness
[P2] Describe the actual unsupported-type fallback
When supportedSinkDataType rejects a shuffle output, convertToShuffleScan returns None. convert does not retry super.convert. The ordinary Scan conversion also rejects the same types. The AQE caller retains the original Spark stage on that failure. This paragraph therefore gives readers the wrong native-conversion boundary when diagnosing unsupported schemas. Please describe the failed native conversion here, and reserve the regular Scan description for inputs/configurations where shouldUseShuffleScan is false.
There was a problem hiding this comment.
You're right, I had the boundary wrong. convertToShuffleScan returns None and convert never falls through to super.convert, and even if it did, the base CometSink.convert gates on the same supportedSinkDataType, so there is nothing to fall back to. The stage just stays as Spark's, via the getOrElse(s) in CometExecRule.
Rewrote the paragraph in 67e734a to say that, and moved the regular Scan description up to the case it actually describes, which is shouldUseShuffleScan returning false.
| defaults to `0`, meaning the fixed limit is disabled and memory pressure is the only trigger. Each | ||
| partition has its own spill file and multiple spills for a partition are concatenated when the | ||
| final output is written. |
There was a problem hiding this comment.
Correctness
[P2] Document one shared spill file for local native shuffle
The current local writer creates one PartitionedSpill for all output partitions. It owns a single spill file and records per-partition byte ranges. The existing spilling_every_partition_creates_one_file test explicitly asserts one file after repeatedly spilling 64 partitions. Teaching reviewers to expect a file per partition misstates the resource ownership and spill lifecycle they need to check. Please describe the shared file and ordered per-partition ranges, and make the same correction to the newly added spill.rs table entry and spill paragraph in native_shuffle.md.
There was a problem hiding this comment.
Confirmed. PartitionedSpill's own doc comment says one file shared by every output partition, and it holds ranges: Vec<Vec<Range<u64>>> rather than a file handle per partition, so spilling_every_partition_creates_one_file is asserting the design rather than an incidental detail.
Fixed in 67e734a, in the skill and in all three places in native_shuffle.md: the spill.rs table row, the Memory Management bullet, and the paragraph at the end of that section. I also noted that the single-partition writer and the RSS writer do not spill at all, since the old wording implied every native path had spill files worth inspecting.
| Complex types are fully supported as **data** columns in both. The primitive-only restriction | ||
| applies to **partition keys** for `HashPartitioning` and `RangePartitioning` only. |
There was a problem hiding this comment.
Correctness
[P2] Include the opt-in nested hash-key path
The primitive-only rule is not unconditional for HashPartitioning. With spark.comet.shuffle.native.partitioning.hash.nested.enabled=true, the type gate admits supported structs and arrays recursively, and maps on Spark 4.0+ subject to normalization/expression support. The existing config test expects a native exchange when enabled and none when disabled. The table and this checklist would classify a valid, covered native plan as an incorrect fallback decision. Please qualify the default-disabled nested-hash path and distinguish it from the complex-key restriction for range partitioning.
There was a problem hiding this comment.
Good catch, and the distinction matters more than I gave it credit for when I wrote that line. Range partitioning really is unconditionally primitive, because native cannot sort nested types, but hash partitioning is only primitive by default. Collapsing the two into one sentence would have had a reviewer flagging a correct native plan as a bug.
Split them in 67e734a, with the config default called out and a note that a map key additionally has to clear the separate CometMapSort check on the mapsort(...) Spark 4.0 inserts. The same unconditional claim was in native_shuffle.md and jvm_shuffle.md, so I corrected those too. The skill sends reviewers to those docs, so leaving them saying the opposite would have defeated the point.
| - [ ] The case ordering is right. `CometSelectionVector` must be matched **before** the general | ||
| `CometVector` case, or the selection is silently dropped and the exported batch has the wrong | ||
| rows. | ||
| - [ ] Selection vectors are applied where `scan.rs` expects them, in `ScanExec::get_next()` |
There was a problem hiding this comment.
Correctness
[P2] Replace obsolete selection-vector review checks
Neither CometSelectionVector nor CometDelegateVector exists in the current Spark/common/native sources. NativeUtil.exportBatch dispatches directly to CometVector or ConstantColumnVector, and ScanExec imports batches through get_next_batch/pull_next, with no get_next() selection-vector path. These mandatory checks would tell the reviewer to flag the current generic CometVector case as dropping rows and to inspect a nonexistent import hook. Please update this section and the preceding hierarchy to the current export/import model.
There was a problem hiding this comment.
Right, neither class exists anywhere in the tree. exportBatch has exactly two cases, CometVector and Spark's ConstantColumnVector, and ScanExec imports through pull_next with no selection hook at all, so both of those checks would have sent a reviewer chasing code that isn't there.
Replaced the section in 67e734a with the current model: the two export cases, the real CometVector hierarchy, and the point that a new subclass needs no new case as long as getValueVector returns an Arrow vector. The checks now cover who owns a materialized vector and the equal-value-count guard, which are the things that can actually go wrong on that path.
…and FFI export dispatch Addresses review feedback on the new review skills and the shuffle docs. - The local shuffle writer spills every partition into one file per task, not one file per partition. PartitionedSpill records per-partition byte ranges in write order and finish_partition copies them back in that order. - HashPartitioning keys are primitive-only only by default. The nested config admits structs and arrays recursively, and maps on Spark 4.0+. Range partitioning keys are primitive unconditionally, which is a different rule. - A shuffle output type that fails supportedSinkDataType does not retry as a regular Scan. convertToShuffleScan returns None and the Spark stage is kept. - NativeUtil.exportBatch has two cases, CometVector and ConstantColumnVector. CometSelectionVector and CometDelegateVector do not exist, and ScanExec imports through pull_next with no selection-vector hook. Also brings the two shuffle docs in line with the corrected key rules, notes that round robin is disabled by default and why, and adds the partitioning and maxBufferBytes configs to the native shuffle config table.
sunchao
left a comment
There was a problem hiding this comment.
Re-reviewed 67e734aad against the previous review at e2f054991. All four prior P2 findings are fixed: the direct-read fallback boundary, shared spill-file ownership, opt-in nested hash keys, and obsolete FFI selection-vector checks. I verified the corrections against unchanged implementation blobs and checked the added spill, export, and configuration guidance.
One new P2 remains in the added strict-mode range-key guidance. Current main already includes #5981, which admits scalar float and double range keys in strict mode. The verified CI merge bbb5c572b contains that behavior, while the new skill and contributor guide say those keys are rejected. The inline comment identifies both locations and the source/test evidence.
Prettier 3.9.6 and local link/anchor checks pass. Current CI has six successful checks and fifteen skipped checks. Its checkout includes all authored documentation changes, with newer main changes recorded separately. Spark/native runtime suites were skipped. No runtime or rendered-documentation validation is claimed, and the previously recorded local Mermaid dependency gap remains.
| and rejects collated strings, because Comet compares raw bytes. It also rejects float and double | ||
| when `spark.comet.exec.strictFloatingPoint` is on. |
There was a problem hiding this comment.
Correctness
[P2] Align the strict-mode range-key rule with current main
Could you update this restriction and the matching paragraph in native_shuffle.md for #5981? That change is already in the current CI merge bbb5c572b, whose main parent is c56cef77a. In that tree, supportedRangePartitioningDataType accepts scalar float and double keys even in strict mode, and CometSortOrder marks them compatible. The existing range-shuffle test expects a native exchange for both strict settings with SortOrder.allowIncompatible=false. The new rule would therefore tell reviewers to reject a supported native plan as soon as this PR lands on current main. The restriction still applies to nested types and collated strings, not scalar floating-point range keys.
There was a problem hiding this comment.
You're right, and I should have caught that the ground had moved under me while the PR was open — #5981 merged about twenty minutes before you posted this. supportedRangePartitioningDataType returns true for float and double with no strict-mode condition now, and CometSortOrder agrees, so both of those sentences would have a reviewer rejecting a plan that CometNativeShuffleSuite explicitly expects to be native.
I've filed #6045 to correct the skill and native_shuffle.md, rather than fix it here, for the same reason as my comment above.
|
Thanks for the reviews @sunchao. Given that these are development process skills and not production code, I'm going to address the remaining feedback in a separate PR. Hope that's ok. |
…6128) (#6222) * fix: let Comet memory pools overcommit on grow instead of panicking MemoryPool::grow must always succeed: DataFusion calls it for memory that already exists, such as a spilled batch the sort-merge join reads back. Both Comet pools implemented it as try_grow().unwrap(), so a partial grant from Spark panicked the task. Record the full amount, carry the ungranted part as overcommit, and repay it on shrink before releasing bytes to Spark, so Spark is still never handed back more than it granted (#1733). Closes #6127. * refactor: share Spark memory calls and overcommit between Comet pools Both pools duplicated the acquire/release JNI calls and wrapped them in the overcommit ledger separately, and had drifted: only one logged a failed acquire, and try_grow still used Spark's reply unclamped. Move the JNI calls, the partial-grant handback and the overcommit into one SparkMemory, generic over the Spark calls so it is tested with a fake instead of only the ledger. Skip the atomic write on release when there is no overcommit, and pass the fair pool its task attempt id for the warning. * fix: make unified try_grow repay overcommit and address review feedback - try_acquire asks Spark for outstanding overcommit on top of the request, so greedy_unified refuses try_grow while in debt instead of letting shortfalls compound unseen by Spark - document the invariant that keeps Spark from being over-released - replace the generic manager with Box<dyn SparkMemoryManager>, keep task_attempt_id in SparkMemory only, add pool test constructors - pool-level tests that fail with the old try_grow().unwrap() grow - saturating_add in both pools' grow - report overcommit in Display and try_grow errors - update the memory management contributor guide * test: drive the try_acquire repayment race and update the memory review skill Add a deterministic test that lands a concurrent release inside try_acquire's call to Spark, after it has read the debt and before it repays it, and a multi-threaded stress test of CometUnifiedMemoryPool that checks Spark is handed back exactly what it granted. SparkMemoryManager is now Send + Sync so the fake can be shared across threads, and the fake uses parking_lot so that one failed assertion does not poison the lock for every other thread. grow no longer panics when Spark grants less than asked, so the memory review skill now tells reviewers to treat a new grow call site as potential unbacked memory, and notes that only grow keeps a partial grant. * chore: drop the memory pools' redundant unsafe Send and Sync impls With SparkMemoryManager now Send + Sync, both pools are Send and Sync without them, and the compiler checks it: removing the bound makes the MemoryPool impls fail to compile. Since SparkMemory started holding a Box<dyn SparkMemoryManager>, these unsafe impls were the only thing satisfying MemoryPool's Send + Sync bound, and would have let a manager that is not thread-safe through. (cherry picked from commit d3406c6) Adapted for branch-1.0: - memory_pools/mod.rs: added the spark_memory module and passed task_attempt_id to CometFairMemoryPool::new by hand. branch-1.0 lacks #5494, which rewrote create_memory_pool on main, so the hunk did not apply. The pool changes themselves apply unchanged. - Dropped the changes to the memory management contributor guide and the review-comet-memory-pr skill. Neither file exists on branch-1.0; they come from #5933 and #6018.
#6228) * fix: check each fair_unified reservation against its own share Since the DataFusion 53 upgrade (#3629), CometFairMemoryPool::try_grow has compared the pool's total reserved bytes against pool_size / num_consumers. That caps the whole task at one consumer's share, and every new registration tightens the ceiling on the consumers already running. The comment that justified the change holds for shrink but not for try_grow. MemoryReservation::try_grow calls pool.try_grow before adding to its size, so reservation.size() is the size before the request. Check the reservation against its share again, and separately keep the pool's total within pool_size, since the shares alone do not bound it once a consumer registers late or one consumer holds sibling reservations. Update the tuning guide, the memory management guide and the memory PR review skill, which described the pool-wide comparison. Closes #5961. * fix: charge sibling fair_unified reservations to their consumer's share Checking each reservation against the share let an operator that holds several reservations under one consumer, such as a sort's streaming merge with one per batch, take other consumers' shares, bounded only by the pool total. Keep a running total for each consumer id instead, so all of a consumer's reservations count against one share. The pool no longer reads reservation.size(), so the check also stops depending on when DataFusion updates it. Note in the tuning guide that tasks can reserve more before they spill than they could in Comet 0.15.0 through 1.0.0. (cherry picked from commit e900844) Adapted for branch-1.0: - Dropped the changes to the memory management contributor guide and the review-comet-memory-pr skill. Neither file exists on branch-1.0; they come from #5933 and #6018. - tuning.md: the new paragraph on 0.15.0 through 1.0.0 no longer links to "Sizing the Overhead from the Memory Usage Log". That section and the native memory usage log it describes come from #6162, which is not on branch-1.0, so the link had no target.
Which issue does this PR close?
No issue. This is a contributor documentation and agent-tooling change.
Rationale for this change
The
review-comet-prskill had grown to cover expression review in considerable depth while sayingalmost nothing about the other subsystems a Comet PR commonly touches. A reviewer working on an FFI,
memory, or shuffle PR got the generic workflow plus a lot of expression-specific material that did not
apply, and no guidance on the invariants that actually matter in those areas.
The second gap was documentation drift. Nothing in the review workflow prompted a reviewer to ask
whether a PR had invalidated the contributor guide. The guide is full of class tables, file paths,
config defaults, and stated invariants, and a rename or a move silently turns a paragraph of it into a
lie. Follow-ups that are not filed as issues do not get done, so the check belongs in the review.
Writing that guidance meant reading the FFI, memory, and shuffle docs closely against the code, which
surfaced three defects in them. Those are fixed here rather than left for later, since they are exactly
the class of drift the new skills tell reviewers to catch.
What changes are included in this PR?
Review skills (
.ai/skills/)review-comet-pris now a generic entry point: PR metadata, existing comments, reading the diff,checks that apply to every PR (Spark compatibility, support levels, version shims, config conventions,
tests, CI), the documentation-freshness contract, the review bar, tone, and output format. Its first
step is a routing table mapping changed-file patterns to the area skill to load, and it notes that more
than one usually applies.
Four sibling skills, each naming the contributor guide docs to read before the diff and ending with the
specific claims in those docs that a PR can falsify:
review-comet-expression-pr, the expression material lifted out of the original skillreview-comet-ffi-pr, ownership per direction, release callbacks on error paths, no unwinding acrossextern "C",exportBatchcase ordering, whyAlignedArrowStreamReaderexists and its exit conditionreview-comet-memory-pr, naming which of the three budgets a change affects,try_growversusgrow,the fair pool's shared-total comparison, task-shared pool lifetime, and what evidence to ask for
review-comet-shuffle-pr, Murmur3 seed 42, why round robin is hash-based, block format changes needingwriter and reader and Celeborn together, and spill triggers per path
Contributor guide fixes
native_shuffle.mdpointed its Rust Side table atnative/core/src/execution/shuffle/, which no longerexists. The code is the
datafusion-comet-shufflecrate undernative/shuffle/, andcodec.rsis gone.Replaced with the current layout.
PartitionBufferandSpillFile, neither of which exists.Replaced with what
MultiPartitionShuffleRepartitioneractually holds, including thepinned_buffersdeduplication that keeps one allocation shared by many sliced batches from being charged per slice.
spark.comet.shuffle.compression.codecdefaults tozstd. It islz4.spark.comet.shuffle.directRead.enabledwas undocumented despite defaulting totrue. Added a DirectRead section covering how the path is selected during plan serialization, why the serialized protobuf
rather than the config is the source of truth for which slots are direct read, the block iterator's
header format and buffer lifetime contract, and the
ShuffleScanExecconstraints on JNI thread affinityand dictionary unpacking. Noted the branch in both read paths and expanded the corresponding aside in
ffi.md, since a direct-read slot does not cross the FFI boundary at all.How are these changes tested?
There are no code changes, so no test suites apply.
prettier --checkpasses on every file touched, which is what CI enforces over**/*.md.Every code reference added to the docs was verified against the source rather than carried over: the Rust
file layout and type names against
native/shuffle/src/, the config defaults againstCometConf.scala,the direct-read path against
CometSink.scala,CometExecRDD.scala,CometShuffleBlockIterator.java,and
shuffle_scan.rs. The Scala and Java paths already cited in both shuffle docs were confirmed to stillexist. The new cross-document anchors follow the
file.md#anchorstyle already used in the guide and arewithin the
myst_heading_anchors = 4depth configured inconf.py.