Repository navigation
perf: stop holding the fair pool lock across blocking memory calls - #5613
dwsmith1983 wants to merge 112 commits into
Conversation
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.
196d709 to
626bc39
Compare
sunchao
left a comment
There was a problem hiding this comment.
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.
| state.used -= subtractive; | ||
| } |
There was a problem hiding this comment.
[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.
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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?
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
|
@sunchao any more feedback here? |
sunchao
left a comment
There was a problem hiding this comment.
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 { |
There was a problem hiding this comment.
[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?
There was a problem hiding this comment.
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.
The workaround for apache#6254 lived in both Comet pools. The fair pool sent its three refusals through a check, and the greedy pool tracked each final aggregate across its reservations. apache#5613 reworks the fair pool's try_grow, and apache#6583 has to take the workaround out again, so both would have had to rework the pools. SpillReplayPool, in spill_replay.rs, now wraps the tracked Comet pool instead. It keeps a total for each final aggregate's consumer. When the pool refuses a request from one while another of its reservations holds memory, it calls the pool's grow, which skips the fair pool's limits and carries what Spark doesn't grant as overcommit. fair_pool.rs and unified_pool.rs are back to main's versions, and overcommit() looks through the wrapper. The Rust tests run against both pool types as createPlan builds them. They now reach the fair pool's pool-limit refusal too, and check that a failed Spark call during the replay is returned rather than recorded. The guide adds the wrapper to the pool stack and keeps a short section. The DataFusion 55.1 facts that an upgrade has to re-check stay in spill_replay.rs, which names apache#6583 next to the removal condition.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: The fair pool held its mutex across blocking Spark calls, delaying concurrent releases and usage reads.
- Design approach: Charge reservations under the mutex, call Spark unlocked, then settle accounting. A lazy anchor preserves Spark’s task entry, and replacement pools wait for teardown.
- Correctness / compatibility analysis: The previously discussed exact-fit concern remains reproducible as a P2 wait regression at
native/core/src/execution/memory_pools/fair_pool.rs:159. In a 100-byte Spark pool with another task holding 90, the firsttry_grow(10)succeeds on the base. Head first acquires its anchor, leaving nine bytes. Spark then parks the 10-byte request because1 + 9 < 25, its minimum share. The allocation cannot return success or a spillable refusal until unrelated memory is released. The native head probe reproduced this, the base control passed, and real Spark 4.1.3 reproduced the sameWAITINGstate. Relevant semantics match Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. This is recorded as an existing concern, without duplicating it inline. - Key design decisions: Reusing
SparkMemorypreserves the overcommit ledger. Keeping the anchor separate avoids treating it as debt. The teardown guard orders pool generations without holding the registry lock across JNI. - Implementation sketch:
growrecords provisional charges,try_growrolls back failures and returns short grants, andshrinksettles after Spark accepts the release. - Behavioral changes worth calling out: Compared with
branch-1.1, intended changes permit concurrent releases, retain one anchor byte, allow conservative refusals during settlement, and make replacement creation wait for teardown. Earlier contention benchmarks and spill comparisons do not cover the reproduced first-allocation boundary. - Suggested improvements: Revisit anchor establishment so an otherwise satisfiable first allocation cannot park solely for the bookkeeping byte. Add a regression without pre-seeding the anchor, which the current exact-fit test does.
Reviewed the complete ten-file diff from 965c8bbe289ff850b614c8c3833ed7802134115b to fb9246632760cff7b652f9db7e329c280c3b2622, including surrounding code and existing discussions. Confirmed non-draft status. Read AGENTS.md and applied review-comet-pr, review-comet-memory-pr, and review-comet-ffi-pr. Beyond the existing concern above, no introduced P1/P2 issues found within this review.
Exact-head CI: Comet CI and CodeQL remain action_required. Labeling passed. No CI build/test verdict is available.
Validation: All 66 native memory-pool tests passed with --offline --no-default-features. Recompiled the head’s Java bridge and Scala suite, and all eight JVM memory-manager tests passed against Spark 4.1.3. The additional bounded native regression failed on head and passed against base source. No full Maven build, JNI query/lifecycle suite, Spark SQL suite, or benchmark ran. Disposable source changes were removed and the checkout is clean.
…grow The anchor byte is taken just before the request of the grow that takes it, so a request for exactly the memory Spark has free comes up one byte short. Below its minimum share the task then waits inside Spark until another task releases memory, and at or above it a try_grow is refused one byte short while a grow carries one byte of overcommit. Add tests for both and describe the trade-off in the memory guide.
|
On the exact-fit wait: it is real, and we are keeping the anchor for it. The anchor is what stops a release from removing the task's entry under a parked acquire, which fails that acquire, and the extra wait only happens when the request would take exactly the memory Spark has free and ends at the next release. 518dc62 pins it: below the minimum share the request parks until another task releases memory and then gets its full grant, and at the share |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem:
fair_unifiedheld its mutex across blocking Spark memory calls, delaying concurrent releases and usage reads. - Design approach: Charge reservations under the mutex, call Spark unlocked, then settle accounting. A lazy anchor preserves Spark’s task entry, and replacement pools wait for teardown.
- Correctness / compatibility analysis: The previously reported P2 exact-fit wait remains at
native/core/src/execution/memory_pools/fair_pool.rs:159. With a 100-byte Spark pool and another task holding 90, the base grantstry_grow(10)immediately. Head takes its anchor first, leaving nine bytes. Spark parks the request because1 + 9 < 25, its minimum share. The request cannot succeed or return a spillable refusal until unrelated memory is released. The new native test confirms this behavior, and a bounded probe against real Spark 4.1.3 reproducedWAITINGwhile the control without an anchor completed. Relevant semantics match Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. Earlier diagnostic-deadlock, native-release and pool-replacement regressions pass their focused tests. - Key design decisions: Atomic admission checks preserve consumer and pool limits. Reusing
SparkMemoryavoids duplicating the overcommit ledger. The teardown guard orders pool generations without holding the registry lock across JNI. - Implementation sketch:
growrecords provisional charges and overcommit.try_growrolls back failures and returns short grants.shrinksettles after Spark accepts the release. - Behavioral changes worth calling out: Compared with
branch-1.1, intended changes permit concurrent releases, retain an anchor byte, allow conservative refusals during settlement, and make replacement creation wait for teardown. Earlier-head contention benchmarks and spill comparisons were reviewed but not independently repeated. They do not resolve the exact-fit regression. - Suggested improvements: Rework anchor establishment so an otherwise satisfiable allocation cannot park solely for the bookkeeping byte. Documenting and testing the wait preserves the regression rather than resolving it.
Reviewed the full ten-file diff from ba08acd815d5fd15f75ce9f314715778b9161cf1 to 9bacfff0de2ac77b44b87e93a67c46018245af41, including surrounding code and existing reviews, comments and threads. Confirmed non-draft status. Read AGENTS.md and applied review-comet-pr, review-comet-memory-pr and review-comet-ffi-pr. Beyond the existing concern above, no introduced P1/P2 issues found within this review.
Exact-head CI: Comet CI, CodeQL and the title check report action_required. Labeling passed. No CI build/test verdict is available.
Validation: All 69 native memory-pool tests passed with --offline --no-default-features. Recompiled the head’s Java bridge and Scala suite, then passed all eight memory-manager tests against Spark 4.1.3. No full Maven build, JNI lifecycle/query suite, Spark SQL suite or benchmark ran. Other Spark versions were checked through source comparison. Project files remain unchanged.
Review state: Request changes for the unresolved existing P2.
… memory share (apache#6544) * fix: let a spilled final aggregate read its spill files back past its memory share DataFusion 55's FinalHashAggregateStream replays its merged spill files through a stream that can't spill, so a refused memory request there fails the task. The merge's read buffers are a sibling reservation of the same consumer and take as many spill files as fit, so the replay often finds the consumer's share already taken (apache#6254). Both Comet pools now record that request as overcommit instead of refusing it. They recognize it as a request from a FinalHashAggregateStream consumer while another of its reservations holds memory, which in DataFusion 55.1 happens only during the replay. The replay emits its finished groups after every batch, so the overcommit stays around one batch of groups, and releases repay it first. Refusals while the aggregate reads its input, and while the merge picks its files, are unchanged. Remove this once Comet's DataFusion includes apache/datafusion#25383. * fix: let an ordered final aggregate read its spill files back too DataFusion runs a final aggregate whose input is sorted on some of its grouping keys as an OrderedFinalAggregateStream. It merges and replays its spill files the same way FinalHashAggregateStream does, so its replay failed the task the same way. Treat its consumer as a final aggregate too. Explain in spill_replay.rs why recording the replay's request is safe: the replay asks for memory only after it has aggregated a batch, so the memory already exists, as it does for a grow. Describe the exception in the memory management guide, which said the fair pool always refuses a request that fails its local checks. * docs: tighten the spill replay comments and guide The fair pool is the only one with a share and a pool total to skip, so say so. The greedy pool takes its tracking lock for other consumers when Spark refuses them, not never. Point whoever removes the workaround at the tests that show whether the replay still needs it. * test: share the spill replay tests between the pools fair_pool.rs and unified_pool.rs each had a copy of the same two spill replay scenarios. Run them once, in spill_replay.rs, against both pool types as createPlan builds them, which also checks that the wrappers pass register through to the greedy pool's tracking. Keep a greedy pool test for its per-consumer map, and drop a test import that the pool's own imports now cover. Share the final aggregate spill checks between the two apache#6254 tests in CometAggregateSuite, and use checkCometAnswer, which collects once and labels the Comet answer as Comet's. * refactor: move the spill replay workaround into a pool wrapper The workaround for apache#6254 lived in both Comet pools. The fair pool sent its three refusals through a check, and the greedy pool tracked each final aggregate across its reservations. apache#5613 reworks the fair pool's try_grow, and apache#6583 has to take the workaround out again, so both would have had to rework the pools. SpillReplayPool, in spill_replay.rs, now wraps the tracked Comet pool instead. It keeps a total for each final aggregate's consumer. When the pool refuses a request from one while another of its reservations holds memory, it calls the pool's grow, which skips the fair pool's limits and carries what Spark doesn't grant as overcommit. fair_pool.rs and unified_pool.rs are back to main's versions, and overcommit() looks through the wrapper. The Rust tests run against both pool types as createPlan builds them. They now reach the fair pool's pool-limit refusal too, and check that a failed Spark call during the replay is returned rather than recorded. The guide adds the wrapper to the pool stack and keeps a short section. The DataFusion 55.1 facts that an upgrade has to re-check stay in spill_replay.rs, which names apache#6583 next to the removal condition.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem:
fair_unifiedheld its mutex across blocking Spark memory calls, delaying concurrent releases and usage reads. - Design approach: Check and charge reservations under the mutex, call Spark unlocked, then settle accounting. A lazy anchor preserves Spark’s task entry, and replacement pools wait for completed teardown.
- Correctness / compatibility analysis: The existing P2 exact-fit wait remains reproducible at
native/core/src/execution/memory_pools/fair_pool.rs:159. With a 100-byte Spark pool and another task holding 90, the base’s unanchored request for 10 succeeds immediately. Head first takes its anchor, leaving nine bytes. Spark parks the request because1 + 9 < 25, its minimum share. It cannot succeed or return a spillable refusal until another release supplies memory. The head’s native test confirms this, and a bounded probe using the head’s Java bridge with real Spark 4.1.3 reproducedWAITING. The control without an anchor completed immediately. Relevant semantics match Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. Earlier diagnostic-deadlock, native-release and replacement-pool regressions pass. Beyond this existing concern, no introduced P1/P2 issues found within this review. - Key design decisions: Atomic admission checks preserve consumer and pool limits. Reusing
SparkMemoryavoids duplicating overcommit accounting. The teardown guard coordinates pool generations without holding the registry lock across JNI. - Implementation sketch:
growrecords provisional charges and overcommit.try_growrolls back failed requests and returns short grants.shrinksettles after Spark accepts the release. - Behavioral changes worth calling out: Compared with
branch-1.1, intended PR changes permit concurrent releases, retain one anchor byte, allow conservative refusals during settlement, and make replacement creation wait for teardown. Earlier-head contention benchmarks and spill comparisons were reviewed but not independently repeated. They do not resolve the exact-fit wait. - Suggested improvements: Rework anchor establishment so its bookkeeping byte cannot park an otherwise satisfiable allocation, while preserving task-entry safety. Documenting and testing the wait does not resolve the existing regression.
Reviewed the complete ten-file diff from 4ea367aab4af430fce5dafa84351bedd288ec074 to 05a5644beff8692a42b467d2a0706ff3cf51ee27, including surrounding code and existing reviews, issue comments, inline comments and threads. Confirmed non-draft status. Read AGENTS.md and applied review-comet-pr, review-comet-memory-pr and review-comet-ffi-pr.
Exact-head CI: Comet CI and CodeQL report action_required. Labeling passed. No CI build/test verdict is available.
Validation: All 69 native memory-pool tests passed with cargo test -p datafusion-comet --lib memory_pools --offline --no-default-features. Recompiled the head’s Java bridge and Scala suite, then passed all eight memory-manager tests against Spark 4.1.3. The bounded exact-fit probe also completed successfully. No full Maven build, JNI lifecycle/query suite, Spark SQL suite or benchmark ran. Other Spark versions were checked through source comparison. Project files remain unchanged.
Review state: Request changes for the unresolved existing P2.
The spill_replay tests from main drive the fair pool through FakeSpark and expected Spark to hold exactly the reserved bytes. This branch's fair pool also holds its one-byte anchor until the pool drops, so the fair_unified expectations add ANCHOR_BYTES and the first test now drops the pool before checking that Spark holds nothing.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem:
fair_unifiedheld its mutex across blocking Spark memory calls, delaying concurrent releases and memory-usage reads. - Design approach: Check and charge reservations under the mutex, call Spark unlocked, then settle accounting. A lazy anchor preserves Spark’s task entry, and replacement pools wait for completed teardown.
- Correctness / compatibility analysis: The existing P2 exact-fit wait remains reproducible at
native/core/src/execution/memory_pools/fair_pool.rs:159. With a 100-byte Spark pool and another task holding 90, the base can granttry_grow(10)immediately. Head takes its anchor first, leaving nine bytes. Spark parks the request because1 + 9 < 25, its minimum share. The request cannot succeed or return a spillable refusal until another release supplies memory. The head’s native test confirms this behavior. A bounded probe using the head’s Java bridge and real Spark 4.1.3 reproducedWAITING, while the control without an anchor completed immediately. Releasing one byte rescued the waiting request. Earlier diagnostic-deadlock, native-release and replacement-pool regressions pass. Beyond this existing concern, no introduced P1/P2 issues found within this review. - Key design decisions: Admission checks and provisional charges remain atomic. Reusing
SparkMemoryavoids duplicating overcommit accounting. The teardown guard coordinates pool generations without holding the registry lock across JNI. - Implementation sketch:
growrecords provisional charges and overcommit.try_growrolls back failures and returns short grants.shrinksettles after Spark accepts the release, and Java usage accounting follows successful release. - Behavioral changes worth calling out: Compared with
branch-1.1, intended PR changes permit concurrent releases, retain an anchor byte, allow conservative refusals during settlement, and make replacement creation wait for teardown. Earlier-head contention benchmarks and spill comparisons were reviewed but not independently repeated. They do not resolve the exact-fit wait. - Suggested improvements: Resolve the existing P2 by establishing task-entry protection without parking an otherwise satisfiable allocation solely for the anchor byte. Documenting and testing the wait does not eliminate the regression.
Reviewed the full eleven-file diff from 9dc8c3ca89962506078831f7064a25c75da83883 to 8aa1c106f69f7cde500b08024e72bc3c0f773243, including surrounding code and existing reviews, issue comments, inline comments and threads. Confirmed non-draft status. Read AGENTS.md and applied review-comet-pr, review-comet-memory-pr and review-comet-ffi-pr.
Exact-head CI: Comet CI and CodeQL report action_required. Labeling passed. No CI build/test verdict is available.
Validation: All 76 native memory-pool tests passed with cargo test -p datafusion-comet --lib memory_pools --offline --no-default-features. Recompiled the head’s Java bridge and Scala suite, then passed all eight memory-manager tests against Spark 4.1.3. Checked relevant Spark locking and task-entry sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. No full Maven build, JNI lifecycle/query suite, Spark SQL suite or benchmark ran. Project files remain unchanged.
Review state: Request changes for the unresolved existing P2.
… memory share (#6544) (#6713) * fix: let a spilled final aggregate read its spill files back past its memory share DataFusion 55's FinalHashAggregateStream replays its merged spill files through a stream that can't spill, so a refused memory request there fails the task. The merge's read buffers are a sibling reservation of the same consumer and take as many spill files as fit, so the replay often finds the consumer's share already taken (#6254). Both Comet pools now record that request as overcommit instead of refusing it. They recognize it as a request from a FinalHashAggregateStream consumer while another of its reservations holds memory, which in DataFusion 55.1 happens only during the replay. The replay emits its finished groups after every batch, so the overcommit stays around one batch of groups, and releases repay it first. Refusals while the aggregate reads its input, and while the merge picks its files, are unchanged. Remove this once Comet's DataFusion includes apache/datafusion#25383. * fix: let an ordered final aggregate read its spill files back too DataFusion runs a final aggregate whose input is sorted on some of its grouping keys as an OrderedFinalAggregateStream. It merges and replays its spill files the same way FinalHashAggregateStream does, so its replay failed the task the same way. Treat its consumer as a final aggregate too. Explain in spill_replay.rs why recording the replay's request is safe: the replay asks for memory only after it has aggregated a batch, so the memory already exists, as it does for a grow. Describe the exception in the memory management guide, which said the fair pool always refuses a request that fails its local checks. * docs: tighten the spill replay comments and guide The fair pool is the only one with a share and a pool total to skip, so say so. The greedy pool takes its tracking lock for other consumers when Spark refuses them, not never. Point whoever removes the workaround at the tests that show whether the replay still needs it. * test: share the spill replay tests between the pools fair_pool.rs and unified_pool.rs each had a copy of the same two spill replay scenarios. Run them once, in spill_replay.rs, against both pool types as createPlan builds them, which also checks that the wrappers pass register through to the greedy pool's tracking. Keep a greedy pool test for its per-consumer map, and drop a test import that the pool's own imports now cover. Share the final aggregate spill checks between the two #6254 tests in CometAggregateSuite, and use checkCometAnswer, which collects once and labels the Comet answer as Comet's. * refactor: move the spill replay workaround into a pool wrapper The workaround for #6254 lived in both Comet pools. The fair pool sent its three refusals through a check, and the greedy pool tracked each final aggregate across its reservations. #5613 reworks the fair pool's try_grow, and #6583 has to take the workaround out again, so both would have had to rework the pools. SpillReplayPool, in spill_replay.rs, now wraps the tracked Comet pool instead. It keeps a total for each final aggregate's consumer. When the pool refuses a request from one while another of its reservations holds memory, it calls the pool's grow, which skips the fair pool's limits and carries what Spark doesn't grant as overcommit. fair_pool.rs and unified_pool.rs are back to main's versions, and overcommit() looks through the wrapper. The Rust tests run against both pool types as createPlan builds them. They now reach the fair pool's pool-limit refusal too, and check that a failed Spark call during the replay is returned rather than recorded. The guide adds the wrapper to the pool stack and keeps a short section. The DataFusion 55.1 facts that an upgrade has to re-check stay in spill_replay.rs, which names #6583 next to the removal condition. (cherry picked from commit ba7c892) Adapted for branch-1.1: - The tests build the pools over a fake Spark through create_pool, a seam that #6271 added to memory_pools. Only that seam is ported, not #6271's memory usage log change, which is not on branch-1.1: create_memory_pool builds the pools through create_pool, with_spark is pub(super), and the pools' new constructors, which nothing calls any more, are removed. - overcommit() and unwrap_task_shared(), which only that log reads, are not ported, so nothing looks through the wrapper. spill_replay.rs drops unwrap_spill_replay, and its tests drop their overcommit(&pool) assertions, which restate pool.reserved() minus what the fake Spark holds. The test that shrinks the replay checks pool.reserved() instead. - memory_management.md: the paragraph about the task's shared CometTaskMemoryManager, from #6261, is not on branch-1.1, so the new section follows the task-shared pools section directly. - CometAggregateSuite: only its imports differ. This branch does not import Column, CometBaseAggregateExec or CometSortAggregateExec.
The spill_replay tests from main drive the fair pool through FakeSpark and expected Spark to hold exactly the reserved bytes. This branch's fair pool also holds its one-byte anchor until the pool drops, so the fair_unified expectations add ANCHOR_BYTES and the first test now drops the pool before checking that Spark holds nothing.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem:
fair_unifiedheld its mutex across blocking Spark memory calls, delaying concurrent releases and memory-usage reads. - Design approach: Check and charge reservations under the mutex, call Spark unlocked, then settle accounting. A lazy anchor preserves Spark’s task entry, and replacement pools wait for completed teardown.
- Correctness / compatibility analysis: The existing P2 exact-fit wait remains reproducible at
native/core/src/execution/memory_pools/fair_pool.rs:159. With a 100-byte Spark pool and another task holding 90, the base’s unanchored 10-byte request succeeds immediately. Head takes its anchor first, leaving nine bytes. Spark parks the request because1 + 9 < 25, its minimum share. The request cannot succeed or return a spillable refusal until another release supplies memory. A bounded probe using the head’s Java bridge and real Spark 4.1.3 reproducedWAITING, while the unanchored control completed immediately. Releasing one byte rescued the request. Relevant locking and task-entry semantics match Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. Beyond this existing concern, no additional introduced P1/P2 issues found within this review. - Key design decisions: Admission checks and provisional charges remain atomic. Reusing
SparkMemoryavoids duplicating overcommit accounting. The teardown guard coordinates pool generations without holding the registry lock across JNI. - Implementation sketch:
growrecords provisional charges and overcommit.try_growrolls back failed acquisitions and returns short grants.shrinksettles after Spark accepts the release, and Java usage accounting follows successful release. - Behavioral changes worth calling out: Compared with
branch-1.1, intended PR changes permit concurrent releases, retain an anchor byte, allow conservative refusals during settlement, and make replacement creation wait for teardown. Earlier-head contention benchmarks and spill comparisons were reviewed but not independently repeated. They do not resolve the exact-fit wait regression. - Suggested improvements: Resolve the existing P2 by preserving task-entry safety without parking an otherwise satisfiable allocation solely for the anchor byte. Documenting and testing the wait does not eliminate the regression.
Reviewed the complete eleven-file diff from 8c783aa88104616dcf0f7876b4f7a9ba71e2bf31 to cef941e424f68e7daf2a3c9bc03d856aa9ab3c91, including surrounding code and existing reviews, issue comments, inline comments and threads. Confirmed non-draft status. Read AGENTS.md and applied review-comet-pr, review-comet-memory-pr and review-comet-ffi-pr.
Exact-head CI: Comet CI and CodeQL report action_required. Labeling passed. No CI build/test verdict is available.
Validation: All 76 native memory-pool tests passed with cargo test -p datafusion-comet --lib memory_pools --offline --no-default-features. Recompiled the head’s Java bridge and Scala suite, then passed all eight memory-manager tests against Spark 4.1.3. The bounded exact-fit component probe confirmed the existing blocker. No full Maven build, JNI lifecycle/query suite, Spark SQL suite or benchmark ran. Other Spark versions were checked through source comparison. Project files remain unchanged.
Review state: Request changes for the unresolved existing P2.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem:
CometFairMemoryPoolheld its mutex across blocking Spark memory calls, preventing concurrent releases and memory-usage reads from progressing. - Design approach: Check limits and charge reservations under the mutex, call Spark unlocked, then settle accounting. A lazy anchor byte protects Spark’s task entry, and replacement pools wait for completed teardown.
- Correctness: No additional introduced P1/P2 issues found within this review. The existing P2 exact-fit wait remains at
native/core/src/execution/memory_pools/fair_pool.rs:159. With a 100-byte Spark pool and another task holding 90, the base can grant a firsttry_grow(10)immediately. Head takes its anchor first, leaving nine bytes. Spark parks the request because1 + 9 < 25, its minimum share. The request cannot succeed or return a spillable refusal until another release supplies memory. The current native regression confirms this behavior. A bounded probe using the freshly compiled head’s Java bridge and real Spark 4.1.3 reproducedWAITING, while the unanchored control completed immediately. Earlier release, diagnostic-deadlock, and pool-replacement regressions pass their focused tests. - Compatibility analysis: Checked
ExecutionMemoryPoolandTaskMemoryManagersources for Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0. Relevant waiting, locking, and task-entry semantics agree. Compared the affected paths withbranch-1.1. Task-wide Java manager sharing and blocking-aware acquisition already exist in the base. No expression semantics, fallback rules, or configuration defaults change here. - Key design decisions: Admission checks and provisional charges remain atomic across consumer registration changes. The anchor stays outside reservation and overcommit totals but counts as an ordinary Java grant. Teardown releases it before removing the registry entry.
- Implementation sketch:
growcharges memory and carries Spark’s shortfall as overcommit.try_growrolls back failed acquisitions and retains short-grant charges until their release succeeds.shrinksettles after Spark accepts the release. Java usage accounting now follows successful release too. - Performance: Earlier-head JNI benchmarks report lower contended release p99 and 1.5–2.5× throughput. The reported off-heap TPC-H comparisons retained identical spill counts. These support the contention improvement but do not cover the exact-fit stall. The implementation adds short bookkeeping lock acquisitions, an anchor acquisition while missing, and an anchor release at teardown. Exact-head performance was not independently benchmarked.
- Design: Keeping a short critical section around coupled limits and charges is straightforward and preserves fairness. Moving blocking calls outside it addresses the original contention problem. The anchor protocol still needs to preserve progress for otherwise satisfiable allocations.
- Abstraction & complexity: Reusing
SparkMemoryavoids duplicating the overcommit ledger.TeardownCompletegives pool replacement an explicit lifetime boundary without holding the registry lock across JNI. No separate P1/P2 abstraction issue was identified. - Behavioral changes worth calling out: Concurrent releases and usage reads can proceed during acquisitions. Provisional charges can conservatively refuse competing requests. The anchor retains one byte and keeps the task active until pool destruction. Replacement-plan creation can wait for teardown. At an exact-fit boundary, the anchor can introduce a wait, short grant, or one byte of overcommit.
- Suggested improvements: Resolve the existing P2 by preserving task-entry safety without parking an otherwise satisfiable allocation solely for the anchor byte. Documenting and testing that wait does not remove the regression. No additional P1/P2 findings are proposed.
Reviewed the complete eleven-file diff from b56349697b786ff2ad1c1bcf6ecf45b809af5f30 to b12b7ce883714455bba609e7ce3f8061cdeeb47d, including surrounding code and existing reviews, issue comments, inline comments, and threads. Confirmed non-draft status. Read AGENTS.md and applied review-comet-pr, review-comet-memory-pr, and review-comet-ffi-pr.
Exact-head CI: Comet CI and CodeQL report action_required. Labeling passed. No CI build/test verdict is available.
Validation: All 76 native memory-pool tests passed with cargo test -p datafusion-comet --lib memory_pools --offline --no-default-features. Freshly compiled the head’s Java bridge and Scala suite, then passed all eight CometTaskMemoryManagerSuite tests against Spark 4.1.3. The bounded exact-fit component probe confirmed the existing blocker. No full Maven build, JNI lifecycle/query suite, Spark SQL suite, or benchmark ran. Other Spark versions were checked through source comparison. Project files remain unchanged.
The anchor byte the fair pool kept with Spark made an exact-fit first allocation park or come back one byte short, and holding a byte back from a parked acquire of the same task can stall it. The pool now makes the same Spark calls as main for every operation, so a request Spark can grant is granted in one call. With nothing held at drop, the teardown wait in the task-shared registry goes too, and task_shared.rs, mod.rs and spill_replay.rs match main again. A release that empties the task's balance under a parked acquire of the same task is handled by the acquire retry in CometTaskMemoryManager (apache#6310), which this change depends on. Adds exact-fit tests, a test that a panicking grow rolls back its charge, and runs the short-grant test with debug logging on so the log path that must not take the task monitor is exercised.
The teardown wait it described is gone, and TaskSharedMemoryPool::drop is back to comparing pointers, which main's wording already describes.
The anchor is gone. The pool now makes the same Spark calls as main for every operation, so with another task holding 90 of 100 bytes a first Task-entry safety now comes from #6310, which you approved: |
|
Spill pressure run on the current head (
Every query result hashes the same across arms and matches the earlier runs. Executor stderr has no lines for This ran in Docker rather than directly on the host as the earlier runs did, with smaller heaps (driver 1g, executor 1536m) to fit the VM; the off-heap pool and cores are unchanged. Main at On wall time, the sum of per-query medians was 34.13 s on head against 34.32 s on base, q18 3.21 s against 3.30 s and q21 5.78 s against 5.90 s. The container swapped during the full-suite arms on both sides, so I read that as no slower rather than faster. The contention gain this PR targets shows in the microbenchmarks in the description, not at this level. |
Which issue does this PR close?
No dedicated issue.
Rationale for this change
CometFairMemoryPoolheld its internal mutex across the JNI acquire and release calls into Spark'sTaskMemoryManager. 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 memory to be released, and a release from another thread of the same task that has to take a lock the parked thread holds cannot land, so the task stalls until an unrelated task frees memory. The sibling unified pool already avoids this.What changes are included in this PR?
The pool builds on main's
SparkMemorywrapper and itsSparkMemoryManagertrait and adds two items to the wrapper.manager()exposes the raw Spark calls without the overcommit ledger, for a short grant the pool hands back itself.try_acquire_leaving_a_short_grantis the variant behindtry_acquirethat leaves a short grant with Spark instead of returning it, so the pool can keep those bytes charged until it has handed them back. Every operation makes the same Spark calls as on main.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.
growcharges under the lock the same way, skips the limit, and carries whatever Spark declines as overcommit, as on main.try_growrolls its charge back if Spark errors or panics inside the JNI frame, andgrowdoes the same on a panic. A short grant stays charged until Spark has it back, and with overcommit owed the request carries the debt, so a refusedtry_growcan receive more bytes than it asked for and hands all of them back. If handing a short grant back fails, the pool logs it, keeps the bytes charged, since Spark keeps them until the task ends, and refuses the request as usual, where main returned the JVM error fromtry_grow.shrinkchecks 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
try_growcalls still cannot jointly take a consumer past its share or the pool pastpool_size, 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 secondtry_growin 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 by the share or pool check 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 atry_growthe limit would have refused, and a refusal is the ordinary error a spillable operator answers by spilling.Spark's
ExecutionMemoryPoolremoves a task's entry the moment its balance reaches zero, and an acquire waiting inside Spark reads that entry when it wakes, so a release that zeroes the balance under a waiting acquire makes it throwNoSuchElementException: key not found: <task attempt id>. With the lock no longer held across the call, one native thread's release can do that to another native thread's waiting acquire in the same task. #6310 makesCometTaskMemoryManager.acquireMemoryask Spark again in that case, which is why this PR should land after it. Its retry is bounded at three attempts and then returns a zero grant, which the pool treats as a refusal or overcommit like any short grant, so the bound is fine for this race too.releaseMemorynow calls Spark before it moves its count, so a release Spark rejects moves nothing andusedkeeps matching what the pool still has charged. The memory guide now explains why a short grant is logged withgetMemoryConsumptionForThisTaskrather thanshowMemoryUsage, which main already does:showMemoryUsagewould take the task's monitor while another acquire of the task may be waiting for those bytes. Thejni_api.rsregistry comment no longer describes the fair pool holding its lock across JNI.task_shared.rs,spill_replay.rs,mod.rsand the memory review skill are unchanged from main. The memory guide gains paragraphs on the lock rule, the two windows and the short-grant log, and its pool-stack diagram now showsPlanMemoryPool.How are these changes tested?
Unit tests in
fair_pool.rsrun the pool against two stand-ins for Spark. Main'sFakeSparkcovers the overcommit paths, and a new test pins that a refusedtry_growwith overcommit outstanding hands the whole grant back with neither the pool's total nor the overcommit moving. A stub of Spark 4.1.3'sExecutionMemoryPoolcovers the rest: the per-task entry lifecycle, the 1/N and 1/(2N) shares, the wait loop andnotifyAllon release, the bridge asking again when the entry is gone (standing in for #6310), 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 no consumer passes its share and the pool never passes its size, and that the pool, Spark and the JVM count all end at zero. They cover the rollback paths: short grants, acquire failure, a panic inside the call for both
growandtry_grow, 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 exact fit: with another task holding 90 of 100 bytes, a firsttry_grow(10)orgrow(10)is granted with one Spark call; with the other task at 70,try_grow(30)andgrow(30)are granted in full; and the same holds after a full release. 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.CometTaskMemoryManagerSuitegains two tests: a release Spark rejects leavesusedunchanged, and a short grant is handed back while another acquire of the task waits in Spark, with debug logging on so the short-grant log runs. The suite passes on Spark 3.4, 3.5, 4.0 and 4.1.Microbenchmarks against a real Spark 4.1.3
TaskMemoryManagerover an off-heapUnifiedMemoryManager, base and head in one binary, are in the discussion. They predate theSparkMemorymerge and the current Spark call pattern, 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 in the settled pass the head to base ratio on the sum of medians was 0.99. A rerun on 88d881d spilled the same bytes the same number of times. Both heads predate the current design, so they cover the lock discipline rather than this exact head.