[ 
https://issues.apache.org/jira/browse/SPARK-59601?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

ASF GitHub Bot updated SPARK-59601:
-----------------------------------
    Labels: pull-request-available  (was: )

> DiskRowQueue holds its output stream open until first read after spilling
> -------------------------------------------------------------------------
>
>                 Key: SPARK-59601
>                 URL: https://issues.apache.org/jira/browse/SPARK-59601
>             Project: Spark
>          Issue Type: Bug
>          Components: SQL
>    Affects Versions: 5.0.0
>            Reporter: Qing Fu
>            Priority: Major
>              Labels: pull-request-available
>
> {{HybridQueue.spill()}} converts every in-memory queue but the last into a 
> disk-backed queue, draining the rows into it. The new {{DiskRowQueue}} is 
> complete at that point -- only the last queue is ever written to -- but its 
> output stream is left open.
> {{DiskRowQueue}} opens its writer eagerly in the constructor:
> {code:scala}
> private var out = new DataOutputStream(serMgr.wrapForEncryption(
>   new BufferedOutputStream(new FileOutputStream(file.toString))))
> {code}
> and only closes it on the queue's *first read*, inside {{remove()}}. So after 
> N spills a task holds N open {{DataOutputStream}} -> {{BufferedOutputStream}} 
> -> {{FileOutputStream}} chains (plus the encryption wrapper when I/O 
> encryption is on), each with its 8 KB buffer, until those queues are read. 
> That is a file descriptor and a buffer per spilled queue held at exactly the 
> moment memory is scarce -- spilling happens because the task is under memory 
> pressure.
> It is also observable as unflushed data: rows written into a spilled queue 
> sit in the output stream's buffer rather than on disk until the first read.
> This is a follow-up on SPARK-53481, which extracted the generic 
> {{HybridQueue}} / {{Queue}} classes.
> h3. Proposed fix
> Add a {{closeOutputStream()}} method to the {{Queue}} trait, defaulting to a 
> no-op (in-memory queues hold no write-only resource), override it in 
> {{DiskRowQueue}} to close and clear {{out}}, and call it from the spill loop 
> once a queue has been drained.
> One subtlety: {{DiskRowQueue.remove()}} currently uses {{if (out != null)}} 
> as its "first read" signal and opens the input stream inside that branch. 
> Once the output stream can already be closed before the first read, that 
> branch no longer fires and {{remove()}} NPEs on {{in.readInt()}}. The check 
> has to be keyed off {{in == null}} instead.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

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

Reply via email to