Apply memory backpressure to the shuffle - #23834
Draft
madsbk wants to merge 1 commit into
Draft
Conversation
|
Auto-sync is disabled for draft pull requests in this repository. Workflows must be run manually. Contributors can view more details about this message here. |
Contributor
Author
|
/ok to test |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Depend on #23833.
The shuffle now takes its device memory from
reserve_memory()on both sides, the fourShuffleManager.Insertermethods going in andextract_chunkcoming out, so a request that cannot be satisfied queues alongside the other actors' and is served by priority. Before this both let the C++ side reserve and spill internally. cudf-polars already reserves this way for scans inactor_graph/io.py.Breaking change
The insert methods and
extract_chunkare coroutines now, andLocalRepartitioner._iter_chunksis an async generator. That updates twenty-one call sites acrossgroupby.py,over.py,sort.py,shuffle.pyand the shuffler tests.Notes
Reserving inside the insert methods rather than at the call sites keeps the accounting in one place. That matters for
insert_hash_keysandinsert_index, whose Python-side reorder allocates a full table copy that nothing reserved before. They reservepartition_and_pack_cost(), which covers a reorder plus a pack, consume the reorder's share once it lands, and hand the remainder tosplit_and_pack()._iter_chunksreserves per piece, so theTODOthere about batching pieces up totarget_partition_sizewould cut the reservation traffic as well as the unpack overhead.