asfimport opened a new issue, #31769:
URL: https://github.com/apache/arrow/issues/31769

   The current implementation of the hash-join node current queues in memory 
the hashtable, the entire build side input, and the entire probe side input 
(e.g. the entire dataset).  This means the current implementation will run out 
of memory and crash if the input dataset is larger than the memory on the 
system.
   
   By spilling to disk when memory starts to fill up we can allow the hash-join 
node to process datasets larger than the available memory on the machine.
   
   **Reporter**: [Weston 
Pace](https://issues.apache.org/jira/browse/ARROW-16389) / @westonpace
   #### Related issues:
   - [[C++] Naive spillover implementation for 
join](https://github.com/apache/arrow/issues/29750) (supercedes)
   #### PRs and other links:
   - [GitHub Pull Request #13669](https://github.com/apache/arrow/pull/13669)
   
   <sub>**Note**: *This issue was originally created as 
[ARROW-16389](https://issues.apache.org/jira/browse/ARROW-16389). Please see 
the [migration documentation](https://github.com/apache/arrow/issues/14542) for 
further details.*</sub>


-- 
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]

Reply via email to