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.
- RDD — низкоуровневый resilient distributed dataset (legacy).
- DataFrame — таблица с именованными колонками, Catalyst optimizer (основной API в PySpark).
- Dataset — typed DataFrame (Scala/Java).
Для 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 (начальный уровень).
spark.executor.memory,spark.executor.coresspark.sql.shuffle.partitions(default 200 — часто нужно менять)- Dynamic allocation — scale executors по нагрузке
Как это выглядит на практике
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).
Типовые проблемы студентов:
AnalysisException: path not found— typo в s3a path или нет credentials.- Driver OOM после
collect()илиtoPandas()на большом DF. - Small files problem — миллионы tiny Parquet files; compact job.
Что сделать после занятия
- [ ] Объясните разницу между transformation и action на примере
filterиcount. - [ ] Почему
groupBy(customer_id).agg(...)может вызвать shuffle? - [ ] Откройте Spark UI (или скриншот из лекции) и найдите stage с наибольшим shuffle read.