Ray Data: Scalable data processing for AI workloads#
Ray Data is a scalable data processing library for AI workloads, built on Ray. It provides APIs for common operations such as batch inference, data preprocessing, and data loading for ML training. Unlike other distributed data systems, Ray Data uses a streaming execution engine to process large datasets efficiently and keep utilization high across both CPU and GPU workloads.
Quickstart#
To learn more about installing Ray and its libraries, see Installing Ray. To install Ray Data, run the following command:
pip install -U 'ray[data]'
The following example runs a batch text classification task with Ray Data:
import ray
import pandas as pd
class ClassificationModel:
def __init__(self):
from transformers import pipeline
self.pipe = pipeline("text-classification")
def __call__(self, batch: pd.DataFrame):
results = self.pipe(list(batch["text"]))
result_df = pd.DataFrame(results)
return pd.concat([batch, result_df], axis=1)
ds = ray.data.read_text("s3://anonymous@ray-example-data/sms_spam_collection_subset.txt")
ds = ds.map_batches(
ClassificationModel,
compute=ray.data.ActorPoolStrategy(size=2),
batch_size=64,
batch_format="pandas"
# num_gpus=1 # this will set 1 GPU per worker
)
ds.show(limit=1)
{'text': 'ham\tGo until jurong point, crazy.. Available only in bugis n great world la e buffet... Cine there got amore wat...', 'label': 'NEGATIVE', 'score': 0.9935141801834106}
Why choose Ray Data?#
AI workloads revolve around deep learning models, which are computationally intensive and often require specialized hardware such as GPUs. Unlike CPUs, GPUs often have less memory, different scheduling semantics, and a much higher cost to run. Systems built for traditional data processing pipelines often don’t use these resources well.
Ray Data treats AI workloads as a first-class use case and offers four advantages:
Faster and cheaper for deep learning: Ray Data streams data between CPU preprocessing tasks and GPU inference or training tasks. Keeping GPUs active maximizes resource utilization and reduces costs.
Framework-friendly: Ray Data integrates with common AI frameworks such as vLLM, PyTorch, Hugging Face, and TensorFlow, and with common cloud providers such as AWS, GCP, and Azure.
Support for multimodal data: Ray Data uses Apache Arrow and pandas and supports many data formats used in ML workloads, such as Parquet, Lance, images, JSON, CSV, audio, and video.
Scalable by default: Ray Data builds on Ray to scale automatically across heterogeneous clusters of CPU and GPU machines. The same code runs unchanged on one machine or on hundreds of nodes processing hundreds of TB of data.
Learn more#
Quickstart
Run a basic example to get started with Ray Data.
Key concepts
Learn the key concepts behind Ray Data, including what a Dataset is and how to use it.
User guides
Learn how to use Ray Data, from basic usage to end-to-end guides.
Examples
Find basic and scaled-out examples of Ray Data workloads.
API
Get in-depth information about the Ray Data API.
Case studies#
The following case studies use Ray Data for training ingest:
Pinterest uses Ray Data to do last mile data processing for model training.
Instacart builds distributed machine learning model training on Ray Data.
The following case studies use Ray Data for batch inference: