[
https://issues.apache.org/jira/browse/SPARK-59827?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Dustin Smith resolved SPARK-59827.
----------------------------------
Resolution: Duplicate
Duplicate of SPARK-59444, which covers the same bug and has an open fix in
https://github.com/apache/spark/pull/58747. Closed
https://github.com/apache/spark/pull/59103 in its favor.
> ExecutionMemoryPool.acquireMemory throws NoSuchElementException when the
> task's entry is removed while it waits
> ---------------------------------------------------------------------------------------------------------------
>
> Key: SPARK-59827
> URL: https://issues.apache.org/jira/browse/SPARK-59827
> Project: Spark
> Issue Type: Bug
> Components: Spark Core
> Affects Versions: 3.4.4, 3.5.8, 4.1.3, 4.0.4
> Reporter: Dustin Smith
> Priority: Major
> Labels: pull-request-available
>
> ExecutionMemoryPool.acquireMemory registers the task's memoryForTask entry
> once, before its wait loop, and reads the entry with
> memoryForTask(taskAttemptId) on every pass of the loop. releaseMemory removes
> the entry when the task's balance reaches zero and calls notifyAll. A caller
> that was parked in lock.wait() and wakes after that removal throws
> java.util.NoSuchElementException: key not found: <taskAttemptId> instead of
> continuing to wait or being granted memory.
> On master
> (core/src/main/scala/org/apache/spark/memory/ExecutionMemoryPool.scala):
> * lines 103-104: the entry is registered only if absent, before the loop
> * line 115: val curMem = memoryForTask(taskAttemptId), no default
> * line 142: lock.wait()
> * lines 165-167: releaseMemory subtracts and removes the entry at zero
> The entry can be removed under a waiting caller because
> TaskMemoryManager.releaseExecutionMemory does not take the TaskMemoryManager
> monitor that acquireExecutionMemory holds while it waits. Any consumer of the
> same task that releases through releaseExecutionMemory
> (MemoryConsumer.freeMemory, or a caller on another thread) can bring the
> balance to zero while another acquire of that task is parked below its share.
> Reproduction (component level, off-heap UnifiedMemoryManager with 100 bytes,
> storageFraction 0):
> # Task 0 acquires 10 bytes. Task 1 acquires 90 bytes.
> # On another thread, task 0 calls acquireExecutionMemory(20). It parks,
> because 10 is below the minimum share of 25 for two active tasks.
> # Task 0's consumer frees its 10 bytes with freeMemory. The balance is zero,
> so releaseMemory removes the entry and wakes the waiter.
> # The waiter throws NoSuchElementException: key not found: 0.
> Expected: the parked acquire keeps waiting, or is granted, exactly as if the
> entry had stayed. The caller has nothing from the failed call, so there is no
> accounting to undo, but the exception is not a memory error and callers such
> as UnsafeExternalSorter do not treat it as a signal to spill.
> Suggested fix: read the entry with
> memoryForTask.getOrElseUpdate(taskAttemptId, 0L) inside the loop (or
> re-register after waking), so a removed entry is recreated at zero and the
> loop continues.
> Seen in Apache DataFusion Comet, whose native memory consumer parks in this
> loop while other consumers of the task release:
> https://github.com/apache/datafusion-comet/issues/6224 and
> https://github.com/apache/datafusion-comet/issues/6304. Comet works around it
> by retrying the acquire.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]