Зачем это нужно
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.