{"id":"ray-data","name":"ray-data","summary":"MLワークロード向けのスケーラブルなデータ処理。CPU/GPUをまたぐストリーミング実行はParquet/CSV/JSON/イメージをサポートしています。","body":"# Ray Data - Scalable ML Data Processing\n\nDistributed data processing library for ML and AI workloads.\n\n## When to use Ray Data\n\n**Use Ray Data when:**\n- Processing large datasets (>100GB) for ML training\n- Need distributed data preprocessing across cluster\n- Building batch inference pipelines\n- Loading multi-modal data (images, audio, video)\n- Scaling data processing from laptop to cluster\n\n**Key features**:\n- **Streaming execution**: Process data larger than memory\n- **GPU support**: Accelerate transforms with GPUs\n- **Framework integration**: PyTorch, TensorFlow, HuggingFace\n- **Multi-modal**: Images, Parquet, CSV, JSON, audio, video\n\n**Use alternatives instead**:\n- **Pandas**: Small data (<1GB) on single machine\n- **Dask**: Tabular data, SQL-like operations\n- **Spark**: Enterprise ETL, SQL queries\n\n## Quick start\n\n### Installation\n\n```bash\npip install -U 'ray[data]'\n```\n\n### Load and transform data\n\n```python\nimport ray\n\n# Read Parquet files\nds = ray.data.read_parquet(\"s3://bucket/data/*.parquet\")\n\n# Transform data (lazy execution)\nds = ds.map_batches(lambda batch: {\"processed\": batch[\"text\"].str.lower()})\n\n# Consume data\nfor batch in ds.iter_batches(batch_size=100):\n    print(batch)\n```\n\n### Integration with Ray Train\n\n```python\nimport ray\nfrom ray.train import ScalingConfig\nfrom ray.train.torch import TorchTrainer\n\n# Create dataset\ntrain_ds = ray.data.read_parquet(\"s3://bucket/train/*.parquet\")\n\ndef train_func(config):\n    # Access dataset in training\n    train_ds = ray.train.get_dataset_shard(\"train\")\n\n    for epoch in range(10):\n        for batch in train_ds.iter_batches(batch_size=32):\n            # Train on batch\n            pass\n\n# Train with Ray\ntrainer = TorchTrainer(\n    train_func,\n    datasets={\"train\": train_ds},\n    scaling_config=ScalingConfig(num_workers=4, use_gpu=True)\n)\ntrainer.fit()\n```\n\n## Reading data\n\n### From cloud storage\n\n```python\nimport ray\n\n# Parquet (recommended for ML)\nds = ray.data.read_parquet(\"s3://bucket/data/*.parquet\")\n\n# CSV\nds = ray.data.read_csv(\"s3://bucket/data/*.csv\")\n\n# JSON\nds = ray.data.read_json(\"gs://bucket/data/*.json\")\n\n# Images\nds = ray.data.read_images(\"s3://bucket/images/\")\n```\n\n### From Python objects\n\n```python\n# From list\nds = ray.data.from_items([{\"id\": i, \"value\": i * 2} for i in range(1000)])\n\n# From range\nds = ray.data.range(1000000)  # Synthetic data\n\n# From pandas\nimport pandas as pd\ndf = pd.DataFrame({\"col1\": [1, 2, 3], \"col2\": [4, 5, 6]})\nds = ray.data.from_pandas(df)\n```\n\n## Transformations\n\n### Map batches (vectorized)\n\n```python\n# Batch transformation (fast)\ndef process_batch(batch):\n    batch[\"doubled\"] = batch[\"value\"] * 2\n    return batch\n\nds = ds.map_batches(process_batch, batch_size=1000)\n```\n\n### Row transformations\n\n```python\n# Row-by-row (slower)\ndef process_row(row):\n    row[\"squared\"] = row[\"value\"] ** 2\n    return row\n\nds = ds.map(process_row)\n```\n\n### Filter\n\n```python\n# Filter rows\nds = ds.filter(lambda row: row[\"value\"] > 100)\n```\n\n### Group by and aggregate\n\n```python\n# Group by column\nds = ds.groupby(\"category\").count()\n\n# Custom aggregation\nds = ds.groupby(\"category\").map_groups(lambda group: {\"sum\": group[\"value\"].sum()})\n```\n\n## GPU-accelerated transforms\n\n```python\n# Use GPU for preprocessing\ndef preprocess_images_gpu(batch):\n    import torch\n    images = torch.tensor(batch[\"image\"]).cuda()\n    # GPU preprocessing\n    processed = images * 255\n    return {\"processed\": processed.cpu().numpy()}\n\nds = ds.map_batches(\n    preprocess_images_gpu,\n    batch_size=64,\n    num_gpus=1  # Request GPU\n)\n```\n\n## Writing data\n\n```python\n# Write to Parquet\nds.write_parquet(\"s3://bucket/output/\")\n\n# Write to CSV\nds.write_csv(\"output/\")\n\n# Write to JSON\nds.write_json(\"output/\")\n```\n\n## Performance optimization\n\n### Repartition\n\n```python\n# Control parallelism\nds = ds.repartition(100)  # 100 blocks for 100-core cluster\n```\n\n### Batch size tuning\n\n```python\n# Larger batches = faster vectorized ops\nds.map_batches(process_fn, batch_size=10000)  # vs batch_size=100\n```\n\n### Streaming execution\n\n```python\n# Process data larger than memory\nds = ray.data.read_parquet(\"s3://huge-dataset/\")\nfor batch in ds.iter_batches(batch_size=1000):\n    process(batch)  # Streamed, not loaded to memory\n```\n\n## Common patterns\n\n### Batch inference\n\n```python\nimport ray\n\n# Load model\ndef load_model():\n    # Load once per worker\n    return MyModel()\n\n# Inference function\nclass BatchInference:\n    def __init__(self):\n        self.model = load_model()\n\n    def __call__(self, batch):\n        predictions = self.model(batch[\"input\"])\n        return {\"prediction\": predictions}\n\n# Run distributed inference\nds = ray.data.read_parquet(\"s3://data/\")\npredictions = ds.map_batches(BatchInference, batch_size=32, num_gpus=1)\npredictions.write_parquet(\"s3://output/\")\n```\n\n### Data preprocessing pipeline\n\n```python\n# Multi-step pipeline\nds = (\n    ray.data.read_parquet(\"s3://raw/\")\n    .map_batches(clean_data)\n    .map_batches(tokenize)\n    .map_batches(augment)\n    .write_parquet(\"s3://processed/\")\n)\n```\n\n## Integration with ML frameworks\n\n### PyTorch\n\n```python\n# Convert to PyTorch\ntorch_ds = ds.to_torch(label_column=\"label\", batch_size=32)\n\nfor batch in torch_ds:\n    # batch is dict with tensors\n    inputs, labels = batch[\"features\"], batch[\"label\"]\n```\n\n### TensorFlow\n\n```python\n# Convert to TensorFlow\ntf_ds = ds.to_tf(feature_columns=[\"image\"], label_column=\"label\", batch_size=32)\n\nfor features, labels in tf_ds:\n    # Train model\n    pass\n```\n\n## Supported data formats\n\n| Format | Read | Write | Use Case |\n|--------|------|-------|----------|\n| Parquet | ✅ | ✅ | ML data (recommended) |\n| CSV | ✅ | ✅ | Tabular data |\n| JSON | ✅ | ✅ | Semi-structured |\n| Images | ✅ | ❌ | Computer vision |\n| NumPy | ✅ | ✅ | Arrays |\n| Pandas | ✅ | ❌ | DataFrames |\n\n## Performance benchmarks\n\n**Scaling** (processing 100GB data):\n- 1 node (16 cores): ~30 minutes\n- 4 nodes (64 cores): ~8 minutes\n- 16 nodes (256 cores): ~2 minutes\n\n**GPU acceleration** (image preprocessing):\n- CPU only: 1,000 images/sec\n- 1 GPU: 5,000 images/sec\n- 4 GPUs: 18,000 images/sec\n\n## Use cases\n\n**Production deployments**:\n- **Pinterest**: Last-mile data processing for model training\n- **ByteDance**: Scaling offline inference with multi-modal LLMs\n- **Spotify**: ML platform for batch inference\n\n## References\n\n- **[Transformations Guide](references/transformations.md)** - Map, filter, groupby operations\n- **[Integration Guide](references/integration.md)** - Ray Train, PyTorch, TensorFlow\n\n## Resources\n\n- **Docs**: https://docs.ray.io/en/latest/data/data.html\n- **GitHub**: https://github.com/ray-project/ray ⭐ 36,000+\n- **Version**: Ray 2.40.0+\n- **Examples**: https://docs.ray.io/en/latest/data/examples/overview.html","author":"@Orchestra-Research","ownerProfile":null,"authorContacts":null,"sourceUrl":"https://github.com/Orchestra-Research/AI-Research-SKILLs/tree/main/05-data-processing/ray-data","license":"MIT","category":"coding","lang":"en","tokens":1777,"stars":0,"calls30d":1,"claimed":false,"visibility":"public","origin":"crawler","version":"0.1.0","createdAt":"2026-08-22","updatedAt":"2026-08-22","files":[{"path":"references/integration.md","size":1851,"sha256":"d182e025b0396779442f16b6461c0bf140a337635051cfb8472698a18f4099ca"},{"path":"references/transformations.md","size":1664,"sha256":"afc01ade53bc95de6dec6de614b52c6f92eb4a0b141740d937b5dd1855de380d"}],"requires":{"mcp":[],"tools":[]},"safety":{"flags":[],"scannedAt":"2026-08-22","hasScripts":false,"networkEndpoints":["docs.ray.io"]}}