06 · Данные для ML6.1–6.4 · Потоки, батч и признаки6.3.1сложный

Распределённая обработка и 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.cores

  • spark.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.

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