Учебник Apache Spark (PySpark) для начинающих
Apache Spark — это унифицированный движок для распределённой обработки больших данных: то, что не влезает в память одной машины и считалось бы вечность на одном процессоре, Spark раскладывает на сотни ядер кластера и считает за минуты. Это де-факто стандарт для дата-инженерии — от ночных ETL-пайплайнов на терабайтах до потоковой аналитики и обучения ML-моделей.
Этот курс — не справочник методов, а глубокий разбор того, как Spark устроен внутри и почему он ведёт себя так, а не иначе. Мы пройдём путь от фундаментальной проблемы «данные не влезают» через архитектуру (driver, executors, partitions, stages), модель RDD и ленивые вычисления, DataFrame API с оптимизатором Catalyst, shuffle и партиционирование как главную тему производительности, joins и оконные функции — до Structured Streaming, MLlib и продакшн-развёртывания. Особое внимание — распределённым ловушкам: shuffle, перекосу данных, OOM и мелким файлам.
Курс рассчитан на тех, кто знает Python и SQL и хочет осознанно работать с большими данными, а не копировать рецепты. После него вы будете понимать, почему один join падает с OutOfMemory, а другой летает; когда cache() помогает, а когда вредит; и как прочитать план выполнения, чтобы найти лишний shuffle.
Курс «Apache Spark (PySpark)» состоит из 6 разделов и 18 уроков: Большие данные и архитектура Spark, RDD и модель вычислений, DataFrame и Dataset API, Spark SQL, агрегации и соединения, Производительность, партиционирование и shuffle и Стриминг, ML и продакшн. Уроки идут по порядку — от основ к более сложным темам, в каждом есть объяснение с примерами, а в конце — вопросы для самопроверки. К урокам привязаны задачи с автоматической проверкой: прочитали тему — сразу закрепили её кодом.
Программа курса
1 Большие данные и архитектура Spark
- Зачем нужен Spark: когда данные перестают влезать
Проблема больших данных: почему один компьютер и pandas не справляются, что такое распределённые вычисления и какую боль решает Apache Spark.
- Архитектура Spark: driver, executors и кластер
Как устроено приложение Spark: driver и executors, cluster manager (YARN/K8s/Standalone), партиции, и как job делится на stages и tasks.
- SparkSession и экосистема Spark
Точка входа SparkSession: как запускается приложение PySpark, ленивость на практике, и обзор экосистемы — Spark SQL, Structured Streaming, MLlib, GraphX.
- Зачем нужен Spark: когда данные перестают влезать
2 RDD и модель вычислений
- RDD: неизменяемость, lineage и отказоустойчивость
Что такое RDD: неизменяемая распределённая коллекция, партиции, граф происхождения (lineage/DAG) и как Spark восстанавливает потерянные данные без репликации.
- Трансформации и действия: ленивость как основа Spark
Главное разделение в Spark: ленивые трансформации (map, filter, flatMap, reduceByKey) против действий (collect, count, reduce), и почему ленивость — фундамент движка.
- Shuffle, узкие/широкие зависимости и кэширование
Узкие против широких трансформаций, почему shuffle — самая дорогая операция в Spark, как его уменьшить, и кэширование: cache/persist и уровни хранения.
- RDD: неизменяемость, lineage и отказоустойчивость
3 DataFrame и Dataset API
- DataFrame против RDD: Catalyst и Tungsten
Почему DataFrame быстрее и удобнее RDD: оптимизатор запросов Catalyst, движок Tungsten с колоночным представлением и whole-stage codegen.
- Создание DataFrame, схемы и базовые операции
Как создавать DataFrame, чем явная схема лучше inferSchema, и базовые операции: select, filter/where, withColumn, drop, distinct.
- Функции, типы, null и чтение/запись (Parquet)
Функции pyspark.sql.functions (col, lit, when, агрегаты), работа с null, и форматы хранения: почему Parquet предпочтительнее CSV/JSON, partitionBy при записи.
- DataFrame против RDD: Catalyst и Tungsten
4 Spark SQL, агрегации и соединения
- Spark SQL, группировки и оконные функции
Spark SQL и createTempView, groupBy/agg для агрегаций и оконные функции (Window) для расчётов внутри групп без схлопывания строк.
- JOIN в Spark: broadcast против sort-merge
Соединения в Spark глубоко: типы join, две главные стратегии — broadcast hash join и sort-merge join, когда какая, broadcast hint и перекос при join.
- UDF и pandas UDF: почему обычные UDF медленные
Пользовательские функции: почему обычные Python-UDF медленные (сериализация, чёрный ящик для Catalyst) и как pandas UDF (Arrow) ускоряют дело.
- Spark SQL, группировки и оконные функции
5 Производительность, партиционирование и shuffle
- Партиционирование: сколько партиций, repartition и coalesce
Что такое партиции и сколько их нужно, spark.sql.shuffle.partitions, разница repartition (полный shuffle) и coalesce (без полного перемешивания).
- Shuffle и перекос данных: диагностика и приёмы
Почему shuffle дорог и что его вызывает, как его уменьшить, broadcast-переменные и борьба с перекосом данных (data skew): диагностика, salting, AQE.
- План выполнения, Catalyst/Tungsten и типичные ошибки
Как Spark оптимизирует план (Catalyst, Tungsten) и как его прочитать через explain (физический план), плюс типичные ошибки: OOM, мелкие файлы, лишние shuffle.
- Партиционирование: сколько партиций, repartition и coalesce
6 Стриминг, ML и продакшн
- Structured Streaming: бесконечная таблица и watermarks
Structured Streaming: модель «бесконечной таблицы», micro-batch обработка, окна по времени события и watermarks для работы с опоздавшими данными.
- MLlib, развёртывание и конфигурация
MLlib и концепция Pipeline, развёртывание через spark-submit (cluster/client mode), динамическое распределение ресурсов и ключевая конфигурация, включая AQE.
- Spark в дата-пайплайне: облако, форматы и что дальше
Место Spark в дата-пайплайне: облачные платформы (Databricks, EMR, Glue), оркестрация Airflow, источники Kafka, табличные форматы Delta/Iceberg и куда двигаться дальше.
- Structured Streaming: бесконечная таблица и watermarks