Skip to content

fix: retry a native acquire that Spark failed after dropping the task's entry - #6310

Open
dwsmith1983 wants to merge 3 commits into
apache:mainfrom
dwsmith1983:fix/6224-parked-acquire-retry
Open

dwsmith1983 wants to merge 3 commits into
apache:mainfrom
dwsmith1983:fix/6224-parked-acquire-retry

Conversation

@dwsmith1983

@dwsmith1983 dwsmith1983 commented Sep 28, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Closes #6224.

Rationale for this change

Spark's ExecutionMemoryPool.acquireMemory registers the task's memoryForTask entry once, before its wait loop, and reads it with memoryForTask(taskAttemptId) on every pass. releaseMemory removes the entry when the task's balance reaches zero and wakes the waiters. A caller that wakes after the removal throws NoSuchElementException: key not found: <taskAttemptId>. The code is the same in Spark 3.4 through 4.1 and on master.

TaskMemoryManager.acquireExecutionMemory holds the task memory manager's monitor while it is parked, but releaseExecutionMemory takes no monitor. So a release can empty the task's balance while an acquire of the same task is parked: under greedy_unified, one native thread's release fails another native thread's parked acquire; under fair_unified, a JVM consumer of the task freeing its last bytes with freeMemory does the same. The exception reaches the native side as a plain error rather than ResourcesExhausted, so the operator cannot spill and the task fails.

What changes are included in this PR?

CometTaskMemoryManager.acquireMemory calls 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 in ExecutionMemoryPool, and calls acquireExecutionMemory again. 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 used and 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 call allocatePage directly, such as CometUnifiedShuffleMemoryAllocator and 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 CometTaskMemoryManagerSuite on a 100-byte off-heap UnifiedMemoryManager:

  • two threads share one 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 throws key not found: 0. With this change it completes with 20 once the other task frees its memory, and getUsed and the task's consumption both read 20.
  • the same, with the 10 bytes held by a sibling MemoryConsumer of the task and released with freeMemory.
  • a NoSuchElementException with another task's id, or without an ExecutionMemoryPool frame, is rethrown unchanged after one call.
  • when every attempt throws the matching exception, the helper returns 0 after the cap with accounting unchanged.

CometTaskMemoryManagerSuite, CometUnboundedShuffleMemoryAllocatorSuite and CometExecIteratorLifecycleSuite pass on Spark 3.4, 3.5, 4.0 and 4.1, the spill tests in CometExecSuite pass on 3.5, and the native memory pool tests pass. With the retry disabled, the two race tests and the cap test fail again.

…'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.
@github-actions github-actions Bot added bug Something isn't working area:memory Memory pools, reservations, OOM handling labels Sep 28, 2026

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Summary

  • Prior state and problem: Releasing a task’s last reserved bytes could remove Spark’s bookkeeping entry while a native acquisition waited, causing NoSuchElementException and 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 ExecutionMemoryPool frame 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.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:memory Memory pools, reservations, OOM handling bug Something isn't working

Projects

None yet

Development

Successfully merging this pull request may close these issues.

fair_unified: a JVM consumer freeing its last bytes can fail a parked native acquire with NoSuchElementException

2 participants