dwsmith1983 opened a new pull request, #6403:
URL: https://github.com/apache/datafusion-comet/pull/6403

   ## Which issue does this PR close?
   
   Closes #6304.
   
   ## 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 from Spark 3.4 through 4.2.
   
   `TaskMemoryManager.releaseExecutionMemory` takes no monitor, so a release 
can run while an acquire of the same task is parked. When 
`CometUnifiedShuffleMemoryAllocator` has 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 expect `SparkOutOfMemoryError` from the allocator, so the task fails 
instead of spilling.
   
   The root cause is in Spark, filed as 
[SPARK-59827](https://issues.apache.org/jira/browse/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.allocate` and a new `allocateArray` 
override run the Spark call through a helper that catches this one exception, 
identified by its message prefix (`key not found: `) and a frame in 
`ExecutionMemoryPool`, 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's `used` and Spark's counters have not moved. 
After three attempts it throws the same `SparkOutOfMemoryError` it 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.
   
   `allocateArray` needs its own override because the inherited 
`MemoryConsumer.allocateArray` asks Spark for the page directly rather than 
through `allocate`, 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, and the memory review skill lists the new suite.
   
   ## How are these changes tested?
   
   A new `CometUnifiedShuffleMemoryAllocatorSuite`, registered in both PR build 
workflows, covers:
   
   - the race from the issue, for a page and for a pointer array: on a 100 byte 
off-heap `UnifiedMemoryManager`, the task holds 10 bytes through a 
`CometTaskMemoryManager` and 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 throws `key 
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.
   - a `NoSuchElementException` with a different message, or without an 
`ExecutionMemoryPool` frame, is rethrown unchanged after a single call.
   - when every attempt throws the matching exception, the allocator gives up 
after the cap with `SparkOutOfMemoryError` (`UNABLE_TO_ACQUIRE_MEMORY`, the 
requested bytes, zero received, and the last exception as the cause), with the 
accounting unchanged.
   - an allocation that loses its entry on the first attempts is granted on the 
last one.
   
   The suite passes on Spark 3.4, 3.5, 4.0 and 4.1, and the race tests fail 
without the retry.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to