Shuffling data#
When you consume or iterate over a Ray Data Dataset, shuffling the order of the data can be useful, for example to randomize the ingest order during ML training. This guide describes several methods for shuffling data with Ray Data and the trade-offs of each.
Choose a shuffling method#
Ray Data provides several options for shuffling data. Each option trades off the granularity of shuffle control against memory consumption and runtime. The following sections describe each option and its trade-offs. Choose the method that fits your use case.
Shuffle the ordering of files#
To randomly shuffle the ordering of input files before reading, call a read function that supports shuffling, such as read_images(), and pass the shuffle="files" parameter. This randomly assigns input files to workers for reading.
This is the fastest shuffle option because it’s purely a metadata operation. Ray Data randomly shuffles the list of files that make up the dataset before read tasks fetch them. However, this option doesn’t shuffle the rows inside each file, so the randomness might not be sufficient for your needs if your files have a large number of rows.
import ray
ds = ray.data.read_images(
"s3://anonymous@ray-example-data/image-datasets/simple",
shuffle="files",
)
Shuffle rows with a local buffer#
To shuffle a subset of rows locally while you iterate with methods such as iter_batches(), iter_torch_batches(), and iter_tf_batches(), specify local_shuffle_buffer_size.
This option shuffles up to local_shuffle_buffer_size rows buffered during iteration. For more details, see Iterate over batches with shuffling.
This option is slower than file order shuffling, and it shuffles rows locally without network transfer. You can combine the local shuffle buffer with file order shuffling. See Shuffle the ordering of files.
import ray
ds = ray.data.read_images("s3://anonymous@ray-example-data/image-datasets/simple")
for batch in ds.iter_batches(
batch_size=2,
batch_format="numpy",
local_shuffle_buffer_size=250,
):
print(batch)
Tip
If throughput drops when you use local_shuffle_buffer_size, check the total time spent in batch creation. In the ds.stats() output, find In batch formatting under Batch iteration time breakdown. If this time is much larger than the time spent in other steps, decrease local_shuffle_buffer_size, or turn off the local shuffle buffer and only shuffle the ordering of files.
Shuffle rows with map_batches#
To shuffle data as a separate stage, use map_batches() with a shuffle function that randomly permutes the rows within each batch. This approach has the following advantages over local buffer shuffle:
It decouples shuffling from the iterator. The shuffle runs as a separate Ray Data operator that doesn’t block downstream CPU or GPU processing.
Ray Data’s resource management automatically schedules shuffle tasks based on the available CPU and memory in the cluster, which avoids resource contention.
The shuffle work can run in parallel across multiple machines, so this approach scales better for large datasets.
The batch_size parameter controls the shuffle window. A larger value shuffles more rows together for better randomness, but it requires more memory.
Important
To avoid out-of-memory errors, always set the memory parameter when you use large batch sizes. Estimate the value as batch_size * row_bytes.
import numpy as np
import pyarrow as pa
import ray
def random_shuffle(batch: pa.Table) -> pa.Table:
indices = np.random.permutation(len(batch))
return batch.take(indices)
row_bytes = 4096
shuffle_memory = int(2**30) # 1 GB shuffle window
batch_size = int(shuffle_memory / row_bytes)
ds = ray.data.range(1000)
ds = ds.map_batches(
random_shuffle,
batch_size=batch_size,
batch_format="pyarrow",
memory=shuffle_memory,
)
ds.take(10)
Tip
Combine map_batches shuffle with file order shuffling for additional randomness. File order shuffling randomizes which files Ray Data reads first, while map_batches shuffle randomizes rows within each shuffle window.
How does local buffer shuffle compare to map_batches shuffle?#
The following benchmark compares steady-state training throughput for local buffer shuffle and map_batches shuffle on a synthetic workload. The workload uses ray.data.range_tensor with about 4 KB per row, four GPU workers, a batch size of 4096, and 200 steps with 100 warmup steps.
Method |
Throughput (rows/s) |
% of baseline |
|---|---|---|
No shuffle (baseline) |
1,759,282 |
100% |
Local buffer shuffle 1 GB |
225,181 |
13% |
Local buffer shuffle 2 GB |
220,644 |
13% |
Local buffer shuffle 3 GB |
153,256 |
9% |
|
1,400,734 |
80% |
|
1,460,037 |
83% |
|
1,588,428 |
90% |
Randomize block order#
This option randomizes the order of blocks in a dataset. The operation alone doesn’t involve heavy computation or communication, but Ray Data must materialize all blocks in memory before it randomizes their order in the queue for the subsequent operation.
Note
By default, Ray Data doesn’t guarantee any particular block order when it reads blocks from different files in parallel, unless you set DataContext.execution_options.preserve_order to true. As a result, this option is mainly relevant when Ray Data yields blocks from a relatively small set of large files.
Note
Use this option only when your dataset is small enough to fit in object store memory.
To shuffle the block order, use randomize_block_order.
import ray
ds = ray.data.read_text(
"s3://anonymous@ray-example-data/sms_spam_collection_subset.txt"
)
# Randomize the block order of this dataset.
ds = ds.randomize_block_order()
Shuffle all rows globally#
Ray Data provides the following options for shuffling all rows globally across the whole dataset:
Random shuffling: Call
random_shuffle()to shuffle individual rows from the existing blocks into new blocks. You can optionally provide a seed.Key-based repartitioning: Call
repartition()with thekeysparameter to shuffle the rows based on the hash of the values in the key columns you provide. This operation co-locates rows with the same key values deterministically. Ray 2.46 introduced this option.
A shuffle is an expensive operation. It requires materializing the whole dataset in memory, and it acts as a synchronization barrier, so subsequent operators can’t start executing until the shuffle completes.
The following example shuffles rows randomly with a seed:
import ray
ds = ray.data.read_images("s3://anonymous@ray-example-data/image-datasets/simple")
# Random shuffle with seed
random_shuffled_ds = ds.random_shuffle(seed=123)
The following example hash shuffles rows based on the id column:
import ray
hash_shuffled_ds = ds.repartition(keys="id", num_blocks=200)
Tip
By default, key-based repartitioning uses shuffle v2, which is ShuffleStrategy.SHUFFLE_V2. For the available settings, see Tune shuffle v2.
To fall back to the previous hash-shuffle implementation, set DataContext.shuffle_strategy to ShuffleStrategy.HASH_SHUFFLE:
from ray.data.context import DataContext, ShuffleStrategy
DataContext.get_current().shuffle_strategy = ShuffleStrategy.HASH_SHUFFLE
Advanced: Optimize shuffles#
Note
Shuffle optimization is an active area of development. If your dataset uses a shuffle operation and you’re having trouble configuring the shuffle, file a Ray Data issue on GitHub.
When should you use global per-epoch shuffling?#
Use global per-epoch shuffling only if your model is sensitive to the randomness of the training data. According to a theoretical foundation, all gradient-descent-based model trainers benefit from improved global shuffle quality. In practice, the benefit is particularly pronounced for tabular data and models. However, the more global the shuffle, the more expensive the shuffling operation. Data transfer costs compound this increase in distributed data-parallel training on a multi-node cluster. This cost can be prohibitive for large datasets.
To find the best tradeoff between preprocessing time and cost and per-epoch shuffle quality, measure the precision gain per training step for your model under different shuffling policies, such as no shuffling, local shuffling, or global shuffling.
As long as your data loading and shuffling throughput is higher than your training throughput, your GPU should saturate. If your model is shuffle-sensitive, push the shuffle quality higher until you reach this threshold.
Enable push-based shuffle#
Note
DataContext.use_push_based_shuffle and the RAY_DATA_PUSH_BASED_SHUFFLE environment variable are deprecated. Select push-based shuffle with DataContext.shuffle_strategy or the RAY_DATA_DEFAULT_SHUFFLE_STRATEGY environment variable instead, as this section shows.
Some Dataset operations require a shuffle operation, which shuffles data from all of the input partitions to all of the output partitions. These operations include Dataset.random_shuffle, Dataset.sort, and Dataset.groupby. For example, a sort operation reorders data between blocks, so it requires shuffling across partitions. Shuffling can be hard to scale to large data sizes and clusters, especially when the total dataset size doesn’t fit in memory.
Ray Data provides an alternative shuffle implementation called push-based shuffle to improve large-scale performance. Try it if your dataset has more than 1,000 blocks or is larger than 1 TB.
To try it locally or on a cluster, start with the nightly release test that Ray runs for Dataset.random_shuffle. The following chart shows run time results for Dataset.random_shuffle on 1 to 10 TB of data, which gives an idea of the performance you can expect. The benchmark ran on 20 m5.4xlarge AWS EC2 instances, each with 16 vCPUs and 64 GB of RAM.
To try push-based shuffle, set the RAY_DATA_DEFAULT_SHUFFLE_STRATEGY=sort_shuffle_push_based environment variable when you run your application:
wget https://raw.githubusercontent.com/ray-project/ray/master/release/nightly_tests/dataset/random_shuffle_benchmark.py
RAY_DATA_DEFAULT_SHUFFLE_STRATEGY=sort_shuffle_push_based python random_shuffle_benchmark.py --num-partitions=10 --partition-size=1e7
You can also set the shuffle strategy while your program runs with DataContext.shuffle_strategy:
import ray
from ray.data.context import ShuffleStrategy
ctx = ray.data.DataContext.get_current()
ctx.shuffle_strategy = ShuffleStrategy.SORT_SHUFFLE_PUSH_BASED
ds = (
ray.data.range(1000)
.random_shuffle()
)
Large-scale shuffles can take a while to finish. For debugging, you can execute only part of a shuffle, so that you can collect an execution profile more quickly. The following example limits a random shuffle operation to two output blocks:
import ray
ctx = ray.data.DataContext.get_current()
ctx.set_config(
"debug_limit_shuffle_execution_to_num_blocks", 2
)
ds = (
ray.data.range(1000, override_num_blocks=10)
.random_shuffle()
.materialize()
)
print(ds.stats())
Operator 1 ReadRange->RandomShuffle: executed in 0.08s
Suboperator 0 ReadRange->RandomShuffleMap: 2/2 blocks executed
...