Skip to content

docs: split the PR review skill by area and correct the shuffle contributor docs - #6018

Merged
andygrove merged 4 commits into
apache:mainfrom
andygrove:review-skills-split
Sep 19, 2026
Merged

andygrove merged 4 commits into
apache:mainfrom
andygrove:review-skills-split

Conversation

@andygrove

Copy link
Copy Markdown
Member

Which issue does this PR close?

No issue. This is a contributor documentation and agent-tooling change.

Rationale for this change

The review-comet-pr skill had grown to cover expression review in considerable depth while saying
almost 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-pr is 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 skill
  • review-comet-ffi-pr, ownership per direction, release callbacks on error paths, no unwinding across
    extern "C", exportBatch case ordering, why AlignedArrowStreamReader exists and its exit condition
  • review-comet-memory-pr, naming which of the three budgets a change affects, try_grow versus grow,
    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 needing
    writer and reader and Celeborn together, and spill triggers per path

Contributor guide fixes

  • native_shuffle.md pointed its Rust Side table at native/core/src/execution/shuffle/, which no longer
    exists. The code is the datafusion-comet-shuffle crate under native/shuffle/, and codec.rs is gone.
    Replaced with the current layout.
  • The same doc's Memory Management section named PartitionBuffer and SpillFile, neither of which exists.
    Replaced with what MultiPartitionShuffleRepartitioner actually holds, including the pinned_buffers
    deduplication that keeps one allocation shared by many sliced batches from being charged per slice.
  • Both shuffle docs said spark.comet.shuffle.compression.codec defaults to zstd. It is lz4.
  • spark.comet.shuffle.directRead.enabled was undocumented despite defaulting to true. Added a Direct
    Read 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 ShuffleScanExec constraints on JNI thread affinity
    and 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 --check passes 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 against CometConf.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 still
exist. The new cross-document anchors follow the file.md#anchor style already used in the guide and are
within the myst_heading_anchors = 4 depth configured in conf.py.

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.
@github-actions github-actions Bot added the documentation Improvements or additions to documentation label Sep 18, 2026

@comphead comphead left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Thanks @andygrove this is a good direction to identify a scope first and get more detailed context for the modified part

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

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.

Comment on lines +208 to +210
`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.

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.

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.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

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.

Comment on lines +107 to +109
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.

@sunchao sunchao Sep 18, 2026 •

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.

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.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

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.

Comment on lines +50 to +51
Complex types are fully supported as **data** columns in both. The primitive-only restriction
applies to **partition keys** for `HashPartitioning` and `RangePartitioning` only.

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.

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.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

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.

Comment thread .ai/skills/review-comet-ffi-pr/SKILL.md Outdated
Comment on lines +98 to +101
- [ ] 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()`

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.

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.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

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

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.

Comment on lines +55 to +56
and rejects collated strings, because Comet compares raw bytes. It also rejects float and double
when `spark.comet.exec.strictFloatingPoint` is on.

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.

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.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

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.

@andygrove
andygrove added this pull request to the merge queue Sep 19, 2026
@andygrove

Copy link
Copy Markdown
Member Author

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.

Merged via the queue into apache:main with commit bce09a6 Sep 19, 2026
21 checks passed
@andygrove
andygrove deleted the review-skills-split branch September 19, 2026 16:58
andygrove added a commit that referenced this pull request Sep 25, 2026
…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.
andygrove added a commit that referenced this pull request Sep 29, 2026
#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.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

documentation Improvements or additions to documentation

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants