Ray + RayDP: Spark ETL → Ray Actor для distributed ML, сравнение с MLlib

Ray + RayDP: Spark ETL → Ray Actor для distributed ML, сравнение с MLlib

optimization

Архитектурный кризис: почему Spark MLlib устарел

Историческая попытка сделать всё в JVM

Когда Spark в 2013–2015 годах завоёвывал рынок, он нёс обещание: одна платформа для всего - и ETL, и Machine Learning. MLlib появился как ответ на это обещание: набор распределённых алгоритмов ML, реализованных поверх Spark RDD, а затем и DataFrame API. Он умел делать Logistic Regression, Random Forest, K-Means clustering, ALS для рекомендаций - и всё это масштабировалось горизонтально на десятки нод.

Прошло десять лет. И реальность оказалась другой. MLlib не исчез - он используется. Но он перестал быть центром ML-экосистемы.

Почему MLlib проиграл современному ML/DL

Проблема 1: JVM как барьер для Python-экосистемы.

Современный ML - это Python. PyTorch, TensorFlow, Hugging Face, scikit-learn, XGBoost, LightGBM - всё это написано на Python/C++. Spark - это JVM. Взаимодействие между ними происходит через Py4J - мост, который сериализует каждый вызов через сокет. Когда нужно передать тензор размером 1 GB из Python-процесса в JVM Executor - это болезненно медленно.

Проблема 2: итеративные алгоритмы против DAG-модели.

Обучение нейронных сетей - итеративный процесс. Один эпох - это прямой проход, обратное распространение ошибки, обновление весов. Повторить 100 раз. Spark построен на DAG-модели: каждая операция - это трансформация, которая материализуется при action. Для одного прохода по данным это отлично. Для 100 итераций с обновлением состояния модели - это создаёт колоссальный overhead на каждый новый DAG.

Проблема 3: GPU-ресурсы и их оркестрация.

Обучение нейронных сетей требует GPU. Spark умеет аллоцировать GPU для SQL-операторов (через RAPIDS), но управлять GPU-ресурсами для тренировки PyTorch-модели - другая задача. Spark Executor получает GPU, но он не понимает концепций gradient synchronization, all-reduce операций между GPU, динамического перераспределения GPU-памяти между эпохами.

Проблема 4: stateless Executor против stateful Actor.

Spark Executor - stateless: каждый таск получает данные, делает преобразование, возвращает результат. Нет состояния между тасками. Но ML-модель в процессе обучения - это состояние. Веса модели живут между батчами и эпохами. В Spark это состояние нужно сериализовывать, передавать, десериализовывать - что полностью разрушает идею инкрементального обучения.

Рождение Ray

Ray - это открытая платформа распределённых вычислений, разработанная в Berkeley RISELab (2017) и написанная на C++ (ядро) и Python (API). Ray не пытается заменить Spark - он решает другую задачу: гибкое распределённое исполнение Python-кода с поддержкой статeful акторов, динамических графов вычислений и первоклассной интеграцией с PyTorch, TensorFlow, XGBoost.

Именно Ray стал стандартом для:

  • Распределённого обучения нейронных сетей (Ray Train)
  • Подбора гиперпараметров (Ray Tune)
  • Distributed inference больших LLM (vLLM работает на Ray)
  • Reinforcement Learning (Ray RLlib)

Но Ray плохо умеет то, в чём Spark силён: SQL-запросы с join, GROUP BY, Catalyst optimizer, работа с lakehouse-форматами (Iceberg, Delta). Так возникла потребность в мосте между двумя мирами.


Мост между мирами: что такое RayDP

Концепция

RayDP (Ray Data Processing) - это open-source библиотека, позволяющая запускать Apache Spark внутри Ray кластера. Вместо того чтобы поднимать отдельный Spark YARN-кластер или Spark on K8s, вы инициализируете SparkSession как распределённое приложение поверх существующего Ray-кластера.

Это инверсия контроля (Inversion of Control): обычно Ray запускают как инструмент анализа данных поверх Spark-кластера. RayDP переворачивает это - Spark запускается как управляемое приложение внутри Ray.

Как Spark Executors становятся Ray Actors

Когда вы вызываете raydp.init_spark(num_executors=10, executor_cores=4), RayDP делает следующее:

  1. Создаёт SparkDriverActor - Ray Actor, внутри которого живёт JVM с Spark Driver (DAGScheduler, SparkContext).
  2. Создаёт SparkExecutorActor × N - Ray Actors, каждый из которых содержит JVM с Spark Executor.
  3. Регистрирует эти акторы как вычислительные ресурсы через Ray Resource Manager.
  4. Настраивает Spark так, чтобы он общался с Executor-ами через Ray Actor protocol вместо стандартного Netty RPC.

Результат: SparkSession работает привычно - API не меняется, PySpark-код не меняется. Но под капотом Executor-ы - это Ray Actors.

Ключевое преимущество: Zero-Copy передача данных через Arrow

После того как ETL-обработка завершена, данные нужно передать из Spark в Ray для обучения модели. Обычный способ: записать на диск (S3/HDFS), прочитать обратно. Это медленно и дорого.

RayDP позволяет передачу через общую память (shared memory) на основе Apache Arrow Plasma Object Store:

Spark DataFrame (Executor JVM heap) 
    ↓
Arrow RecordBatch (нативная память, shared via Plasma)
    ↓
Ray Dataset (Python, zero-copy read через memmap)

Plasma Object Store - это компонент Apache Arrow, реализующий разделяемую память между процессами. Один процесс записывает Arrow RecordBatch в Plasma по идентификатору (ObjectID). Другой процесс читает его по тому же ObjectID через memory-mapped файл - без копирования! Данные физически один раз в RAM, доступны всем Ray-рабочим.

Для датасета в 100 GB это означает: передача данных из Spark в Ray занимает секунды (время сериализации в Arrow и записи указателей в Plasma) вместо минут (запись на S3 + чтение).

Полная архитектура RayDP

Архитектура чётко разделяет ответственности:

  • Spark ETL Layer: SQL, JOIN, GROUP BY, feature engineering - всё, что делает Spark хорошо через Catalyst.
  • Plasma Object Store: общая память, устраняющая необходимость диска при передаче данных.
  • Ray ML Layer: Python-нативные акторы с PyTorch, XGBoost, Hugging Face - на GPU.
  • MLflow: трекинг экспериментов, версионирование моделей.

Ray: архитектура под капотом

Основные примитивы Ray

Прежде чем перейти к RayDP, нужно понять фундаментальные концепции Ray.

Ray Task - это функция, которая выполняется асинхронно в распределённой среде. Аналог Spark RDD-операции, но без DAG-модели - результат возвращается как ObjectRef (ссылка на объект в Plasma Store).

import ray

ray.init()

@ray.remote
def compute_features(batch: dict) -> dict:
    """Эта функция выполняется на одном из Ray-воркеров."""
    import numpy as np
    features = np.log1p(np.array(batch["amount"]))
    return {"log_amount": features.tolist()}

# Вызов: немедленно возвращает ObjectRef (не ждёт результата)
future = compute_features.remote({"amount": [100, 500, 1000]})

# Получение результата (блокирующий вызов)
result = ray.get(future)
print(result)  # {"log_amount": [4.615, 6.216, 6.909]}

Ray Actor - это класс, экземпляр которого живёт как долгоживущий процесс на воркер-ноде. Актор сохраняет состояние между вызовами - именно это делает его идеальным для ML: загрузил веса модели один раз, делаешь inference или обучение много раз.

@ray.remote(num_gpus=1)
class ModelTrainer:
    """Ray Actor: живёт на ноде с GPU, хранит состояние модели."""

    def __init__(self, model_config: dict):
        import torch
        self.model = torch.nn.Sequential(
            torch.nn.Linear(model_config["input_dim"], 256),
            torch.nn.ReLU(),
            torch.nn.Linear(256, model_config["output_dim"])
        ).cuda()
        self.optimizer = torch.optim.Adam(self.model.parameters(), lr=1e-3)
        self.epoch = 0

    def train_batch(self, batch_data: dict) -> float:
        """Обучение на одном батче - вызывается многократно."""
        import torch
        X = torch.tensor(batch_data["features"]).cuda()
        y = torch.tensor(batch_data["labels"]).cuda()

        self.optimizer.zero_grad()
        pred = self.model(X)
        loss = torch.nn.functional.mse_loss(pred, y)
        loss.backward()
        self.optimizer.step()
        return float(loss)

    def get_epoch(self) -> int:
        return self.epoch

    def increment_epoch(self):
        self.epoch += 1

# Создание Актора (занимает 1 GPU на одной из нод)
trainer = ModelTrainer.remote({"input_dim": 50, "output_dim": 1})

# Многократные вызовы - актор живёт, состояние сохраняется
for batch in training_batches:
    loss_ref = trainer.train_batch.remote(batch)
    loss = ray.get(loss_ref)
    print(f"Batch loss: {loss:.4f}")

Ключевое отличие от Spark Executor: Spark Executor - stateless, он не помнит ничего между тасками. Ray Actor - stateful, он помнит веса модели, оптимизатор, номер эпохи. Это принципиально для iterative ML.

Ray Object Store (Plasma)

Ray Object Store - это распределённый key-value store, хранящий объекты в native памяти нод. Ключевые свойства:

  • Shared Memory на одной ноде: если два Actor-а на одной ноде читают один объект - они оба получают доступ через memory-mapped файл, без копирования.
  • IPC через Arrow: объекты хранятся в Arrow-совместимом формате - NumPy массивы, pandas DataFrame, Arrow RecordBatch.
  • Distributed: объекты могут находиться на любой ноде. Ray автоматически передаёт их по сети, если нужно.
# Сохранение большого массива в Object Store
import numpy as np

big_array = np.random.randn(10_000_000)  # ~80 MB

# put(): сериализует и помещает в Plasma (на текущей ноде)
ref = ray.put(big_array)

# Два Actor-а на той же ноде читают без копирования
@ray.remote
def process_array(array_ref):
    arr = ray.get(array_ref)  # zero-copy если на той же ноде
    return arr.mean()

result1 = process_array.remote(ref)
result2 = process_array.remote(ref)

Инициализация и настройка RayDP

Установка и конфигурация

# Установка зависимостей
pip install raydp[torch]
# Включает: ray[default], pyspark, torch, pyarrow

# Для XGBoost-интеграции:
pip install raydp[xgboost]

Инициализация Ray и Spark

import ray
import raydp
from pyspark.sql import SparkSession

# Шаг 1: Инициализируем Ray кластер
# На локальной машине: ray.init() запускает локальный кластер
# На K8s: ray.init(address="ray://ray-head-svc:10001") подключается к существующему
ray.init(
    num_cpus=32,          # суммарные CPU ресурсы (все ноды)
    num_gpus=4,           # суммарные GPU
    object_store_memory=40 * 1024**3,  # 40 GB для Plasma Object Store
    dashboard_host="0.0.0.0",          # Ray Dashboard доступен снаружи
)

# Шаг 2: Инициализируем Spark ВНУТРИ Ray
# Это создаёт SparkDriverActor и N SparkExecutorActor-ов
spark = raydp.init_spark(
    app_name="recommendation-pipeline",
    num_executors=8,              # 8 Spark Executor-ов (Ray Actors)
    executor_cores=4,             # 4 CPU на каждый Executor
    executor_memory="16g",        # 16 GB RAM на каждый Executor
    configs={
        # Iceberg для чтения lakehouse
        "spark.sql.extensions": (
            "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions"
        ),
        "spark.sql.catalog.minio": "org.apache.iceberg.spark.SparkCatalog",
        "spark.sql.catalog.minio.type": "hadoop",
        "spark.sql.catalog.minio.warehouse": "s3a://lakehouse/warehouse/",
        # S3A для MinIO
        "spark.hadoop.fs.s3a.endpoint": "http://minio:9000",
        "spark.hadoop.fs.s3a.access.key": "minioadmin",
        "spark.hadoop.fs.s3a.secret.key": "minioadmin",
        "spark.hadoop.fs.s3a.path.style.access": "true",
        # AQE
        "spark.sql.adaptive.enabled": "true",
    }
)

print(f"Spark version: {spark.version}")
print(f"Ray resources: {ray.cluster_resources()}")

От DataFrame к Ray Actors: практический пайплайн

Полная система рекомендаций: от MinIO до PyTorch

Разберём end-to-end пример: корпоративный пайплайн рекомендательной системы. Входные данные - логи транзакций в MinIO (Parquet). Цель - обучить модель, которая предсказывает вероятность покупки пользователем товара из определённой категории.

Шаг 1: Spark ETL - Feature Engineering

from pyspark.sql import functions as F
from pyspark.sql.window import Window
from datetime import date, timedelta

# Читаем исходные таблицы из Iceberg
df_transactions = (
    spark.read.format("iceberg")
    .load("minio.silver.transactions")
    .filter(
        F.col("event_date") >= (date.today() - timedelta(days=90)).isoformat()
    )
)

df_users = spark.read.format("iceberg").load("minio.silver.users")
df_items = spark.read.format("iceberg").load("minio.silver.items")

# --- Генерация признаков на Spark (именно здесь Spark силён) ---

# Оконные функции: поведение пользователя за последние 30/90 дней
w_user_30d = Window.partitionBy("user_id").orderBy("event_date") \
    .rangeBetween(-30, 0)
w_user_90d = Window.partitionBy("user_id").orderBy("event_date") \
    .rangeBetween(-90, 0)

df_user_features = (
    df_transactions
    # Агрегаты поведения
    .withColumn("spend_30d", F.sum("amount").over(w_user_30d))
    .withColumn("spend_90d", F.sum("amount").over(w_user_90d))
    .withColumn("order_count_30d", F.count("order_id").over(w_user_30d))
    .withColumn("avg_basket_90d", F.avg("amount").over(w_user_90d))
    # Последняя транзакция
    .withColumn("days_since_last_order",
        F.datediff(F.current_date(), F.max("event_date").over(w_user_90d)))
    # Категориальные предпочтения (top category за 90 дней)
    .join(df_items.select("item_id", "category_l1", "price_tier"), "item_id")
    .groupBy("user_id", "event_date")
    .agg(
        F.first("spend_30d").alias("spend_30d"),
        F.first("spend_90d").alias("spend_90d"),
        F.first("order_count_30d").alias("order_count_30d"),
        F.first("avg_basket_90d").alias("avg_basket_90d"),
        F.first("days_since_last_order").alias("days_since_last_order"),
        F.mode("category_l1").alias("top_category"),  # Spark 3.4+
        F.mode("price_tier").alias("preferred_tier"),
        F.sum(F.when(F.col("is_return"), 1).otherwise(0))
         .alias("return_count_90d"),
    )
)

# Обогащение профилем пользователя
features_df = (
    df_user_features
    .join(
        df_users.select("user_id", "age_bucket", "region", "days_since_signup"),
        "user_id"
    )
    # Нормализация числовых признаков через встроенные функции Spark
    .withColumn("log_spend_90d", F.log1p(F.col("spend_90d")))
    .withColumn("log_order_count", F.log1p(F.col("order_count_30d")))
    # Целевая переменная: купит ли пользователь завтра (label = 1/0)
    .withColumn("label",
        F.when(F.col("event_date") == date.today().isoformat(), 1).otherwise(0))
    # Убираем строки с null в ключевых признаках
    .dropna(subset=["spend_30d", "spend_90d", "order_count_30d"])
)

print(f"Feature dataset: {features_df.count():,} rows")
features_df.printSchema()

Шаг 2: Конвертация Spark DataFrame → Ray Dataset

Это ключевой момент интеграции. raydp.spark.to_ray_dataset() конвертирует Spark DataFrame в Ray Dataset через Arrow, используя Plasma Object Store:

import raydp.spark

# Конвертация: Spark DataFrame → Ray Dataset
# Внутри: Spark сериализует партиции в Arrow RecordBatch,
# помещает их в Ray Object Store (Plasma), Ray Dataset читает по ObjectRef
ray_dataset = raydp.spark.to_ray_dataset(
    df=features_df,
    parallelism=32,           # на сколько партиций разбить Ray Dataset
)

print(f"Ray Dataset schema: {ray_dataset.schema()}")
print(f"Ray Dataset size: {ray_dataset.count():,} rows")
# Ray Dataset schema: user_id: string, spend_30d: double, label: int64, ...
# Ray Dataset size: 24,500,000 rows

# Разбивка на train/validation (без записи на диск!)
train_ds, val_ds = ray_dataset.train_test_split(test_size=0.2, shuffle=True, seed=42)

print(f"Train: {train_ds.count():,}, Validation: {val_ds.count():,}")

Шаг 3: Распределённое обучение с Ray Train

Ray Train предоставляет высокоуровневый API для distributed ML, скрывая сложность gradient synchronization и GPU orchestration.

import ray.train
from ray.train.torch import TorchTrainer
from ray.train import ScalingConfig, CheckpointConfig
import torch
import torch.nn as nn
import numpy as np
from typing import Dict

# Определяем архитектуру модели (простой двухслойный MLP для демонстрации)
class RecommendationModel(nn.Module):
    def __init__(self, input_dim: int):
        super().__init__()
        self.network = nn.Sequential(
            nn.Linear(input_dim, 256),
            nn.BatchNorm1d(256),
            nn.ReLU(),
            nn.Dropout(0.3),
            nn.Linear(256, 64),
            nn.ReLU(),
            nn.Linear(64, 1),
            nn.Sigmoid()
        )

    def forward(self, x: torch.Tensor) -> torch.Tensor:
        return self.network(x).squeeze(1)


# ЧИСЛОВЫЕ ПРИЗНАКИ: список колонок для обучения
FEATURE_COLS = [
    "log_spend_90d", "log_order_count", "avg_basket_90d",
    "days_since_last_order", "days_since_signup", "return_count_90d",
    "spend_30d",
]
LABEL_COL = "label"


def training_loop_per_worker(config: Dict):
    """
    Эта функция выполняется на КАЖДОМ Ray Train Worker (каждом GPU).
    Ray Train запускает её параллельно на всех воркерах и обеспечивает
    синхронизацию градиентов через PyTorch DistributedDataParallel (DDP).
    """
    import ray.train.torch

    # Получаем датасет для текущего воркера (автоматически шардируется)
    train_shard = ray.train.get_dataset_shard("train")
    val_shard = ray.train.get_dataset_shard("val")

    # Создаём модель и оборачиваем в DDP (Ray делает это автоматически)
    model = RecommendationModel(input_dim=len(FEATURE_COLS))
    model = ray.train.torch.prepare_model(model)  # оборачивает в DDP

    optimizer = torch.optim.AdamW(model.parameters(), lr=config["lr"],
                                  weight_decay=1e-4)
    scheduler = torch.optim.lr_scheduler.CosineAnnealingLR(
        optimizer, T_max=config["num_epochs"]
    )
    criterion = nn.BCELoss()

    for epoch in range(config["num_epochs"]):
        # --- Тренировочная фаза ---
        model.train()
        train_loss_sum = 0.0
        num_train_batches = 0

        for batch in train_shard.iter_torch_batches(
            batch_size=config["batch_size"],
            dtypes=torch.float32,
            prefetch_batches=2,   # prefetch следующих 2 батчей пока GPU считает
        ):
            # Извлекаем признаки и метки из батча
            features = torch.stack(
                [batch[col] for col in FEATURE_COLS], dim=1
            )
            labels = batch[LABEL_COL].float()

            # Forward pass
            optimizer.zero_grad()
            predictions = model(features)
            loss = criterion(predictions, labels)

            # Backward pass + gradient synchronization (DDP делает это автоматически)
            loss.backward()
            torch.nn.utils.clip_grad_norm_(model.parameters(), max_norm=1.0)
            optimizer.step()

            train_loss_sum += loss.item()
            num_train_batches += 1

        scheduler.step()
        avg_train_loss = train_loss_sum / num_train_batches

        # --- Валидация (только на rank 0 воркере) ---
        model.eval()
        val_loss_sum = 0.0
        num_val_batches = 0
        correct = 0
        total = 0

        with torch.no_grad():
            for batch in val_shard.iter_torch_batches(
                batch_size=config["batch_size"] * 2,
                dtypes=torch.float32,
            ):
                features = torch.stack(
                    [batch[col] for col in FEATURE_COLS], dim=1
                )
                labels = batch[LABEL_COL].float()
                preds = model(features)
                val_loss_sum += criterion(preds, labels).item()
                correct += ((preds > 0.5).float() == labels).sum().item()
                total += labels.shape[0]
                num_val_batches += 1

        avg_val_loss = val_loss_sum / num_val_batches
        accuracy = correct / total

        # Отчёт метрик в Ray Train (агрегируется со всех воркеров)
        ray.train.report(
            metrics={
                "epoch": epoch,
                "train_loss": avg_train_loss,
                "val_loss": avg_val_loss,
                "val_accuracy": accuracy,
                "lr": scheduler.get_last_lr()[0],
            },
            checkpoint=ray.train.Checkpoint.from_dict(
                {"model_state": model.state_dict(),
                 "optimizer_state": optimizer.state_dict(),
                 "epoch": epoch}
            ) if epoch % 5 == 0 else None,  # чекпоинт каждые 5 эпох
        )

        print(f"Epoch {epoch}: train_loss={avg_train_loss:.4f}, "
              f"val_loss={avg_val_loss:.4f}, accuracy={accuracy:.4f}")


# Конфигурация масштабирования: 4 GPU-воркера
scaling_config = ScalingConfig(
    num_workers=4,                  # 4 воркера = 4 GPU
    use_gpu=True,                   # каждый воркер получает 1 GPU
    resources_per_worker={"CPU": 4, "GPU": 1},
)

# Запуск обучения
trainer = TorchTrainer(
    train_loop_per_worker=training_loop_per_worker,
    train_loop_config={
        "lr": 1e-3,
        "num_epochs": 30,
        "batch_size": 4096,
    },
    scaling_config=scaling_config,
    datasets={
        "train": train_ds,   # Ray Dataset передаётся напрямую - без файлов!
        "val": val_ds,
    },
    run_config=ray.train.RunConfig(
        name="recommendation-v1",
        storage_path="s3://my-bucket/ray-results/",
        checkpoint_config=CheckpointConfig(
            num_to_keep=3,              # хранить 3 последних чекпоинта
            checkpoint_score_attribute="val_accuracy",
            checkpoint_score_order="max",
        ),
    ),
)

result = trainer.fit()
print(f"Best checkpoint: {result.best_checkpoints}")
print(f"Final metrics: {result.metrics}")

Шаг 4: Управление ресурсами между ETL и ML

Один из ключевых сценариев RayDP - динамическое переключение ресурсов между фазами. Во время Spark ETL GPU простаивают. Во время ML-тренировки Spark Executor-ы простаивают.

import raydp

# Фаза 1: Spark ETL (CPU-heavy)
# Spark Executors занимают все CPU, GPU свободны
spark = raydp.init_spark(
    app_name="etl-phase",
    num_executors=16,
    executor_cores=4,
    executor_memory="16g",
)

features_df = run_spark_etl(spark)   # тяжёлый ETL
ray_dataset = raydp.spark.to_ray_dataset(features_df)

# Освобождаем Spark-ресурсы после ETL
raydp.stop_spark()
# Теперь 64 CPU свободны - Ray переделит их под ML

# Фаза 2: ML Training (GPU-heavy)
# Все CPU и GPU доступны для Ray Train
trainer = TorchTrainer(
    train_loop_per_worker=training_loop_per_worker,
    scaling_config=ScalingConfig(num_workers=4, use_gpu=True),
    datasets={"train": ray_dataset},
)
result = trainer.fit()

# Опционально: снова запустить Spark для записи результатов
spark = raydp.init_spark(
    app_name="write-phase",
    num_executors=4,
    executor_cores=4,
    executor_memory="8g",
)
# Записать предсказания обратно в Iceberg...
raydp.stop_spark()

ray.shutdown()

Сравнительный анализ: Spark MLlib vs Spark + RayDP

Детальное сравнение

Критерий Spark MLlib Spark + RayDP + Ray Train
Среда выполнения JVM (Scala/Java) Python-native C++ (Ray core)
Deep Learning Практически отсутствует Полная нативная интеграция (PyTorch, Lightning, TF)
GPU orchestration Сложная настройка, неэффективна Fine-grained Actor-based scheduling
Передача данных Нет (всё внутри Spark) Zero-copy Arrow через Plasma
Итеративное обучение Overhead на DAG per iteration Ray Actors: persistent state, нет overhead
Экосистема ~30 встроенных алгоритмов Весь Python ML stack (Hugging Face, vLLM...)
Hyperparameter tuning Нет (сторонние интеграции) Ray Tune (встроен, 20+ алгоритмов поиска)
Model serving Нет Ray Serve (встроен)
Fault tolerance Spark lineage + retry Ray Actor recovery + checkpointing
Сложность настройки Низкая (всё в Spark) Средняя (два кластера или RayDP)
Когда уместен Простая табличная ML + ETL Сложный DL + LLM + GPU-heavy ML

Когда MLlib по-прежнему правильный выбор

MLlib не мёртв. Для ряда сценариев он остаётся оптимальным:

Простые табличные алгоритмы рядом с данными: если вам нужен Random Forest или Gradient Boosted Trees на данных, которые уже в Spark DataFrame, MLlib - самое простое решение. Нет необходимости в отдельном Ray-кластере.

ETL + классификация как одна задача: когда ML-логика - это не отдельная система, а часть трансформационного пайплайна. Например: df.filter(...).join(...).mllib_predict(model) - одна Spark job.

XGBoost через XGBoost4J-Spark: XGBoost имеет нативную Spark-интеграцию, работающую без RayDP. Для табличных задач это часто достаточно.

# Пример: когда MLlib достаточно
from pyspark.ml.classification import GBTClassifier
from pyspark.ml.feature import VectorAssembler
from pyspark.ml import Pipeline

assembler = VectorAssembler(
    inputCols=FEATURE_COLS,
    outputCol="features"
)

gbt = GBTClassifier(
    featuresCol="features",
    labelCol="label",
    maxIter=50,
    maxDepth=5,
)

pipeline = Pipeline(stages=[assembler, gbt])
model = pipeline.fit(features_df)  # всё внутри Spark - просто и надёжно

# Оценка
predictions = model.transform(features_df_test)

Когда эта простота правильная:

  • Команда хорошо знает Spark, не знает Ray
  • Модель - это tabular ML без DL (нет нейросетей, нет GPU)
  • Данных < 100 GB: MLlib хорошо масштабируется на такие объёмы
  • Нет требований к онлайн-инференсу, experimentation tracking и т.д.

Когда Ray - единственный правильный выбор

  • Нейронные сети: PyTorch, TensorFlow, JAX - всё это работает на Ray Train из коробки.
  • LLM fine-tuning: Hugging Face Trainer + Ray Train = стандарт индустрии для fine-tuning больших моделей.
  • Reinforcement Learning: Ray RLlib - лучший RL-фреймворк в экосистеме.
  • Hyperparameter optimization: Ray Tune с 20+ алгоритмами поиска (Bayesian Opt, PBT, ASHA).
  • Online model serving: Ray Serve для A/B тестирования и feature-rich inference pipelines.

Лабораторная практика: мониторинг и отладка

Ray Dashboard

Ray Dashboard (http://localhost:8265) - центральный инструмент мониторинга. Ключевые разделы:

Actors: показывает все живые Actor-ы (SparkDriverActor, SparkExecutorActor-ы, ModelTrainer-ы), их состояние, потребление ресурсов и логи.

Tasks: текущие и завершённые задачи с временем выполнения. Позволяет найти медленные батчи или задачи с OOM.

Metrics: CPU/GPU/Memory утилизация по нодам. Для ML-фазы ожидаем GPU > 80%, для ETL-фазы - GPU = 0%, CPU = 90%+.

Object Store: состояние Plasma - сколько объектов, суммарный размер, spill на диск.

# Программный доступ к статистике Ray
import ray

# Статистика кластера
cluster_resources = ray.cluster_resources()
available_resources = ray.available_resources()
print(f"Total GPUs: {cluster_resources.get('GPU', 0)}")
print(f"Available GPUs: {available_resources.get('GPU', 0)}")

# Статистика Object Store
nodes = ray.nodes()
for node in nodes:
    obj_store = node.get("ObjectStoreSocketName", "")
    print(f"Node {node['NodeID'][:8]}: "
          f"CPU={node['Resources'].get('CPU', 0)}, "
          f"GPU={node['Resources'].get('GPU', 0)}")

Диагностика проблем с памятью Object Store

Когда Ray Dataset не вмещается в Plasma Object Store, Ray начинает spill на диск - что многократно замедляет работу:

# В логах Ray при spill:
INFO raylet: Spilling objects of total size 12GB to external storage

# Диагностика: сколько данных в Object Store?
import ray._private.internal_api as ray_internal

# Смотрим на Object Store usage
for node in ray.nodes():
    state = node["ObjectStoreSocketName"]
    print(f"Object Store: {state}")

Решение: увеличить object_store_memory при ray.init() или уменьшить parallelism при to_ray_dataset() (меньше партиций = меньше данных одновременно в памяти).

Batch Inference: Spark ETL → Ray для предсказаний

После обучения модели часто нужно делать batch inference на новых данных: читаем данные Spark-ом, делаем inference через Ray Actors.

import ray
from ray.data import ActorPoolStrategy

@ray.remote
class InferenceActor:
    """Долгоживущий Actor с загруженной моделью - state сохраняется."""

    def __init__(self, checkpoint_path: str):
        import torch
        self.model = RecommendationModel(input_dim=len(FEATURE_COLS))

        # Загружаем чекпоинт из Ray Train
        checkpoint = ray.train.Checkpoint(checkpoint_path)
        with checkpoint.as_directory() as d:
            state = torch.load(f"{d}/model.pt", map_location="cuda")
            self.model.load_state_dict(state["model_state"])

        self.model.eval().cuda()
        print(f"Model loaded on GPU from {checkpoint_path}")

    def predict_batch(self, batch: dict) -> dict:
        import torch
        features = torch.stack(
            [torch.tensor(batch[col]) for col in FEATURE_COLS], dim=1
        ).float().cuda()

        with torch.no_grad():
            probs = self.model(features).cpu().numpy()

        return {
            "user_id": batch["user_id"],
            "buy_probability": probs.tolist(),
        }


# Пайплайн inference
def run_batch_inference(spark: SparkSession, checkpoint_path: str):
    # 1. Spark ETL: готовим данные для предсказания
    new_data_df = (
        spark.read.format("iceberg").load("minio.silver.users_today")
        .join(spark.read.format("iceberg").load("minio.silver.features"), "user_id")
        .select("user_id", *FEATURE_COLS)
        .dropna()
    )

    # 2. Конвертация в Ray Dataset (zero-copy)
    inference_ds = raydp.spark.to_ray_dataset(new_data_df, parallelism=16)

    # 3. Batch inference через Actor Pool (2 GPU Actor-а параллельно)
    results_ds = inference_ds.map_batches(
        InferenceActor,
        fn_constructor_kwargs={"checkpoint_path": checkpoint_path},
        concurrency=2,                          # 2 Actor-а параллельно
        num_gpus=1,                             # каждый Actor: 1 GPU
        batch_size=8192,                        # большой батч = эффективнее GPU
        batch_format="numpy",
    )

    # 4. Записываем результаты обратно (через Spark или напрямую через Arrow)
    results_arrow = results_ds.to_arrow()
    results_spark_df = spark.createDataFrame(results_arrow.to_pandas())
    results_spark_df.write.format("iceberg").mode("overwrite") \
        .saveAsTable("minio.gold.user_buy_probabilities")

    print(f"Inference complete: {results_ds.count():,} predictions written")

Best Practices и матрица применимости

Когда связка Spark + Ray необходима

Тяжёлый ETL + сложный ML/DL: датасет больше RAM одной машины, модель - нейросеть или ансамбль с GPU, нужен полный Python ML stack. Это главный кейс для RayDP.

LLM fine-tuning на проприетарных данных: схема: Spark читает и чистит текстовые данные из Iceberg → RayDP передаёт в Ray → Hugging Face Trainer на GPU.

Reinforcement Learning с реальными данными: среда для RL создаётся из данных, подготовленных Spark-ом; Ray RLlib обучает агента.

Feature Store + Online Serving: Spark строит features в offline-режиме → Ray Serve раздаёт предсказания в реальном времени.

Когда это избыточно

Простая табличная ML с XGBoost/LightGBM: используйте XGBoost4J-Spark или LightGBM-Spark напрямую. Нет нейросетей = нет нужды в Ray.

ML на данных < 10 GB: sklearn + pandas на одной машине справятся. Ray добавит сложность без пользы.

Streaming ML: Spark Structured Streaming + онлайн-модель через predict() вызов - проще и надёжнее, чем Ray для этого сценария.

Анти-паттерны

Анти-паттерн 1: Использовать Ray для SQL. Ray Data умеет фильтровать и проецировать, но у него нет Catalyst optimizer, нет AQE, нет нативного Parquet predicate pushdown. Для SQL - Spark.

Анти-паттерн 2: Запускать Spark и Ray как два отдельных кластера без RayDP. Тогда передача данных идёт через S3 - медленно и дорого. RayDP именно для того, чтобы избежать этого.

Анти-паттерн 3: Слишком маленькие батчи в Ray Dataset. Маленькие батчи недозагружают GPU. Для обучения нейросетей: батч > 1024, лучше 4096–16384 элементов.


Домашнее задание

Условие:

Вам дан legacy-скрипт машинного обучения на Spark MLlib:

# УСТАРЕВШИЙ СКРИПТ НА SPARK MLLIB (стартовая точка)
from pyspark.sql import SparkSession, functions as F
from pyspark.ml.classification import RandomForestClassifier
from pyspark.ml.feature import VectorAssembler, StringIndexer
from pyspark.ml.evaluation import BinaryClassificationEvaluator
from pyspark.ml import Pipeline

spark = SparkSession.builder.appName("churn-prediction-legacy").getOrCreate()

# Чтение данных (CSV - неоптимально, но дано как есть)
df = spark.read.csv("/data/customer_features/*.csv", header=True, inferSchema=True)

# Ограниченная feature engineering (нет оконных функций)
df_features = (
    df
    .filter(F.col("is_active") == True)
    .select("customer_id", "age", "tenure_months", "avg_monthly_spend",
            "support_tickets", "product_category", "churn_label")
    .dropna()
)

# StringIndexer для категориального признака
indexer = StringIndexer(inputCol="product_category", outputCol="category_idx")
assembler = VectorAssembler(
    inputCols=["age", "tenure_months", "avg_monthly_spend",
               "support_tickets", "category_idx"],
    outputCol="features"
)

# Random Forest - ограниченная точность, нет GPU, нет DL
rf = RandomForestClassifier(featuresCol="features", labelCol="churn_label",
                             numTrees=100, maxDepth=10)
pipeline = Pipeline(stages=[indexer, assembler, rf])
model = pipeline.fit(df_features)

# Оценка
evaluator = BinaryClassificationEvaluator(labelCol="churn_label")
auc = evaluator.evaluate(model.transform(df_features))
print(f"AUC: {auc:.4f}")

Скрипт работает, но имеет ограничения: нет оконных признаков, нет GPU, Random Forest упирается в точность, нет hyperparameter tuning.

Задание (три части):

Часть 1: Улучшенный Spark ETL

Перепишите feature engineering: добавьте оконные функции (spend_90d, ticket_rate_30d, days_since_last_login), нормализацию через F.log1p, обогащение из второй таблицы (product_details). Выходной датасет должен содержать минимум 10 признаков.

Часть 2: Перевод на RayDP + Ray Train

Реализуйте:

  • raydp.init_spark() вместо обычного SparkSession.builder
  • Конвертацию через raydp.spark.to_ray_dataset()
  • Замену Random Forest на распределённый XGBoostTrainer из ray.train.xgboost:
from ray.train.xgboost import XGBoostTrainer

trainer = XGBoostTrainer(
    scaling_config=ScalingConfig(num_workers=4, use_gpu=False),
    label_column="churn_label",
    num_boost_round=200,
    params={"objective": "binary:logistic", "eval_metric": "auc"},
    datasets={"train": train_ds, "validation": val_ds},
)
result = trainer.fit()

Часть 3: Мониторинг и отчёт

Запустите оба скрипта (MLlib и RayDP+XGBoost) на одном датасете (минимум 1M строк). Приложите:

  • Скриншот Ray Dashboard с метриками воркеров во время тренировки
  • Сравнительную таблицу: время обучения, AUC на validation, потребление RAM
  • Обоснование: в каких сценариях MLlib остаётся правильным выбором для данного типа задачи?