fix: retry a native acquire that Spark failed after dropping the task's entry - #6310
dwsmith1983 wants to merge 3 commits into
Conversation
…'s entry Spark removes a task's execution pool entry when its balance reaches zero, and an acquire that was parked in the pool then throws NoSuchElementException when it wakes up. Releases take no task monitor, so another native thread or a JVM consumer of the same task can empty the balance while a native acquire waits. The error is not a memory refusal, so the operator could not spill and the task failed. CometTaskMemoryManager now recognises that one exception and calls Spark again. A missing entry means the task held nothing from Spark, so the failed call has no grant to reconcile, and the retry registers the task again and waits for its share as the first call would have. After three attempts it returns a zero grant, which the native side treats as a refusal and spills. Closes apache#6224.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Releasing a task’s last reserved bytes could remove Spark’s bookkeeping entry while a native acquisition waited, causing
NoSuchElementExceptionand task failure. - Design approach: Retry the specifically identified missing-entry exception, with three attempts before returning a zero grant.
- Correctness / compatibility analysis: Checked Spark’s memory-pool and task-manager sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. The zero-balance removal rule supports retrying without losing a partial grant. Unrelated exceptions still propagate.
- Key design decisions: Matching both the task-specific message and an
ExecutionMemoryPoolframe limits recovery to the intended failure. Successful acquisitions still make one Spark call. Exception inspection and retry logging occur only on failure, and no new locks are introduced. - Implementation sketch: Two private helpers wrap acquisition. Four regression tests extend the existing suite. Reviewed all five files in the supplied base-to-head diff, including both window-documentation edits.
- Behavioral changes worth calling out: A missing entry becomes recoverable. Repeated failures return zero through the existing native refusal/overcommit handling. Direct JVM page allocations remain outside this fix, as documented in #6304.
- Suggested improvements: None meeting the P1/P2 reporting threshold. No introduced P1/P2 issues found within this review.
Reviewed full SHA 6bc7575621819b7dd5ff21cab5c07c84e58756fc against base 76493618769ff85e7a43e6f6c79f8b58d9e0b1ad. The PR remains non-draft. Routed skills: review-comet-pr and review-comet-memory-pr. The snapshot and live discussion checks contained no existing review concerns.
Validation: Freshly compiled the changed Java class and Scala suite against cached Spark 4.1.3 dependencies. All seven tests passed. Running the same suite against the base implementation reproduced failures in both race tests and the retry-cap test. This was isolated component validation. Full Maven/native builds, end-to-end spill tests, other runtime profiles and performance benchmarks were not run.
Exact-head CI: The label job passed. Comet CI, CodeQL and Check PR Title report action_required, so no build/test CI verdict is available.
Which issue does this PR close?
Closes #6224.
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 in Spark 3.4 through 4.1 and on master.TaskMemoryManager.acquireExecutionMemoryholds the task memory manager's monitor while it is parked, butreleaseExecutionMemorytakes no monitor. So a release can empty the task's balance while an acquire of the same task is parked: undergreedy_unified, one native thread's release fails another native thread's parked acquire; underfair_unified, a JVM consumer of the task freeing its last bytes withfreeMemorydoes the same. The exception reaches the native side as a plain error rather thanResourcesExhausted, so the operator cannot spill and the task fails.What changes are included in this PR?
CometTaskMemoryManager.acquireMemorycalls Spark through a helper that catches this one exception, identified by its message (key not found:plus this task's attempt id) and a frame inExecutionMemoryPool, and callsacquireExecutionMemoryagain. The retry re-registers the entry, re-checks the task's share and parks again if it has to, so the caller ends up with what Spark would have granted had the entry stayed. After three attempts it returns 0, which the native side treats as a refusal and spills on. Any other exception is rethrown. A retry logs at info level, and giving up logs a warning.The retry is safe because Spark only removes the entry at a zero balance, which cannot hold while the failed call has a partial grant. So the failed call has nothing to reconcile, and
usedand Spark's counters have not moved. The helper holds no lock of its own across attempts.Retrying from the acquire side reaches the same end state as fixing the loop in Spark (reading the entry with
getOrElseUpdate(taskAttemptId, 0L), filed as SPARK-59827), without coordinating the release side, which does not know whether an acquire is parked. The Spark fix would also cover JVM consumers that callallocatePagedirectly, such asCometUnifiedShuffleMemoryAllocatorand Spark's own operators; they are not covered here and are tracked in #6304.The contributor guide's memory management page describes the new failure and why the retry is safe.
How are these changes tested?
Four tests in
CometTaskMemoryManagerSuiteon a 100-byte off-heapUnifiedMemoryManager:CometTaskMemoryManager; the task holds 10 bytes and another task holds 90; an acquire of 20 parks below the minimum share; releasing the 10 bytes empties the balance. On main the parked acquire throwskey not found: 0. With this change it completes with 20 once the other task frees its memory, andgetUsedand the task's consumption both read 20.MemoryConsumerof the task and released withfreeMemory.NoSuchElementExceptionwith another task's id, or without anExecutionMemoryPoolframe, is rethrown unchanged after one call.CometTaskMemoryManagerSuite,CometUnboundedShuffleMemoryAllocatorSuiteandCometExecIteratorLifecycleSuitepass on Spark 3.4, 3.5, 4.0 and 4.1, the spill tests inCometExecSuitepass on 3.5, and the native memory pool tests pass. With the retry disabled, the two race tests and the cap test fail again.