MLOps Path

6.3.1 · блок 6

Распределённая обработка и Apache Spark

Распределённая обработка и Apache Spark

Зачем это нужно

Dataset из 500 строк обучается в pandas на ноутбуке. Dataset из 500 миллионов строк — уже нет. Apache Spark — движок распределённой обработки данных: тот же SQL-подобный или DataFrame API, но вычисления идут на кластере worker'ов. Для MLOps Spark — workhorse batch feature engineering, ETL в data lake и подготовки training snapshots.

Без понимания базовой модели Spark вы не сможете отладить «почему job 6 часов» или «почему OOM на executor».

Основные идеи

Driver vs Executors.

| Компонент | Роль |

|-----------|------|

| **Driver** | Планирует job, держит SparkContext, собирает результаты |

| **Executor** | JVM-процесс на worker-ноде, выполняет tasks, хранит cache |

| **Cluster Manager** | YARN, Kubernetes, standalone — выделяет ресурсы |

RDD vs DataFrame vs Dataset.

Для ML в PySpark почти всегда DataFrame API.

Lazy evaluation. Трансформации (filter, select, groupBy) строят DAG; action (count, write, collect) запускает выполнение. collect() на большом DF — антипаттерн (тянет всё на driver).

Partitioning. Данные разбиты на partitions; один task ≈ одна partition. Слишком мало partitions — нет parallelism; слишком много — overhead. Правило большого пальца: 2–4 tasks на CPU core executor.

Shuffle — дорогая операция. groupBy, join, repartition перемешивают данные между nodes. Shuffle spill на disk → медленно. Профилируйте Spark UI: Stages with large shuffle read/write.

Spark on Kubernetes. Spark Operator или spark-submit --master k8s://... поднимает driver + executors как Pods. Интеграция с object storage (MinIO/S3) через s3a:// paths.

Форматы хранения.

| Формат | Плюсы для ML |

|--------|--------------|

| **Parquet** | Columnar, compression, schema |

| **Delta Lake / Iceberg** | ACID, time travel, merge |

| **CSV** | Только для tiny datasets / debug |

Structured Streaming (preview). Тот же Spark API для streaming источников (Kafka, files). Micro-batch каждые N секунд. Связь с модулем 6.2.

Resource tuning (начальный уровень).

Как это выглядит на практике

Batch feature job (PySpark, упрощённо):


from pyspark.sql import SparkSession
from pyspark.sql import functions as F

spark = SparkSession.builder.appName("churn-features-daily").getOrCreate()

events = spark.read.parquet("s3a://ml-bronze/customer_events/")
customers = spark.read.parquet("s3a://ml-silver/customers/")

features = (
    events.filter(F.col("event_date") == "2026-01-15")
    .groupBy("customer_id")
    .agg(
        F.count("*").alias("event_count_1d"),
        F.sum("amount").alias("amount_sum_1d"),
    )
    .join(customers, "customer_id", "left")
)

features.write.mode("overwrite").parquet(
    "s3a://ml-gold/churn/features/dt=2026-01-15/"
)

Запуск в K8s (концепт): Argo Workflow step с образом spark:3.5-py и ServiceAccount с доступом к MinIO.

Spark UI: port-forward на driver Pod → Stages, skew (одна task 10x дольше других → data skew, нужен repartition/salting).

Типовые проблемы студентов:

Что сделать после занятия

Официальные материалы

Открыть интерактивную версию