Ray + RayDP: Spark ETL → Ray Actor для distributed ML, сравнение с MLlib
Ray + RayDP: Spark ETL → Ray Actor для distributed ML, сравнение с MLlib
Архитектурный кризис: почему 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 делает следующее:
- Создаёт
SparkDriverActor- Ray Actor, внутри которого живёт JVM с Spark Driver (DAGScheduler, SparkContext). - Создаёт
SparkExecutorActor× N - Ray Actors, каждый из которых содержит JVM с Spark Executor. - Регистрирует эти акторы как вычислительные ресурсы через Ray Resource Manager.
- Настраивает 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 остаётся правильным выбором для данного типа задачи?