Skip to content

perf: stop holding the fair pool lock across blocking memory calls - #5613

Open
dwsmith1983 wants to merge 80 commits into
apache:mainfrom
dwsmith1983:perf/fair-pool-lock-free
Open

dwsmith1983 wants to merge 80 commits into
apache:mainfrom
dwsmith1983:perf/fair-pool-lock-free

Conversation

@dwsmith1983

@dwsmith1983 dwsmith1983 commented Sep 1, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

No dedicated issue. Adjacent to #5494, which made task-shared pools the norm and so widened the blast radius of this lock.

Rationale for this change

CometFairMemoryPool held its internal mutex across the JNI acquire and release calls into Spark's TaskMemoryManager. That call can block for a long time while Spark spills other consumers, and while it blocked, every other native thread sharing the task pool sat behind the lock, including plain releases that needed nothing from the JVM. An acquire parked inside Spark waits for another thread of the same task to release, and a release that has to take a lock the parked thread holds can never land. The sibling unified pool already avoids this.

What changes are included in this PR?

The pool builds on main's SparkMemory wrapper and its SparkMemoryManager trait, which arrived with the overcommit change, and adds two items to the wrapper. manager() exposes the raw Spark calls without the overcommit ledger, for bytes the pool never records: the anchor byte and a short grant it hands back itself. try_acquire_leaving_a_short_grant is the variant behind try_acquire that leaves a short grant with Spark instead of returning it, so the pool can keep those bytes charged until it has handed them back.

The fair limit is checked and the bytes are charged in one short locked step, since the limit couples the used total and the consumer count. The lock is dropped before any Spark call and the bookkeeping is settled after it returns. grow charges under the lock the same way, skips the limit, and carries whatever Spark declines as overcommit, as on main. try_grow rolls its charge back if Spark errors or panics inside the JNI frame. A short grant stays charged until Spark has it back, and with overcommit owed the request carries the debt, so a refused try_grow can receive more bytes than it asked for and hands all of them back. shrink checks its bounds under the lock, releases to Spark with the lock dropped, and only then takes the bytes off the pool's total. Going fully lock-free like the unified pool was rejected, since separate atomics would let a register or unregister slip between reading the count and committing the charge.

Concurrent grows still cannot jointly exceed pool_size divided by the consumer count, because each is charged before Spark is asked. Settling after the call opens two windows, both on the conservative side, and the memory guide documents them. A grow's charge is visible before Spark answers, so a second try_grow in that window is checked against a total that includes the first and can be refused where waiting would have let it through. A release is visible to Spark before the pool's total drops, so a grow that only fits once those bytes are free is refused at the fair limit rather than sent to Spark ahead of the release, while a grow that fits without them can still reach Spark first and come back short if the task is at its Spark share. Neither window admits a try_grow the limit would have refused, and a refusal is the ordinary error a spillable operator answers by spilling.

Spark's ExecutionMemoryPool removes a task's accounting entry the moment its balance reaches zero, and an acquire parked inside Spark reads that entry when it wakes, so a release that zeroes the balance under a parked acquire crashes the waiter. With the lock no longer held across the call that window is real. The pool therefore takes one byte from Spark on the first try_grow that passes the fair limit, or the first grow, before that call's own request, and keeps it until the pool drops. Creating the pool makes no JVM call, so a plan that never allocates natively never touches Spark's memory manager, and a grow the pool refuses locally never reaches Spark. Spark declines the byte with a zero grant for a task already at its share. The pool then runs without the anchor and each grow retries the byte as a request of its own until it is held, so the task spills on the short grant instead of failing, and the extra JNI call is paid only while the anchor is missing. On grow the anchor is best effort, and a try_grow whose anchor request fails rolls its charge back. The byte is never part of the pool's total or the overcommit, so it can never be repaid as debt. It is taken through acquireAnchor and releaseAnchor on CometTaskMemoryManager, which count it toward the task's balance and the consumer's usage but not toward the usage close checks, so a plan closing while a sibling still holds the pool does not log a one byte leak and a real leak is reported at its exact size. Until it is held, a sibling's release can still zero the balance under the grow's own parked request. That is the one window the anchor does not cover.

The registry is untouched. task_shared.rs and mod.rs differ from main only by doc sentences saying that creating a pool makes no JVM call. The memory guide gains three paragraphs on the anchor, the lock rule and the two windows.

How are these changes tested?

Unit tests in fair_pool.rs run the pool against two stand-ins for Spark. Main's FakeSpark covers the overcommit paths: the existing test of a grow past the fair limit now sees 99 of 100 bytes granted toward the grow because the anchor takes one first, and a new test pins that a refused try_grow with overcommit outstanding hands the whole grant back with neither the pool's total nor the overcommit moving. A stub of Spark 4.1.3's ExecutionMemoryPool covers the rest: the per-task entry lifecycle, the 1/N and 1/(2N) shares, the wait loop and notifyAll on release, with a bounded wait that fails a test instead of hanging, and gates that pause a call so a second thread can be interleaved at a precise point.

Those tests cover the lock discipline: a release does not stall behind a parked acquire on another thread, a second acquire does not queue behind a parked one inside the pool, and an eight-thread stress test asserts the total never exceeds the fair limit and nets to zero. They cover the rollback paths: short grants, acquire failure, a panic inside the call, and an over-shrink. They cover Spark's wait loop: an acquire parked below its minimum share woken by a full release, a late acquire surviving a release already on its way, and a sibling consumer freeing its last page under a parked acquire. They cover the anchor: taken on the first grow and returned once at drop, kept out of the usage a closing plan checks, a pool that never grows never calling Spark, a declined anchor leaving a usable pool that takes the byte once the share frees, two grows retrying at once keeping exactly one byte, and a pool from the real registry taking its anchor on its first grow. They cover the two windows: a grow racing a shrink is refused without a JVM call rather than admitted on bytes Spark still holds, a grow that fits without those bytes gets a short grant at the Spark share, a short-grant rollback keeps its bytes charged until Spark has them back, and an in-flight grow counts against the limit until it settles.

Microbenchmarks against a real Spark 4.1.3 TaskMemoryManager over an off-heap UnifiedMemoryManager, base and head in one binary, are in the discussion. They predate the lazy anchor and the SparkMemory merge, and the contended rows are the claim: at 4 and 8 native threads the release p99 falls from tens of microseconds to single digits because a release no longer waits behind an in-flight acquire, throughput improves 1.5x to 2.5x, grow p99 does not improve since that is Spark's own work, and errors and final balances are zero on every row. TPC-H SF10 on c2d5a28 with 2g off-heap, base and head alternated, spilled identically on every arm with no task failures, and the settled pass ran at a head to base wall-time ratio of 0.99. The same pair rerun on the merged head against main, base and head alternated, spilled the same bytes and the same number of times on every arm, with no failed tasks, no executor loss, and no overcommit, fair limit or short grant line in any executor log. The wall times of that rerun are not comparable because other builds shared the machine.

CometFairMemoryPool held its mutex across the JNI calls into Spark
task memory manager, which can block for seconds while Spark spills,
so every native thread sharing a task pool serialized behind whichever
thread was acquiring. The fairness check and the reservation are now
one short locked step, the blocking call runs unlocked, and the
reservation rolls back if the JVM fails to back it or the call panics.
Fairness semantics are unchanged: concurrent grows still cannot
jointly exceed pool_size divided by the consumer count.

The JNI boundary moved behind a small trait so the pool finally has
tests: fairness rejection, limit tightening on register, partial-grant
rollback, panic rollback, a blocking test that took ten seconds on the
old code and 30ms now, and an eight-thread stress test.
@dwsmith1983
dwsmith1983 force-pushed the perf/fair-pool-lock-free branch from 196d709 to 626bc39 Compare September 2, 2026 02:13

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

Reviewed 626bc395091f313ae93b8867c1cb1c846dc7fc53 against 8729f6e6adf7091e18a48670e790d4ba8fd41e51. I found one P2 in the new concurrent release path, detailed inline.

The focused Spark memory-pool component probe reproduced the missing-task exception. No full Comet native/JNI query suite ran. CI, CodeQL, and Delta Contrib Build Gate currently report action_required, with no test checks recorded.

For this performance change, please include matched BASE/HEAD microbenchmarks with one and multiple native threads, full and partial grants, and consumer registration changes. Report completed operations, grow/release latency, error counts, and final Rust/Spark balances. The parked stub test establishes lock behavior but does not measure production JNI throughput.

Comment on lines 176 to 177
state.used -= subtractive;
}

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P2] Keep waiting Spark acquisitions registered during a full release

Could you handle Spark's same-task wait/release contract before allowing this release to overlap an acquire? With fair_unified, two native consumers can now enter the same CometTaskMemoryManager concurrently. In a 100-unit executor pool, let another task hold 90 and this task hold 10. This task's next 10-unit grow waits below Spark's 1/(2N) minimum. Freeing its last 10 units on another native thread removes its entry from ExecutionMemoryPool.memoryForTask and wakes the grower. The grower then indexes the removed entry and throws NoSuchElementException: key not found. Spark's release bypasses the task monitor held by the waiting acquire, so that monitor does not prevent this interleaving. The Rust provisional reservation does not keep Spark's entry alive.

I reproduced the failure using unchanged Spark 3.5.9 pool source with only logging/annotation/memory-mode scaffolding. A scheduling control modeling the previous serialization completed after the other task freed its memory. The relevant map lifecycle is also present in 4.0.4 source. This was a component probe plus JNI source tracing, not a full Comet query reproduction. Please make full releases safe while grows are pending and add a regression that exercises Spark's memory manager.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Confirmed against the pool source, thanks for the repro. Blocking the release until the in-flight acquires drain turned out to deadlock in testing, since the parked acquire can be waiting on exactly the memory that release frees. So the fix defers instead: a release that would zero the task's balance while acquires are in flight frees n-1 bytes right away (that is what wakes the waiter) and holds the last byte, which the final completing acquire pays off. At most one byte is ever deferred and it always settles once the acquires finish. The test stub now models the entry lifecycle (created on acquire, removed at zero, a woken waiter fails if the entry is gone) and reproduced this crash before the fix. It also turned out one of our existing tests was exercising the same broken pattern.

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.

The new deferral fixes the already-pending case, but I can still reproduce the same missing-task failure through a late-arriving acquire on current head eb410e51.

At current lines 252-259, plan_release can see pending_acquires == 0, schedule the whole balance, and drop the state lock before release reaches Spark. A new try_grow can then increment pending_acquires and park in Spark while the old balance is still present. The already-planned release removes the task entry, and the waiter resumes with NoSuchElementException.

I reproduced this with the exact head's production state machine and a gated bridge: hold 10, plan and pause the full release, start and park a grow of 10, then resume the release. The waiter panics with key not found: task entry removed while acquire waited. paying_deferred fences only deferred payments, so could you also coordinate ordinary zeroing releases with newly starting acquires?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Also covered by db1f1bc. With the anchor there is no zeroing release left to coordinate: once any reservation exists the anchor is held, so a full release leaves the task at one byte and the entry survives.

Two related windows came out of review and are closed in the same commit. A grant that covers the request but not the extra byte is handed back as a short grant rather than running unanchored, and every acquire that starts while the anchor request is still parked carries its own extra byte, with the first full grant keeping it and later ones returning theirs. Your gated repro is pinned as late_acquire_survives_a_release_already_on_its_way, alongside a test for the concurrent in-flight case; both failed on eb410e5 with key not found and pass now. 30 loops each in debug and release are clean.

The earlier build-gate failure compared dylib sizes on a change confined to fair_pool.rs, so it looks like the size check rather than this branch.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

One more window closed in 353541d, found while probing the anchor bootstrap with a third task. Before the anchor lands, a short grant is handed back whole; if another task's entry disappears in that gap (raising this task's minimum share) and a sibling acquire of this task then parks, the rollback release zeroes the entry under it. A small bootstrap mutex now serializes anchor carriers from the bridge acquire through the rollback release. Releases never take it, so a parked carrier cannot starve anyone, and non-carriers only exist once the anchor is held. Pinned as short_grant_rollback_cannot_land_under_a_sibling_parked_in_spark, which failed with key not found before the change.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

The existing P2 missing-task-entry concern also remains reproducible after a declined anchor: a sibling releases 100 bytes, another task takes 90, and the real grow acquires 10 without an anchor. The next anchor retry parks. Releasing those 10 native bytes removes the entry and produces NoSuchElementException.

Confirmed. When Spark frees up between a declined anchor retry and the real request, the pool holds bytes from Spark without its anchor. A later retry can then park, and releasing those bytes took the task's balance to zero under it.

Now the first release that hands bytes back to Spark while the anchor is missing keeps one of them as the anchor. That covers a shrink and the rollback of a short grant. The new CometTaskMemoryManager.releaseKeepingAnchor releases all but one byte and moves that byte from used to the anchor count, so it is returned once at drop and close warnings stay quiet. The claim is made under the pool lock, so only one release keeps a byte, and a retry that lands afterwards hands its byte back as a duplicate, as before. The anchor stays outside the pool total and the overcommit ledger.

releasing the last bytes held without the anchor keeps a parked anchor retry's entry in CometTaskMemoryManagerSuite runs your sequence against a real UnifiedMemoryManager and TaskMemoryManager. With a plain releaseMemory(10) it fails with NoSuchElementException on Spark 3.5 and 4.1, and with the new call the parked retry is granted. release_while_unanchored_keeps_a_byte_under_a_parked_anchor_retry in fair_pool.rs drives the same sequence through the pool itself and fails with key not found without the change. The memory management guide now describes the release rule, and names the one window left, a JVM consumer of the task freeing its last bytes while the pool holds nothing from Spark. That window is tracked in #6224.

A release that would zero the JVM-side balance while other acquires are
still in flight frees all but one byte immediately and holds the last
byte until the in-flight acquires complete. Spark drops the task's
accounting entry when its balance hits zero, so a parked acquire waking
after that point indexes a missing entry and fails. Blocking the release
instead can deadlock because the parked acquire may be waiting for the
very memory the release frees. The stub task memory now models the
entry lifecycle so the regression is covered.
@dwsmith1983
dwsmith1983 requested a review from sunchao September 2, 2026 10:37
@dwsmith1983

Copy link
Copy Markdown
Contributor Author

@sunchao any more feedback here?

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

Reviewed eb410e5175417256f0223deaee1b122d760ae077 against ef62b46306e925bc51e7d7f29922c1870eb729e7. I found two actionable issues in the revised release/acquire protocol: the held-back byte can leave a valid Spark acquire waiting indefinitely, and an ordinary full release can still be overtaken by a newly starting acquire. The first is inline below. I added the second as a follow-up on the existing unresolved thread because it is the same missing-task failure under a different ordering.

Local validation: the focused native fair-pool suite passed 12/12 tests with 241 filtered out. An exact-head component probe reproduced the ordinary-release ordering failure. A probe using unmodified Spark 4.1.3 ExecutionMemoryPool reproduced both the missing-key failure and the one-byte wait. These are component probes, not a full Comet native/JNI query run.

Current CI has 63 successful checks, 9 skipped checks, and one failed Delta build gate. Its log fails the contrib-enabled-versus-default libcomet size invariant; this PR changes only fair_pool.rs.

.jvm_held
.checked_sub(bytes)
.expect("released more bytes than the JVM side holds");
if bytes > 0 && state.jvm_held == 0 && state.pending_acquires > 0 {

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.

[P1] Avoid retaining the byte a waiting acquire may need

Could you avoid withholding one byte until finish_acquire? This can deadlock the default off-heap pool. In a 1 GiB Spark execution pool, let another task hold 900 MiB, let this task hold 100 MiB, and let a second native consumer for this task request 100 MiB. Spark parks that request because this task is below its 1/(2N) minimum share. When the holder frees 100 MiB, this branch sends only 100 MiB - 1. On wake, Spark computes toGrant = 100 MiB - 1; because that is short of the request and curMem + toGrant = 100 MiB is still below 256 MiB, it waits again. The deferred byte is paid only by finish_acquire, which cannot run while this acquire is waiting.

I reproduced this against unmodified Spark 4.1.3 ExecutionMemoryPool: retaining one byte left the grower in WAITING, and freeing one additional byte from the other task let it complete. The new stub test misses this because it grants the full request after any release without reapplying Spark's free-memory and minimum-share checks. Could you preserve the task entry without withholding capacity needed by the waiter?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Fixed in db1f1bc. The deferral is gone: every release now goes to the JVM whole, so a parked acquire sees the full freed amount and Spark's minimum-share check passes. The task entry is kept alive by a permanent anchor instead. The first acquire asks for one extra byte and the pool holds it until it drops, so the balance never reaches zero mid task. That means the task stays in Spark's active-task set for the pool's lifetime and retains one byte, which the header comment now states.

Your scenario is pinned as min_share_wait_is_granted_after_the_holder_frees_its_memory_in_full. The stub now models ExecutionMemoryPool.acquireMemory from 4.1.3 with the per-task entry lifecycle, the 1/N and 1/(2N) shares, the wait loop, and notifyAll on release, with a bounded wait that fails the test instead of hanging. It failed on the previous head with the waiter timing out and passes now.

Spark's ExecutionMemoryPool removes a task's entry when its balance hits
zero, and an acquire parked inside Spark indexes that entry on wake. The
previous fix held one byte back from a zeroing release while acquires
were in flight and paid it back later. That withheld byte starved
Spark's minimum share check, which grants a parked request only when the
freed bytes cover it in full, so the waiter slept forever; and a release
planned before a late acquire parked could still remove the entry.

The pool now asks for one extra byte with every acquire that starts
before the anchor lands, keeps the first one until the pool drops, and
sends every release whole. The task's balance never returns to zero mid
task, so both failure modes are impossible by construction. The test
double now models ExecutionMemoryPool's entry lifecycle, share rules and
wait loop, and covers the minimum share scenario, the late acquire race
and the in flight anchor race.
…e task entry

Until the anchor lands, a short grant is handed back to Spark whole. In
the gap between that grant returning and its rollback release landing,
a sibling acquire of the same task can enter Spark and park if another
task has left the pool meanwhile, and the release then removes the
task's entry under it. Carriers now run one at a time from the bridge
acquire through the rollback release. Releases never take that lock,
so a parked carrier cannot hold up the release it waits for, and Spark
serializes a task's acquires anyway.
@dwsmith1983

Copy link
Copy Markdown
Contributor Author

@sunchao both of your findings are addressed on the branch head (353541d), with replies on each thread. Ready for another look whenever you have time.

@dwsmith1983
dwsmith1983 requested a review from sunchao September 4, 2026 04:23
@andygrove andygrove added enhancement New feature or request performance area:memory Memory pools, reservations, OOM handling labels Sep 6, 2026

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Summary

  • Prior state and problem: fair_unified held its mutex across blocking Spark memory calls, preventing concurrent native releases from progressing.
  • Design approach: Check and charge reservations under the mutex, call Spark unlocked, then settle accounting. Reuse SparkMemory for overcommit and retain a separate anchor byte.
  • Correctness / compatibility analysis: Found one new P1 deadlock, detailed inline. The existing P2 missing-task-entry concern also remains reproducible after a declined anchor: a sibling releases 100 bytes, another task takes 90, and the real grow acquires 10 without an anchor. The next anchor retry parks. Releasing those 10 native bytes removes the entry and produces NoSuchElementException. Both an exact-head native probe and Spark 4.1.3 reproduced this. It is recorded as an existing blocker rather than duplicated inline.
  • Key design decisions: Keeping anchor accounting separate fixes the reported one-byte close warning. Reusing the existing ledger and retaining the registry implementation keeps the production changes focused.
  • Implementation sketch: grow records overcommit, try_grow rolls back failed requests and returns short grants, and shrink settles after releasing to Spark.
  • Behavioral changes worth calling out: Provisional charges can cause additional conservative refusals. The anchor keeps the task active until pool destruction. The discussion contains performance measurements from earlier heads, but throughput was not independently measured here.
  • Suggested improvements: Fix the JVM diagnostic lock cycle and close the remaining unanchored native-release window. Add regressions using Spark's actual TaskMemoryManager, whose monitor behavior the native stub does not model.

Reviewed the entire nine-file diff from 88a1f48cc8a5c8f017016737c86f0912bbaefbe9 to bb273df5218dd42a3113b366011e465f158591cd. Read AGENTS.md and all supplied discussion and review threads. Routed skills: review-comet-pr, review-comet-memory-pr, and review-comet-ffi-pr.

Exact-head CI: Comet CI and CodeQL remain action_required. Only labeling passed. No exact-head CI build or test verdict is available.

Validation: cargo test -p datafusion-comet --lib memory_pools --offline --no-default-features passed all 50 tests. Bounded JVM probes used the unchanged head's Java bridge and real Spark 4.1.3 memory managers. The deadlock's control emulating the previous serialization completed. Relevant Spark sources were checked for 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0. No full Comet JVM/query suite, Spark SQL suite, or benchmark ran. Disposable test changes were removed and the checkout is clean.

// repays the overcommit. A short grant stays with Spark until this pool hands it
// back below, so the bytes can stay charged meanwhile.
let refusal = match std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
self.spark.try_acquire_leaving_a_short_grant(additional)

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.

[P1] [P1] Prevent showMemoryUsage from blocking short-grant rollback

Allowing these calls to overlap exposes a JVM lock cycle. CometTaskMemoryManager.acquireMemory releases Spark's task monitor after acquiring memory, then reacquires it in showMemoryUsage when the grant is short. In a 100-byte executor pool, let other tasks hold 82 and 1 bytes and this pool hold its anchor. A 30-byte request receives 16. Before its diagnostic reacquires the monitor, release the other task's 1 byte and start a second 10-byte request. The second request parks below its new 25-byte minimum while holding the task monitor. The first blocks in showMemoryUsage, so it cannot return to Rust and release the 16-byte short grant that would satisfy the second request. Both requests fit the local fair limit with two consumers.

Expected behavior is to return the short grant, spill, and let the second request proceed. Instead, the task can hang until an unrelated task frees memory. The previous fair-pool mutex excluded this interleaving. Could the bridge avoid reacquiring the task monitor after obtaining the grant, or cover acquisition and diagnostics with one monitor scope while keeping releases independent?

Evidence: Bounded reproduction: /tmp/pr5613-bb273df5-jvm-probe/ShortGrantProbe.scala, compiled with the unchanged head's CometTaskMemoryManager.java against Spark 4.1.3. A logging latch schedules the interleaving and is released before observing the failure. Output in result.log: head: first=BLOCKED, second=WAITING, held=16+anchor. Stacks identify TaskMemoryManager.showMemoryUsage:325 and ExecutionMemoryPool.acquireMemory:142. Releasing 50 bytes from the unrelated task rescues both threads. A control emulating the previous serialization completes without that release and ends with zero balance. Source inspection confirms the same monitor structure across all five checked Spark versions.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Could the bridge avoid reacquiring the task monitor after obtaining the grant, or cover acquisition and diagnostics with one monitor scope while keeping releases independent?

It now avoids it. CometTaskMemoryManager.acquireMemory no longer calls showMemoryUsage on a short grant. The warning still logs the task, the request, the grant, this manager's total and getMemoryConsumptionForThisTask. That last call takes only the memory manager's monitor, which a waiting acquire gives up in lock.wait(), so it cannot close the cycle.

One monitor scope would also work, since the monitor is reentrant. It would mean Comet locking on internal and relying on TaskMemoryManager using this as its lock. That holds from 3.4 through 4.1 but is not an API. It would also keep a thread that has bytes to hand back inside the task monitor for the whole dump. Dropping the dump avoids both. Releases were already independent. releaseExecutionMemory takes the memory manager's monitor, and on 4.x offHeapMemoryLock, but never the task's.

CometTaskMemoryManagerSuite now runs your interleaving against a real UnifiedMemoryManager with a 100 byte off-heap pool. Other tasks hold 82 bytes and 1 byte, and this task holds the anchor. A 30 byte request is granted 16, the 1 byte task leaves, and a 10 byte request waits in ExecutionMemoryPool. A TaskMemoryManager subclass holds the first thread right after its grant until the second has parked. With the old bridge the second acquire never finishes, and the test fails with the first thread BLOCKED and the second WAITING on both Spark 3.5 and 4.1. With the change the first thread hands its 16 bytes back and the second gets its 10.

Main makes the same call. There the fair pool's mutex keeps two native acquires of a task apart. greedy_unified has no such lock, so the same interleaving is open to it whenever two threads of one task acquire at once. The change is in the bridge, so it covers both pools.

The memory management guide now states the rule next to the one for getUsed and spill.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Summary

  • Prior state and problem: fair_unified held its mutex across blocking Spark memory calls, preventing concurrent native releases from progressing.
  • Design approach: Charge reservations under the mutex, call Spark unlocked, then settle accounting. A lazy anchor byte preserves Spark’s task entry once acquired.
  • Correctness / compatibility analysis: No additional introduced P1/P2 issues found within this review. The existing P1 diagnostic deadlock and P2 missing-task-entry concern remain reproducible. The JVM probes reproduced first=BLOCKED, second=WAITING, held=16+anchor and NoSuchElementException after an unanchored full release. Relevant locking and task-entry semantics match across Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0.
  • Key design decisions: Reusing SparkMemory and retaining the existing registry limits added complexity. Separate anchor accounting prevents false one-byte leak warnings without masking reservation leaks.
  • Implementation sketch: grow records overcommit, try_grow rolls back failures and returns short grants, and shrink settles after releasing to Spark.
  • Behavioral changes worth calling out: Provisional charges can cause additional conservative refusals. The anchor retains one byte and keeps the task active until pool destruction. Earlier-head performance measurements were reviewed but not independently repeated.
  • Suggested improvements: Resolve the two existing blockers and add regressions using Spark’s actual TaskMemoryManager, whose monitor behavior the native stub omits.

Reviewed the full nine-file diff from 88a1f48cc8a5c8f017016737c86f0912bbaefbe9 to bb273df5218dd42a3113b366011e465f158591cd, including surrounding code and supplied discussion. Confirmed the PR is not a draft. Read AGENTS.md. Routed skills: review-comet-pr, review-comet-memory-pr, and review-comet-ffi-pr.

Exact-head CI: Comet CI and CodeQL remain action_required. Only labeling passed. No CI build or test verdict is available for this head.

Validation: cargo test -p datafusion-comet --lib memory_pools --offline --no-default-features passed all 50 tests. Recompiled the unchanged head’s Java bridge and reran bounded probes against real Spark 4.1.3 memory managers. The control emulating the previous serialization completed without external memory release. No full Comet JVM/query suite, Spark SQL suite, or benchmark ran. The checkout remains clean.

Main now checks each fair_unified reservation against its own consumer's share
as well as the pool total. The pool keeps both totals and charges them in the
same locked step, the lock is still dropped before any Spark call, and settling
after the call updates both. The anchor stays out of both totals. Tests that
pinned the old shared comparison now race a sibling reservation of the same
consumer, and main's new tests count the anchor byte.
On a short grant the bridge logged Spark's per-consumer memory dump through
TaskMemoryManager.showMemoryUsage, which takes the task monitor. A second acquire
of the same task can park inside Spark while holding that monitor, waiting for
memory the first thread was about to hand back, so neither moved until another
task freed memory. The warn line with the requested, granted and current totals
stays and the dump is gone from the grant path. A test on a real
UnifiedMemoryManager drives that interleaving and hung before this change.
@dwsmith1983

Copy link
Copy Markdown
Contributor Author

Merged main, which brings in #6205. The pool keeps both of its checks: each consumer against its share, using the running total per consumer id, and the pool's total against pool_size. Both run under the lock in the same step that charges the bytes, and the charge now goes to the pool's total and the consumer's total together. Settling after the Spark call adjusts both. So a grant in flight, a shrink on its way to Spark, and a short grant being handed back count against the consumer's share as well as the pool until Spark has answered. shrink checks against the consumer's total, as main does. The anchor byte stays out of both totals.

A few tests here had pinned the old comparison, where one consumer's bytes refused another consumer at the share. They now race a sibling reservation of the same consumer instead. The test for a grow in flight registers the second consumer late, so that the pool total is what refuses it. The concurrency test checks each consumer against its share and the pool against its size. Main's new tests count the anchor in what Spark holds. The memory management guide's paragraphs on the lock and the windows it leaves now describe both checks.

…anchored

The anchor byte is taken lazily and Spark can decline it. After a decline the
pool could hold real bytes with no anchor, and a release of all of them while an
anchor retry was parked in Spark removed the task's entry under the retry, which
then failed with key not found. The first release that returns bytes to Spark
while the anchor is missing now keeps one of them as the anchor, through
CometTaskMemoryManager.releaseKeepingAnchor. That method releases to Spark before
it moves its counters, and a failed call clears the pool's claim so later grows
take the anchor as before.
@dwsmith1983

Copy link
Copy Markdown
Contributor Author

@sunchao both points from your last review are addressed. The showMemoryUsage lock cycle is fixed in 46ed0614c, and the release while the anchor is missing now keeps a byte as the anchor in 996e86f8f. Each has a regression test on a real TaskMemoryManager, and the replies in both threads have the details. The remaining window where a JVM consumer frees the task's last bytes is tracked in #6224. Could you take another look?

@andygrove

Copy link
Copy Markdown
Member

One more symptom of holding the fair pool's lock across the JNI call. While a fair-pool acquire is parked in Spark's ExecutionMemoryPool.acquireMemory, the executor's comet-memory-usage-log thread blocks in Native.getMemoryUsage on the same lock, because reading the pool's reservation takes it. So the periodic memory log goes quiet exactly when memory is contended. I saw it in a thread dump from a repro where the only Tokio worker was parked in that wait.

…ked acquire

The executor's memory usage logger reads the pool's reservation from its own
thread. While another task's acquire is parked in Spark, that read must return
at once, which only holds if no lock taken by reserved() is held across the
Spark call.
@dwsmith1983

Copy link
Copy Markdown
Contributor Author

While a fair-pool acquire is parked in Spark's ExecutionMemoryPool.acquireMemory, the executor's comet-memory-usage-log thread blocks in Native.getMemoryUsage on the same lock, because reading the pool's reservation takes it.

That is the lock this PR stops holding across the Spark call. On main, try_grow holds state across try_acquire, and reserved() needs the same lock, so the logger waits for as long as the acquire is parked. Here the limit check and the charge finish under the lock, and the anchor, acquire and release calls all run after it is dropped, so reserved() returns straight away with the in-flight bytes counted. The greedy pool keeps its usage in an atomic and TrackConsumersPool only forwards the call, so no other pool type has this wait.

2987a934f adds reserved_is_not_blocked_by_an_acquire_parked_in_spark, which parks an acquire in the test stub and expects reserved() to return within a second, including the parked task's bytes. With the lock held across the Spark call it times out. Could you approve a CI run on that head?

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Summary

  • Prior state and problem: fair_unified held its mutex across blocking Spark calls, delaying releases and memory-usage reads.
  • Design approach: Check and charge reservations under the mutex, call Spark unlocked, then settle accounting. Retain a separately accounted anchor byte.
  • Correctness / compatibility analysis: Found one P2 in anchor teardown during pool replacement, detailed below. The earlier diagnostic-deadlock and unanchored-release regressions pass. The JVM-consumer case tracked in #6224 is inherited. Checked relevant Spark memory-manager sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0.
  • Key design decisions: Reusing SparkMemory preserves overcommit handling, and separate anchor accounting avoids false leak warnings. The anchor's lifetime still needs coordination across replacement pools.
  • Implementation sketch: grow records provisional charges, try_grow rolls back failures and returns short grants, and shrink settles after releasing to Spark.
  • Behavioral changes worth calling out: Provisional charges can cause conservative refusals. The anchor consumes one byte and keeps the task active until pool destruction. Earlier-head performance measurements were reviewed but not independently repeated.
  • Suggested improvements: Coordinate anchor teardown with replacement-pool creation and add a regression covering that handoff through the task registry.

Reviewed the full nine-file diff from ce455f32d948355e073638e81009ecc3e5dea349 to f3dc4efd287958c9bab877af3cc6868fb15ad10a, including surrounding code and existing discussion. Confirmed the PR is not a draft. Read AGENTS.md. Routed skills: review-comet-pr, review-comet-memory-pr, and review-comet-ffi-pr.

Exact-head CI: Comet CI and CodeQL report action_required. Only labeling passed. No exact-head CI build or test verdict is available.

Validation: All 57 native memory-pool tests passed with --offline --no-default-features. Recompiled the head's Java bridge and ran all 6 CometTaskMemoryManagerSuite tests against Spark 4.1.3 successfully. Disposable native and JVM probes reproduced the teardown failure, while the ordered-release control passed. No full Maven build, JNI lifecycle/query suite, Spark SQL suite, or benchmark ran. Disposable source changes were removed and the checkout is clean.

return;
}
let released = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
self.spark.manager().release_anchor(ANCHOR_BYTES)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P2] Coordinate the final anchor release with replacement pools for the same task. TaskSharedMemoryPool permits replacement once its Weak expires, before its inner fair pool finishes dropping. In a 100-byte Spark pool, let other tasks hold 74 and 25 bytes and the dying pool hold only its anchor. Pause this release, create a replacement through acquire_task_shared_pool, and start its first grow. The replacement's anchor request parks. Resuming the old release removes the task's memoryForTask entry, so the replacement wakes with NoSuchElementException instead of acquiring memory or returning a spillable refusal. This can fail a task when final pool destruction overlaps creation of its next native plan. Base has no teardown anchor release. Could the anchor be shared across pool generations, or could replacement acquisition be coordinated with completed teardown without holding the global registry lock across JNI?

Evidence: The disposable native test review_replacement_pool_survives_the_old_pools_anchor_release exercised the actual head's registry and fair-pool implementation and failed with key not found: 0. A separate probe compiled with the unchanged head's Java bridge against real Spark 4.1.3 produced releaseOldAfterStart=true failure=java.util.NoSuchElementException: key not found: 0, at ExecutionMemoryPool.scala:115. Releasing the old anchor before starting the replacement produced failure=null. Logs: /tmp/comet-5613-review-xdjznr_c/native-probes.log and /tmp/comet-5613-review-xdjznr_c/jvm-probes.log.

mohitgurav20 pushed a commit to mohitgurav20/datafusion-comet that referenced this pull request Sep 28, 2026
apache#6269)

CometTaskMemoryManager.acquireMemory logged a warning, then called
TaskMemoryManager.showMemoryUsage, every time Spark granted less memory than a
native pool asked for. A partial grant is how a native operator learns to spill,
so a query that spills logged hundreds of them. Log the partial grant at DEBUG
instead. A refusal that fails the task already says in its error what Spark
granted and lists the pool's top consumers.

Drop the dump rather than moving it to DEBUG. showMemoryUsage takes the
TaskMemoryManager monitor. Another acquire of the same task can hold that
monitor while it waits inside Spark for memory, and this thread still holds its
partial grant. The task then hangs until some other task frees memory.
greedy_unified can reach this on main, because it calls Spark without a lock.
sunchao found the cycle while reviewing apache#5613.

Closes apache#6257.
@andygrove

Copy link
Copy Markdown
Member

This is a light fully automated review since there are so many PRs open.

native/core/src/execution/jni_api.rs:156-161 still explains why pool reservations are read outside the registry lock by saying that CometFairMemoryPool holds its own lock across the JNI call that acquires memory from Spark. After this PR that is no longer true. try_grow, grow and shrink all release state before any Spark call, reserved() no longer waits on a parked acquire, and the new paragraph in memory_management.md says the pool mutex is never held across a JNI call. Could that comment be updated in this PR so the two don't contradict each other? The rule itself can stay as a precaution. It just shouldn't cite the fair pool's lock as the reason.

The new helper in spark/src/test/scala/org/apache/spark/CometExecIteratorLifecycleSuite.scala:282 captures on classOf[CometExecIterator].getName. The test log4j config has no entry for that logger, so withLogAppender makes log4j create a non-additive config for org.apache.comet.CometExecIterator, and that config outlives the test. For the rest of that test JVM, every CometExecIterator log line, including the close warning and the memory usage log, reaches no appender. That includes any later capture on org.apache.comet. CometPluginsSuite.scala:338-341 describes this trap and listens on the package logger for that reason. Could this test do the same with Seq("org.apache.comet")? It already filters on the message text, so the assertions would not change.

andygrove added a commit that referenced this pull request Sep 28, 2026
#6269) (#6346)

CometTaskMemoryManager.acquireMemory logged a warning, then called
TaskMemoryManager.showMemoryUsage, every time Spark granted less memory than a
native pool asked for. A partial grant is how a native operator learns to spill,
so a query that spills logged hundreds of them. Log the partial grant at DEBUG
instead. A refusal that fails the task already says in its error what Spark
granted and lists the pool's top consumers.

Drop the dump rather than moving it to DEBUG. showMemoryUsage takes the
TaskMemoryManager monitor. Another acquire of the same task can hold that
monitor while it waits inside Spark for memory, and this thread still holds its
partial grant. The task then hangs until some other task frees memory.
greedy_unified can reach this on main, because it calls Spark without a lock.
sunchao found the cycle while reviewing #5613.

Closes #6257.

(cherry picked from commit e1d2c11)
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:ffi Arrow FFI / JNI boundary area:memory Memory pools, reservations, OOM handling enhancement New feature or request performance

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants