perf: stop holding the fair pool lock across blocking memory calls - #5613
dwsmith1983 wants to merge 80 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.
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, preventing concurrent native releases from progressing. - Design approach: Check and charge reservations under the mutex, call Spark unlocked, then settle accounting. Reuse
SparkMemoryfor 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:
growrecords overcommit,try_growrolls back failed requests and returns short grants, andshrinksettles 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) |
There was a problem hiding this comment.
[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.
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem:
fair_unifiedheld 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+anchorandNoSuchElementExceptionafter 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
SparkMemoryand retaining the existing registry limits added complexity. Separate anchor accounting prevents false one-byte leak warnings without masking reservation leaks. - Implementation sketch:
growrecords overcommit,try_growrolls back failures and returns short grants, andshrinksettles 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.
|
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 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.
|
@sunchao both points from your last review are addressed. The |
|
One more symptom of holding the fair pool's lock across the JNI call. While a fair-pool acquire is parked in Spark's |
…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.
That is the lock this PR stops holding across the Spark call. On main,
|
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem:
fair_unifiedheld 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
SparkMemorypreserves overcommit handling, and separate anchor accounting avoids false leak warnings. The anchor's lifetime still needs coordination across replacement pools. - Implementation sketch:
growrecords provisional charges,try_growrolls back failures and returns short grants, andshrinksettles 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) |
There was a problem hiding this comment.
[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.
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.
|
This is a light fully automated review since there are so many PRs open.
The new helper in |
#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)
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
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 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
SparkMemorywrapper and itsSparkMemoryManagertrait, 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_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.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. 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.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 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_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 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 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 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 firsttry_growthat passes the fair limit, or the firstgrow, 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. Ongrowthe anchor is best effort, and atry_growwhose 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 throughacquireAnchorandreleaseAnchoronCometTaskMemoryManager, which count it toward the task's balance and the consumer's usage but not toward the usageclosechecks, 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.rsandmod.rsdiffer 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.rsrun the pool against two stand-ins for Spark. Main'sFakeSparkcovers the overcommit paths: the existing test of agrowpast 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 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, 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
TaskMemoryManagerover an off-heapUnifiedMemoryManager, base and head in one binary, are in the discussion. They predate the lazy anchor and theSparkMemorymerge, 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.