Ray Data

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

作者 orchestra-research773a52944ba4MIT13K 個星標收錄於 2026年10月8日更新於 2026年10月8日儲存庫3 個月前更新

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 產生的概覽

指導使用 Ray Data 進行面向機器學習工作負載的分散式資料處理,涵蓋大型資料集的讀取、轉換與寫入。

功能
此技能提供使用 Ray Data 在叢集上讀取、轉換與寫入大型資料集的說明與程式碼範例。內容涵蓋串流執行、批次與逐列轉換、篩選、分組聚合、GPU 加速前處理以及批次推論。它也展示與 Ray Train、PyTorch 和 TensorFlow 的整合,以及重新分割與批次大小調整等效能調校方式。
適用情境
適用於處理超出記憶體或超過 100GB 的機器學習訓練資料集、分散式前處理、批次推論、多模態資料載入或分散式 ETL 管線。也適合將資料處理從單機擴展到叢集的場景。
執行需求
需要 ray[data] 套件以及 pyarrow 和 pandas,並能執行 Python 程式碼。分散式或遠端資料需要叢集或雲端儲存存取(如 S3 或 GCS)。此技能不附帶指令碼,僅為說明文件。

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

來源與署名

來源:orchestra-research/ai-research-skills位於05-data-processing/ray-data提交773a529

授權條款: MIT

內容歸原作者所有。SourceWeft 從公開儲存庫中收錄這些內容。

檢舉或申請下架