Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
21 commits
Select commit Hold shift + click to select a range
fda7bd6
feat: report JVM Arrow allocations to Spark's memory manager
andygrove Sep 17, 2026
0e3dfa7
refactor: reduce lock traffic in Arrow allocation accounting
andygrove Sep 17, 2026
55903fd
fix: exclude FFI-imported buffers from Arrow allocation accounting
andygrove Sep 17, 2026
2d531ec
docs: record that JVM Arrow allocations are now reported to Spark
andygrove Sep 17, 2026
4d7657d
style: keep the new config doc within the 100 character line limit
andygrove Sep 17, 2026
a2518d9
style: give the imported Arrow allocator an explicit result type
andygrove Sep 17, 2026
e58b177
fix: bind JVM Arrow accounting to a per-task allocator
andygrove Sep 17, 2026
094863f
test: cover release ownership, task completion and failed acquisitions
andygrove Sep 17, 2026
d591b8f
test: benchmark the cost of reporting Arrow allocations to Spark
andygrove Sep 17, 2026
054fe9a
fix: contain acquisition failures and stop double-charging FFI buffers
andygrove Sep 17, 2026
c2a542d
fix: measure a lost grant under the task memory manager monitor
andygrove Sep 18, 2026
a08908b
Merge remote-tracking branch 'apache/main' into arrow-allocator-accou…
andygrove Sep 18, 2026
9c9e147
fix: allocate codegen UDF output from the unaccounted root
andygrove Sep 18, 2026
d0e4210
Merge apache/main into arrow-allocator-accounting
andygrove Sep 26, 2026
fc4e6e5
fix: size the Arrow reservation from what the task allocator owns
andygrove Sep 26, 2026
e1da8de
docs: describe per-task Arrow accounting in the guides, skills and co…
andygrove Sep 26, 2026
1e13a9a
refactor: rename the accounting config to spark.comet.memory.jvmArrow…
andygrove Sep 26, 2026
6514a06
feat: reserve JVM Arrow memory before allocating, and refuse what Spa…
andygrove Sep 26, 2026
1f1e24e
feat: replace the accounting switch with spark.comet.legacy.unbounded…
andygrove Sep 26, 2026
62b1806
refactor: simplify the JVM Arrow listener and its tests
andygrove Sep 26, 2026
fd611d8
Merge remote-tracking branch 'apache/main' into arrow-allocator-accou…
andygrove Sep 26, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 7 additions & 3 deletions .ai/skills/review-comet-ffi-pr/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -135,9 +135,13 @@ a new subclass needs no new case as long as `getValueVector` returns an Arrow ve
## 5. Memory Accounting

An FFI change is usually also a memory change, and the two are easy to review separately and miss
the interaction. Imported JVM buffers come from `CometArrowAllocator`, which is an unbounded
`RootAllocator(Long.MaxValue)` that no budget sees. Exported native batches are usually no longer
pool-reserved by the time the JVM receives them, but they stay resident until the JVM closes them.
the interaction. The JVM buffers native imports belong to a child of `CometArrowAllocator`, an
unbounded `RootAllocator(Long.MaxValue)` that no budget sees, and a native operator that retains
them reserves them itself. Batches the JVM owns within a task come from `CometTaskArrowAllocator`,
whose listener charges them to Spark and refuses what Spark cannot cover, so a new path that exports
them without going through `CometArrowStream`, which moves the charge off the task, charges them on
both sides. Exported native batches are usually no longer pool-reserved by the time the JVM receives
them, but they stay resident until the JVM closes them.

If the PR changes how long either side holds a batch, use `review-comet-memory-pr` as well.

Expand Down
22 changes: 17 additions & 5 deletions .ai/skills/review-comet-memory-pr/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -158,15 +158,27 @@ memory_limit = spark.memory.offHeap.size * spark.comet.exec.memoryPool.fraction
`spark.executor.memoryOverhead` is the only real slack in the container, and JVM non-heap
usage already consumes much of it.

## 6. Unaccounted Allocators
## 6. JVM-Side Allocators

Two allocators sit outside Comet's pool entirely, and they are easy to confuse:

- **`CometArrowAllocator`** (`spark/src/main/scala/org/apache/comet/package.scala`) is a
process-wide `RootAllocator(Long.MaxValue)`. It is unbounded and no budget sees it. Child
allocators are cut from it for FFI stream export, broadcast coalescing, and
`CometSparkToColumnarExec`. A PR that adds a child allocator, or makes an existing one hold more,
is adding uncounted container RSS. Say so even if the volume is small.
process-wide `RootAllocator(Long.MaxValue)` with no allocation listener. It is unbounded. Buffers
allocated from it, or from a child cut with the three-argument `newChildAllocator`, are seen by no
budget: FFI stream export, `NativeUtil`, the JVM UDF result, and the import path's
`CometArrowImportAllocator`. JVM-owned allocations inside a task go through
`CometTaskArrowAllocator.forCurrentTask()` instead, a per-task child whose
`CometArrowAllocationListener` charges what it owns to that task's `TaskMemoryManager` and refuses
an allocation Spark cannot cover, so it fails with Arrow's `OutOfMemoryException`. A PR that
allocates from the root inside a task, or makes a root-backed allocator hold more, is adding
uncounted container RSS. Say so even if the volume is small. A PR that allocates something native
will retain from the task allocator is charging it twice, unless the batch reaches native through
`CometArrowStream`, whose reader moves the charge off the task allocator. A new call site on the
task allocator also has to cope with that refusal: nothing spills these buffers, so it fails the
task. Check that the listener stays without a lock in `getUsed` and `spill`, that only
`onPreAllocation` throws, and that it sizes the reservation from what the allocator owns rather
than from a sum of its callbacks, because an ownership transfer between allocators calls no
listener.
- **JVM shuffle pages** go through `CometShuffleMemoryAllocator.getInstance`, which returns
`CometUnifiedShuffleMemoryAllocator`, an ordinary Spark `MemoryConsumer` drawing from
`spark.memory.offHeap.size`. These **are** arbitrated by Spark against its other consumers. Do
Expand Down
1 change: 1 addition & 0 deletions .github/workflows/pr_build_linux.yml
Original file line number Diff line number Diff line change
Expand Up @@ -543,6 +543,7 @@ jobs:
org.apache.spark.CometPluginsSuite
org.apache.spark.CometRuntimeShutdownSuite
org.apache.spark.CometTaskMemoryManagerSuite
org.apache.spark.comet.CometArrowAllocationListenerSuite
org.apache.spark.CometExecIteratorLifecycleSuite
org.apache.spark.CometPluginsDefaultSuite
org.apache.spark.CometPluginsMemoryOverheadWarningSuite
Expand Down
1 change: 1 addition & 0 deletions .github/workflows/pr_build_macos.yml
Original file line number Diff line number Diff line change
Expand Up @@ -191,6 +191,7 @@ jobs:
org.apache.spark.CometPluginsSuite
org.apache.spark.CometRuntimeShutdownSuite
org.apache.spark.CometTaskMemoryManagerSuite
org.apache.spark.comet.CometArrowAllocationListenerSuite
org.apache.spark.CometExecIteratorLifecycleSuite
org.apache.spark.CometPluginsDefaultSuite
org.apache.spark.CometPluginsMemoryOverheadWarningSuite
Expand Down
Loading
Loading