Repository navigation
fix: retry a shuffle page allocation that Spark failed after dropping the task's entry - #6403
dwsmith1983 wants to merge 17 commits into
Conversation
… the task's entry Spark's ExecutionMemoryPool registers a task's memoryForTask entry once, before its wait loop, and reads it on every pass, while releaseMemory removes the entry when the task's balance reaches zero. When the JVM shuffle allocator's page or pointer array request was parked below the task's minimum share and Comet's native memory consumer released the task's last bytes, the request woke to NoSuchElementException "key not found: <taskAttemptId>" and failed the task, because the shuffle writers only expect SparkOutOfMemoryError. CometUnifiedShuffleMemoryAllocator now retries allocate and allocateArray when Spark throws that exception from ExecutionMemoryPool. The retry registers the task again and waits as the first attempt would have. Nothing is counted or allocated before the exception, so retrying cannot leak. After three attempts it refuses the request with the same SparkOutOfMemoryError it uses for a refused page, with the last exception as its cause. The Spark issue is SPARK-59827.
Only the shuffle allocator's callers are guarded; Spark's own operators in the task can still hit the removed entry until Spark re-registers a waiting task (apache/spark#59103).
ad2225d to
0de3286
Compare
|
@andygrove could you approve a CI run on |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: A waiting JVM shuffle allocation could fail with
NoSuchElementExceptionwhen another consumer released the task’s last reserved bytes. - Design approach: Retry the specific missing-entry failure for both shuffle pages and pointer arrays through one private helper.
- Correctness / compatibility analysis: Checked Spark’s memory-pool, allocation and exception implementations in 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. The relevant failure occurs before a grant or page allocation, so retrying preserves accounting and Spark’s fair-share checks.
- Key design decisions: Message-and-stack matching limits retries to the intended failure. The helper keeps the abstraction local. Successful allocations still make one Spark call, with exception inspection and logging confined to the failure path.
- Implementation sketch: Wrap
allocatePageand inheritedallocateArray, cap attempts at three, preserve the final exception as the cause, register the new suite in both CI workflows, and update the memory and shuffle documentation. - Behavioral changes worth calling out: Exhausted retries become
SparkOutOfMemoryError, reaching existing shuffle refusal/spill handling. Unrelated exceptions propagate unchanged. Spark-owned callers remain outside this workaround. - Suggested improvements: None at P1/P2. No introduced P1/P2 issues found within this review.
Reviewed the entire six-file diff from 9f68a4144fdefc70ec26f3d80b0e4eafe72a8df4 to 20e1327dd353b32aec9be2833077766573c6057b. The PR is not a draft. Read the existing discussion and verified there are no reviews, inline comments or review threads containing unresolved concerns. Routed skills: review-comet-pr, review-comet-memory-pr, review-comet-shuffle-pr.
Exact-head CI: Comet CI and CodeQL report action_required, awaiting approval. Only the labeling check passed. There is no completed build/test verdict.
Validation: Independently compiled the exact-head allocator, supporting classes and new suite against cached Spark 4.1.3 dependencies using JDK 17. All five tests passed. Both race tests failed against the base allocator with NoSuchElementException: key not found: 0, confirming regression coverage. Full Maven/native builds, end-to-end shuffle tests and performance benchmarks were not run. Other supported Spark versions were source-checked only.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: A waiting JVM shuffle allocation could throw
NoSuchElementExceptionwhen another consumer released the task’s last reserved bytes. - Design approach: Retry that specific Spark failure for both pages and pointer arrays through one private helper.
- Correctness / compatibility analysis: Checked the relevant Spark sources in 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. The missing-entry lookup fails before a grant or page allocation. Retrying re-registers the task and preserves Spark’s fair-share checks and allocation accounting.
- Key design decisions: Message-and-stack matching narrows the retry to the intended failure. The shared helper keeps both allocation paths consistent. Successful requests still make one Spark allocation call through a supplier wrapper. Stack inspection and logging occur only on failure. Performance was not benchmarked.
- Implementation sketch: Wrap
allocatePageand inheritedallocateArray, cap attempts at three, preserve the final exception as the cause, register the new suite in both CI workflows, and update the memory and shuffle documentation. - Behavioral changes worth calling out: Exhausted retries become
SparkOutOfMemoryError, reaching existing allocation-refusal handling. Unrelated exceptions propagate unchanged. Compared withbranch-1.1at992c806a7e38c2e88bd018aa5774164b0850e1fa, whose allocator matches the supplied base, this is an intended recovery improvement. Spark-owned callers remain exposed to the upstream race. - Suggested improvements: None at P1/P2. No introduced P1/P2 issues found within this review.
Reviewed the entire six-file diff from fef94f6cd78b18151dff57b7a936798385356de5 to 48458d0840317700ce4febb184e847243c8aee05. The PR is not a draft. Read the existing review and issue comment. There are no inline comments, review threads or substantiated unresolved P1/P2 concerns. Routed skills: review-comet-pr, review-comet-memory-pr, review-comet-shuffle-pr.
Exact-head CI: Comet CI and CodeQL report action_required, awaiting approval. The labeling check passed. There is no completed build/test CI verdict.
Validation: Independently compiled the exact-head allocator, supporting classes and suite against cached Spark 4.1.3 dependencies using JDK 21 with Java 17 source/target settings. All five tests passed. Both race tests failed against the supplied base allocator with NoSuchElementException: key not found: 0, confirming regression coverage. Full Maven/native builds, end-to-end shuffle tests and benchmarks were not run. Other supported Spark versions were source-checked only.
|
@andygrove could you take a look at this one when you have a moment? It is merged up with |
Which issue does this PR close?
Closes #6304.
Rationale for this change
Spark's
ExecutionMemoryPool.acquireMemoryregisters the task'smemoryForTaskentry once, before its wait loop, and reads it withmemoryForTask(taskAttemptId)on every pass.releaseMemoryremoves the entry when the task's balance reaches zero and wakes the waiters. A caller that wakes after the removal throwsNoSuchElementException: key not found: <taskAttemptId>. The code is the same from Spark 3.4 through 4.2.TaskMemoryManager.releaseExecutionMemorytakes no monitor, so a release can run while an acquire of the same task is parked. WhenCometUnifiedShuffleMemoryAllocatorhas a page or pointer array request parked below the task's minimum share and Comet's native consumer releases the task's last bytes, the parked request fails with that exception. The shuffle writers only expectSparkOutOfMemoryErrorfrom the allocator, so the task fails instead of spilling.The root cause is in Spark, filed as SPARK-59827 with a fix in apache/spark#59103. Until that lands in the versions Comet supports, Comet can only guard its own callers, and Spark's own operators in the task stay exposed.
What changes are included in this PR?
CometUnifiedShuffleMemoryAllocator.allocateand a newallocateArrayoverride run the Spark call through a helper that catches this one exception, identified by its message prefix (key not found:) and a frame inExecutionMemoryPool, and calls Spark again. The retry registers the task's entry again, checks its share and parks again if it has to, so the caller ends up with what Spark would have granted had the entry stayed. Retrying is safe because Spark only removes the entry at a zero balance, so the failed call was granted nothing and the allocator'susedand Spark's counters have not moved. After three attempts it throws the sameSparkOutOfMemoryErrorit uses for a refused page, with the last exception as the cause, and the writers spill as they would for any other refusal. Any other exception is rethrown. A retry logs at info level and giving up logs a warning.allocateArrayneeds its own override because the inheritedMemoryConsumer.allocateArrayasks Spark for the page directly rather than throughallocate, and the sorters' pointer arrays come through it.#6310 adds the same retry for Comet's native acquires in
CometTaskMemoryManager, with its own matcher. Whichever of the two lands second will move both onto one shared helper.The contributor guide's memory management and JVM shuffle pages describe the failure and the retry.
How are these changes tested?
A new
CometUnifiedShuffleMemoryAllocatorSuite, registered in both PR build workflows, covers:UnifiedMemoryManager, the task holds 10 bytes through aCometTaskMemoryManagerand another task holds 90. An allocation on a second thread parks in Spark below the task's minimum share, and the native consumer releases the task's last 10 bytes. On main the parked allocation throwskey not found: 0. With this change it is granted in full once the other task frees its memory, with one retry logged and the allocator's accounting back at zero afterwards.NoSuchElementExceptionwith a different message, or without anExecutionMemoryPoolframe, is rethrown unchanged after a single call.SparkOutOfMemoryError(UNABLE_TO_ACQUIRE_MEMORY, the requested bytes, zero received, and the last exception as the cause), with the accounting unchanged.The suite passes on Spark 3.4, 3.5, 4.0 and 4.1, and the race tests fail without the retry.