Ray Data

orchestra-research/ai-research-skills/05-data-processing/ray-data

by orchestra-research773a52944ba4MIT13K starsListed Oct 8, 2026Updated Oct 8, 2026Repository updated 3 months ago

Scalable data processing for ML workloads. Streaming execution across CPU/GPU, supports Parquet/CSV/JSON/images. Integrates with Ray Train, PyTorch, TensorFlow. Scales from single machine to 100s of nodes. Use for batch inference, data preprocessing, multi-modal data loading, or distributed ETL pipelines.

AI-generated overview

Guides distributed data processing with Ray Data for ML pipelines, covering loading, transforming, and writing large datasets.

What it does
This skill provides instructions and code examples for using Ray Data to read, transform, and write large datasets across a cluster. It covers streaming execution, batch and row transformations, filtering, group-by aggregation, GPU-accelerated preprocessing, and batch inference. It also shows integration with Ray Train, PyTorch, and TensorFlow, plus performance tuning such as repartitioning and batch size adjustment.
When to use it
Use it when processing datasets larger than memory or over 100GB for ML training, distributed preprocessing, batch inference, multi-modal data loading, or distributed ETL pipelines. It is also relevant when scaling data processing from a single machine to a cluster.
Requirements
Requires the ray[data] package along with pyarrow and pandas, and Python code execution. Cluster or cloud storage access (such as S3 or GCS) is needed for distributed or remote data. It ships no scripts; it is instructions only.

Ray Data - Scalable ML Data Processing

Distributed data processing library for ML and AI workloads.

When to use Ray Data

Use Ray Data when:

  • Processing large datasets (>100GB) for ML training
  • Need distributed data preprocessing across cluster
  • Building batch inference pipelines
  • Loading multi-modal data (images, audio, video)
  • Scaling data processing from laptop to cluster

Key features:

  • Streaming execution: Process data larger than memory
  • GPU support: Accelerate transforms with GPUs
  • Framework integration: PyTorch, TensorFlow, HuggingFace
  • Multi-modal: Images, Parquet, CSV, JSON, audio, video

Use alternatives instead:

  • Pandas: Small data (<1GB) on single machine
  • Dask: Tabular data, SQL-like operations
  • Spark: Enterprise ETL, SQL queries

Quick start

Installation

bash
pip install -U 'ray[data]'

Load and transform data

python
import ray
# Read Parquet filesds = ray.data.read_parquet("s3://bucket/data/*.parquet")
# Transform data (lazy execution)ds = ds.map_batches(lambda batch: {"processed": batch["text"].str.lower()})
# Consume datafor batch in ds.iter_batches(batch_size=100):    print(batch)

Integration with Ray Train

python
import rayfrom ray.train import ScalingConfigfrom ray.train.torch import TorchTrainer
# Create datasettrain_ds = ray.data.read_parquet("s3://bucket/train/*.parquet")
def train_func(config):    # Access dataset in training    train_ds = ray.train.get_dataset_shard("train")
    for epoch in range(10):        for batch in train_ds.iter_batches(batch_size=32):            # Train on batch            pass
# Train with Raytrainer = TorchTrainer(    train_func,    datasets={"train": train_ds},    scaling_config=ScalingConfig(num_workers=4, use_gpu=True))trainer.fit()

Reading data

From cloud storage

python
import ray
# Parquet (recommended for ML)ds = ray.data.read_parquet("s3://bucket/data/*.parquet")
# CSVds = ray.data.read_csv("s3://bucket/data/*.csv")
# JSONds = ray.data.read_json("gs://bucket/data/*.json")
# Imagesds = ray.data.read_images("s3://bucket/images/")

From Python objects

python
# From listds = ray.data.from_items([{"id": i, "value": i * 2} for i in range(1000)])
# From rangeds = ray.data.range(1000000)  # Synthetic data
# From pandasimport pandas as pddf = pd.DataFrame({"col1": [1, 2, 3], "col2": [4, 5, 6]})ds = ray.data.from_pandas(df)

Transformations

Map batches (vectorized)

python
# Batch transformation (fast)def process_batch(batch):    batch["doubled"] = batch["value"] * 2    return batch
ds = ds.map_batches(process_batch, batch_size=1000)

Row transformations

python
# Row-by-row (slower)def process_row(row):    row["squared"] = row["value"] ** 2    return row
ds = ds.map(process_row)

Filter

python
# Filter rowsds = ds.filter(lambda row: row["value"] > 100)

Group by and aggregate

python
# Group by columnds = ds.groupby("category").count()
# Custom aggregationds = ds.groupby("category").map_groups(lambda group: {"sum": group["value"].sum()})

GPU-accelerated transforms

python
# Use GPU for preprocessingdef preprocess_images_gpu(batch):    import torch    images = torch.tensor(batch["image"]).cuda()    # GPU preprocessing    processed = images * 255    return {"processed": processed.cpu().numpy()}
ds = ds.map_batches(    preprocess_images_gpu,    batch_size=64,    num_gpus=1  # Request GPU)

Writing data

python
# Write to Parquetds.write_parquet("s3://bucket/output/")
# Write to CSVds.write_csv("output/")
# Write to JSONds.write_json("output/")

Performance optimization

Repartition

python
# Control parallelismds = ds.repartition(100)  # 100 blocks for 100-core cluster

Batch size tuning

python
# Larger batches = faster vectorized opsds.map_batches(process_fn, batch_size=10000)  # vs batch_size=100

Streaming execution

python
# Process data larger than memoryds = ray.data.read_parquet("s3://huge-dataset/")for batch in ds.iter_batches(batch_size=1000):    process(batch)  # Streamed, not loaded to memory

Common patterns

Batch inference

python
import ray
# Load modeldef load_model():    # Load once per worker    return MyModel()
# Inference functionclass BatchInference:    def __init__(self):        self.model = load_model()
    def __call__(self, batch):        predictions = self.model(batch["input"])        return {"prediction": predictions}
# Run distributed inferenceds = ray.data.read_parquet("s3://data/")predictions = ds.map_batches(BatchInference, batch_size=32, num_gpus=1)predictions.write_parquet("s3://output/")

Data preprocessing pipeline

python
# Multi-step pipelineds = (    ray.data.read_parquet("s3://raw/")    .map_batches(clean_data)    .map_batches(tokenize)    .map_batches(augment)    .write_parquet("s3://processed/"))

Integration with ML frameworks

PyTorch

python
# Convert to PyTorchtorch_ds = ds.to_torch(label_column="label", batch_size=32)
for batch in torch_ds:    # batch is dict with tensors    inputs, labels = batch["features"], batch["label"]

TensorFlow

python
# Convert to TensorFlowtf_ds = ds.to_tf(feature_columns=["image"], label_column="label", batch_size=32)
for features, labels in tf_ds:    # Train model    pass

Supported data formats

FormatReadWriteUse Case
Parquet✅✅ML data (recommended)
CSV✅✅Tabular data
JSON✅✅Semi-structured
Images✅❌Computer vision
NumPy✅✅Arrays
Pandas✅❌DataFrames

Performance benchmarks

Scaling (processing 100GB data):

  • 1 node (16 cores): ~30 minutes
  • 4 nodes (64 cores): ~8 minutes
  • 16 nodes (256 cores): ~2 minutes

GPU acceleration (image preprocessing):

  • CPU only: 1,000 images/sec
  • 1 GPU: 5,000 images/sec
  • 4 GPUs: 18,000 images/sec

Use cases

Production deployments:

  • Pinterest: Last-mile data processing for model training
  • ByteDance: Scaling offline inference with multi-modal LLMs
  • Spotify: ML platform for batch inference

References

  • Transformations Guide [blocked] - Map, filter, groupby operations
  • Integration Guide [blocked] - Ray Train, PyTorch, TensorFlow

Resources

Source and attribution

Source:orchestra-research/ai-research-skillsin05-data-processing/ray-dataat commit773a529

License: MIT

Content belongs to its original authors. SourceWeft indexes it from a public repository.

Report or request removal