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]
