From ef0b44ec47bae3375acdb9792a97fd803212e247 Mon Sep 17 00:00:00 2001 From: You-Cheng Lin Date: Fri, 25 Sep 2026 12:53:52 +0800 Subject: [PATCH 1/4] [Data][Docs] Document disk-based shuffle in Data internals Add a disk-based shuffle subsection under shuffle v2: what it is, the per-node file-server actor serving shards over Arrow Flight, when to use it, and how to enable it. Signed-off-by: You-Cheng Lin Signed-off-by: You-Cheng Lin --- doc/source/data/internals.md | 31 +++++++++++++++++++++++++++++++ 1 file changed, 31 insertions(+) diff --git a/doc/source/data/internals.md b/doc/source/data/internals.md index 88ca56046337..bd503aea5c23 100644 --- a/doc/source/data/internals.md +++ b/doc/source/data/internals.md @@ -104,6 +104,37 @@ Shuffle v2 supports the following operations: Shuffle v2 doesn't yet support {meth}`Dataset.sort ` or {meth}`Dataset.random_shuffle `, which use the {ref}`range-partitioning shuffle `. +(disk-based-shuffle)= + +##### Disk-based shuffle + +Shuffle v2 can optionally transport intermediate shuffle data through node-local disk files instead of the Ray object store. In this *disk-based shuffle* mode, shuffle bytes never enter the object store; only small file-handle metadata travels through it. + +Disk-based shuffle works as follows: + +1. Each map task partitions its input and writes all of the resulting shards into a single file on its node's local disk. +2. Each node runs a lightweight file-server actor that serves those files over [Arrow Flight](https://arrow.apache.org/docs/format/Flight.html), a high-throughput RPC protocol for Arrow data. +3. Each reduce task fetches its partition's byte ranges from every mapper node's file server over Arrow Flight, then reduces them. + +Ray Data deletes the shuffle files when the shuffle completes. + +Because intermediate data goes straight to disk instead of filling the object store and relying on reactive spilling, disk-based shuffle is a good fit when: + +- The shuffled data is much larger than the cluster's aggregate object-store memory, for example large joins and aggregations over out-of-core datasets. +- Object-store pressure from shuffle intermediates causes spilling that interferes with other operators in the pipeline. + +For shuffles that fit comfortably in object-store memory, the default in-memory path avoids the disk round trip and is typically faster. + +To enable disk-based shuffle for key-based repartitioning, aggregations, and joins under the `SHUFFLE_V2` strategy, set `DataContext.use_disk_based_hash_shuffle` before creating a `Dataset`: + +```python +import ray + +ray.data.DataContext.get_current().use_disk_based_hash_shuffle = True +``` + +Alternatively, set the environment variable `RAY_DATA_ENABLE_DISK_SHUFFLE=1`. + (tuning-shuffle-v2)= ##### Tuning shuffle v2 From 179f3fc0c6fcfb24e355794bb884f14d6aef15ef Mon Sep 17 00:00:00 2001 From: You-Cheng Lin Date: Fri, 25 Sep 2026 13:12:32 +0800 Subject: [PATCH 2/4] [Data][Docs] Cross-reference disk-based shuffle from key-based shuffle APIs Add a tip to Dataset.repartition (keys), Dataset.join, and Dataset.groupby docstrings pointing to the disk-based shuffle section for datasets larger than aggregate object-store memory. Signed-off-by: You-Cheng Lin --- python/ray/data/dataset.py | 24 ++++++++++++++++++++++++ 1 file changed, 24 insertions(+) diff --git a/python/ray/data/dataset.py b/python/ray/data/dataset.py index 8484817f9c0e..9ddba0a0f2bd 100644 --- a/python/ray/data/dataset.py +++ b/python/ray/data/dataset.py @@ -2334,6 +2334,14 @@ def repartition( .. https://docs.google.com/drawings/d/132jhE3KXZsf29ho1yUdPrCHB9uheHBWHJhDQMXqIVPA/edit + .. tip:: + + Repartitioning with ``keys`` hash-shuffles the whole dataset. If the + dataset is much larger than the cluster's aggregate object-store + memory, enable :ref:`disk-based shuffle ` by + setting ``DataContext.use_disk_based_hash_shuffle = True`` or the + environment variable ``RAY_DATA_ENABLE_DISK_SHUFFLE=1``. + Examples: >>> import ray >>> ds = ray.data.range(100).repartition(10).materialize() @@ -3613,6 +3621,14 @@ def join( Joins require the ``polars`` package. + .. tip:: + + Joins hash-shuffle both input datasets by key. If the inputs are + much larger than the cluster's aggregate object-store memory, + enable :ref:`disk-based shuffle ` by setting + ``DataContext.use_disk_based_hash_shuffle = True`` or the + environment variable ``RAY_DATA_ENABLE_DISK_SHUFFLE=1``. + Args: ds: Other dataset to join against join_type: The kind of join that should be performed, one of ("inner", @@ -3790,6 +3806,14 @@ def groupby( Use this method to transform data based on a categorical variable. + .. tip:: + + Grouping hash-shuffles the dataset by key. If the dataset is much + larger than the cluster's aggregate object-store memory, enable + :ref:`disk-based shuffle ` by setting + ``DataContext.use_disk_based_hash_shuffle = True`` or the + environment variable ``RAY_DATA_ENABLE_DISK_SHUFFLE=1``. + Examples: .. testcode:: From 05fd0e2bc446eacf3765f52d795243b4d189f136 Mon Sep 17 00:00:00 2001 From: You-Cheng Lin Date: Fri, 25 Sep 2026 14:00:08 +0800 Subject: [PATCH 3/4] remove italic Signed-off-by: You-Cheng Lin --- doc/source/data/internals.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/doc/source/data/internals.md b/doc/source/data/internals.md index bd503aea5c23..8e89edfd80bc 100644 --- a/doc/source/data/internals.md +++ b/doc/source/data/internals.md @@ -108,7 +108,7 @@ Shuffle v2 doesn't yet support {meth}`Dataset.sort ` or { ##### Disk-based shuffle -Shuffle v2 can optionally transport intermediate shuffle data through node-local disk files instead of the Ray object store. In this *disk-based shuffle* mode, shuffle bytes never enter the object store; only small file-handle metadata travels through it. +Shuffle v2 can optionally transport intermediate shuffle data through node-local disk files instead of the Ray object store. In this disk-based shuffle mode, shuffle bytes never enter the object store; only small file-handle metadata travels through it. Disk-based shuffle works as follows: From 57cf0e6e824fa4336d112be80f945f7b96a6a2c0 Mon Sep 17 00:00:00 2001 From: You-Cheng Lin <106612301+owenowenisme@users.noreply.github.com> Date: Thu, 24 Sep 2026 23:05:47 -0700 Subject: [PATCH 4/4] Apply batched suggestions from code review Co-authored-by: gemini-code-assist[bot] <176961590+gemini-code-assist[bot]@users.noreply.github.com> Signed-off-by: You-Cheng Lin <106612301+owenowenisme@users.noreply.github.com> --- python/ray/data/dataset.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/python/ray/data/dataset.py b/python/ray/data/dataset.py index 9ddba0a0f2bd..7db77b4dfac9 100644 --- a/python/ray/data/dataset.py +++ b/python/ray/data/dataset.py @@ -2339,7 +2339,7 @@ def repartition( Repartitioning with ``keys`` hash-shuffles the whole dataset. If the dataset is much larger than the cluster's aggregate object-store memory, enable :ref:`disk-based shuffle ` by - setting ``DataContext.use_disk_based_hash_shuffle = True`` or the + setting ``ray.data.DataContext.get_current().use_disk_based_hash_shuffle = True`` or the environment variable ``RAY_DATA_ENABLE_DISK_SHUFFLE=1``. Examples: @@ -3626,7 +3626,7 @@ def join( Joins hash-shuffle both input datasets by key. If the inputs are much larger than the cluster's aggregate object-store memory, enable :ref:`disk-based shuffle ` by setting - ``DataContext.use_disk_based_hash_shuffle = True`` or the + ``ray.data.DataContext.get_current().use_disk_based_hash_shuffle = True`` or the environment variable ``RAY_DATA_ENABLE_DISK_SHUFFLE=1``. Args: @@ -3811,7 +3811,7 @@ def groupby( Grouping hash-shuffles the dataset by key. If the dataset is much larger than the cluster's aggregate object-store memory, enable :ref:`disk-based shuffle ` by setting - ``DataContext.use_disk_based_hash_shuffle = True`` or the + ``ray.data.DataContext.get_current().use_disk_based_hash_shuffle = True`` or the environment variable ``RAY_DATA_ENABLE_DISK_SHUFFLE=1``. Examples: