Блок 1 - Архитектура и основы
20 вопросов о внутреннем устройстве Spark: Driver, Executor, RDD, DAG, Catalyst, Tungsten и fault tolerance.
Что такое Apache Spark и чем он отличается от MapReduce?¶
Apache Spark - это unified analytics engine для крупномасштабной обработки данных. Ключевое отличие от Hadoop MapReduce: Spark хранит промежуточные результаты в памяти (а не записывает каждый шаг на диск) и строит произвольный DAG операций, тогда как MapReduce ограничен двумя фазами - Map и Reduce. Практический результат: итеративные алгоритмы (ML) на Spark в 10–100 раз быстрее. Spark не хранит данные - он читает из внешних источников (HDFS, S3, Kafka) и пишет результат обратно.
Опишите компоненты архитектуры: Driver, Executor, Cluster Manager¶
- Driver - процесс пользовательского кода. Строит DAG, планирует задачи, координирует выполнение. Работает на одной машине.
- Executor - JVM-процесс на рабочих узлах. Выполняет Tasks и хранит закешированные данные. Число и размер задаются при submit.
- Cluster Manager - управляет ресурсами кластера. Spark поддерживает YARN, Kubernetes, Mesos, Standalone. Не знает о внутренней логике Spark - просто выдаёт контейнеры.
Что такое RDD и каковы его свойства?¶
RDD (Resilient Distributed Dataset) - базовая абстракция Spark, неизменяемая коллекция записей, распределённая по партициям на узлах кластера. Ключевые свойства:
| Свойство | Суть |
|---|---|
| Immutable | После создания RDD не изменяется; трансформации создают новый RDD |
| Distributed | Данные партиционированы по узлам кластера |
| Fault-tolerant | Каждый RDD помнит линейку (lineage) - как он был создан; при сбое пересчитывается |
| Lazy | Вычисляется только при вызове Action |
В современном коде RDD используют редко - DataFrame/Dataset перекрывает 95% задач и оптимизируется через Catalyst.
В чём разница между DataFrame и Dataset? Почему в PySpark нет Dataset API?¶
DataFrame - распределённая таблица с именованными колонками. Тип строки - Row, схема известна только в runtime. Dataset - типизированная версия DataFrame: Spark знает тип каждой записи в compile-time (Dataset[Person]), что позволяет ловить ошибки типов до запуска.
В PySpark Dataset API нет, потому что Python - динамически типизированный язык. Compile-time проверка типов здесь невозможна - это преимущество Java/Scala. В Python весь кодё идёт через DataFrame, а типы проверяются через схему StructType в runtime.
Как работает Py4J в контексте PySpark?¶
Py4J - библиотека, создающая мост между Python-процессом и JVM. Когда вы вызываете spark.read.parquet(...) в Python, вызов транслируется через Py4J в Java-объекты SparkContext и выполняется в JVM. Данные остаются в JVM-памяти (Executor'ы). Python-процесс - только тонкий клиент, который формирует инструкции.
Исключение - Python UDF: тогда данные сериализуются из JVM в Python-процесс на воркере, обрабатываются и возвращаются обратно. Это главная причина медлительности Python UDF.
Что такое Lazy Evaluation?¶
Spark не выполняет трансформации (filter, select, join) в момент их вызова - только записывает план. Реальное выполнение запускается при вызове Action (count(), show(), write()). Это позволяет Catalyst оптимизировать весь план целиком: переставить фильтры, убрать ненужные колонки, выбрать стратегию join.
df = spark.read.parquet("s3://...") # Ничего не читается
df2 = df.filter("age > 18").select("id") # Только план
df2.count() # Вот тут Spark реально работает
Разница между Transformations и Actions¶
| Transformations | Actions | |
|---|---|---|
| Выполнение | Lazy - строят план | Eager - запускают выполнение |
| Возвращает | Новый DataFrame/RDD | Результат или запись в хранилище |
| Примеры | filter, select, join, groupBy |
count, show, collect, write |
Трансформации делятся на Narrow (каждый output-раздел зависит от одного input-раздела, shuffle не нужен) и Wide (данные перемещаются между узлами - shuffle).
Что такое DAG?¶
DAG (Directed Acyclic Graph) - ориентированный граф без циклов, в котором узлы - это операции (трансформации), а рёбра - зависимости между ними. Spark строит DAG из всех трансформаций до момента Action. DAG Scheduler делит его на Stages (по границам shuffle), Stages - на Tasks (по числу партиций). DAG виден на вкладке SQL в Spark UI.
Как Spark разбивает Job на Stages и Tasks?¶
- Job создаётся при вызове Action.
- DAG Scheduler находит в DAG границы Shuffle (Wide трансформации -
groupBy,join,repartition). Каждая граница = новый Stage. - Каждый Stage состоит из Tasks - по одной на каждую партицию входных данных.
- Task Scheduler распределяет Tasks по Executor'ам с учётом data locality.
Правило: N партиций на входе Stage → N Tasks в этом Stage.
Что такое Narrow и Wide зависимости?¶
Narrow: каждый output-раздел зависит от одного (или фиксированного числа) input-разделов. Shuffle не нужен. Примеры: map, filter, union, select.
Wide: output-раздел может зависеть от всех input-разделов. Требует Shuffle - перераспределения данных по сети. Примеры: groupBy, join, distinct, repartition. Wide-трансформации - главная причина медленной работы и вызывают Stage Boundary.
Что такое Shuffle и почему он дорогой?¶
Shuffle - физическое перемещение данных между Executor'ами для группировки строк с одинаковыми ключами на одном узле. Этапы: (1) Shuffle Write - каждый Executor записывает свои данные по ключам на локальный диск, (2) Shuffle Read - каждый Executor читает чужие данные по сети.
Почему дорого: сетевой I/O, дисковый I/O, сериализация/десериализация данных. При неравномерном распределении ключей (data skew) один Executor пишет/читает несоразмерно больше других.
Роль SparkContext и SparkSession¶
SparkContext (Spark 1.x) - точка входа низкого уровня, создаёт соединение с кластером, позволяет работать с RDD. SparkSession (Spark 2.x+) - единая точка входа, объединяет SparkContext, SQLContext и HiveContext. В современном коде используется только SparkSession:
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("MyApp").getOrCreate()
sc = spark.sparkContext # Доступен через SparkSession при необходимости
Важное ограничение: в одной JVM может существовать только один активный SparkContext. getOrCreate() возвращает существующий экземпляр, если он уже создан - это singleton-паттерн. Попытка создать второй SparkContext явно приведёт к ошибке или остановке первого. Это означает:
- Одно Spark-приложение = одна JVM = один SparkSession/SparkContext
- Для параллельной обработки разных датасетов нужно параллельно запускать отдельные spark-submit задания, а не несколько SparkSession внутри одного
- SQLContext и HiveContext (Spark 1.x) - устаревшие альтернативы; SparkSession поглощает оба
Что такое партиции? Как определить их оптимальное количество?¶
Партиция - логический срез данных, обрабатываемый одной Task. По умолчанию spark.sql.shuffle.partitions = 200. Практическое правило: целевой размер партиции после shuffle - 100–200 МБ. Формула: N = общий_объём_shuffle / 150MB. Слишком мало партиций → каждая Task тяжелее, риск OOM. Слишком много → overhead на планирование, много мелких файлов при записи.
Как работает Catalyst Optimizer?¶
Catalyst - движок оптимизации запросов. Работает в 4 фазы:
- Analysis - разрешение имён колонок и таблиц (Unresolved → Resolved LogicalPlan)
- Logical Optimization - rule-based: predicate pushdown, column pruning, constant folding
- Physical Planning - выбор стратегий join (broadcast vs sort-merge), порядок операций
- Code Generation (Tungsten) - генерация Java bytecode для выполнения без интерпретатора
Что такое Project Tungsten?¶
Tungsten - движок физического выполнения Spark, нацеленный на максимальное использование аппаратуры:
- Off-heap memory - управление памятью без GC, через
sun.misc.Unsafe - Cache-aware data structures - данные в памяти расположены для максимального cache hit
- Whole-stage code generation - весь Stage компилируется в единый Java-метод (нет virtual dispatch)
Результат: CPU-bound операции (join, aggregation) выполняются в 2–10× быстрее чем в Spark 1.x.
Как Spark обеспечивает fault tolerance?¶
Через Lineage (линейку): каждый RDD/DataFrame знает, как он был создан (из какого RDD и через какие трансформации). При сбое Executor'а Spark не перезапускает Job целиком - только пересчитывает потерянные партиции, пройдя по lineage с последнего checkpoint.
При использовании .cache() данные могут быть недоступны после сбоя - Spark пересчитает их заново. checkpoint() обрезает lineage и записывает данные на надёжное хранилище (HDFS/S3).
Client mode vs Cluster mode¶
| Client mode | Cluster mode | |
|---|---|---|
| Driver работает | На машине, откуда запускается spark-submit |
Внутри кластера (один из узлов) |
| Использование | Интерактивная разработка, Jupyter | Production-задачи |
| Риск | Если машина упадёт - Job падает | Отказоустойчив - Driver в кластере |
| Сетевой трафик | Данные Driver ↔ Executor через сеть клиента | Всё внутри кластера |
Какие менеджеры ресурсов поддерживает Spark?¶
- YARN - Hadoop Resource Manager, стандарт для On-premise кластеров
- Kubernetes - современный стандарт, контейнеры, auto-scaling, cloud-native
- Apache Mesos - устаревает, практически не используется в новых проектах
- Standalone - встроенный менеджер Spark без внешних зависимостей, для dev/test
Что такое Local mode?¶
SparkSession.builder.master("local[*]") запускает Spark в одном JVM-процессе без кластера. Driver и Executor - один и тот же процесс. local[N] - N потоков (Threads) вместо Executor'ов. Используется для юнит-тестов, отладки и небольших датасетов. В local режиме нет shuffle сети - всё в памяти одного процесса.
Как PySpark взаимодействует с Python-воркерами?¶
При встроенных функциях (functions.col, groupBy, join) весь расчёт происходит в JVM Executor'ах. Python-процесс отправляет только инструкции через Py4J. При Python UDF Spark запускает дочерний Python-процесс на каждом Executor'е, сериализует данные (через pickle или Arrow), передаёт их в Python, получает результат и десериализует обратно. Именно поэтому Pandas UDF (Arrow) намного быстрее Python UDF (pickle, строка-за-строкой).
Подробно опишите жизненный цикл приложения Spark от spark-submit до завершения¶
spark-submitпередаёт JAR/Python-файл и конфигурацию в Cluster Manager- Cluster Manager запускает Driver (в cluster-mode - на одном из узлов кластера, в client-mode - на клиентской машине)
- Driver инициализирует
SparkSession/SparkContext, подключается к Cluster Manager - Driver запрашивает ресурсы → Cluster Manager запускает Executor'ы на worker-узлах
- Executor'ы регистрируются у Driver через heartbeat
- При вызове Action: Driver строит DAG, разбивает на Stages и Tasks
- Task Scheduler отправляет Tasks на Executor'ы с учётом data locality
- Executor'ы выполняют Tasks, обмениваются shuffle-данными напрямую между собой
- Результат Actions возвращается на Driver (или записывается в хранилище)
- При завершении Driver освобождает ресурсы: Executor'ы останавливаются, Cluster Manager забирает контейнеры
Что такое сериализация данных в Spark и почему Kryo эффективнее Java-сериализации?¶
Сериализация - преобразование объектов Java в байты для передачи по сети (shuffle) или записи на диск (cache, spill). Spark поддерживает два механизма:
| Java Serialization | Kryo Serialization | |
|---|---|---|
| Компактность | Высокий overhead | ~5× меньше байт |
| Скорость | Медленная рефлексия | В разы быстрее |
| Совместимость | Все Serializable объекты | Требует регистрации |
conf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
# Опциональная регистрация классов для максимальной эффективности
conf.set("spark.kryo.classesToRegister", "com.example.MyClass")
Kryo критично важен при работе с RDD. Для DataFrame-операций Tungsten использует собственный бинарный формат (UnsafeRow) - Kryo там не задействован.
Как работает механизм Data Locality (локальность данных)?¶
Spark пытается запустить Task на том узле, где физически хранятся обрабатываемые данные, чтобы избежать сетевой передачи.
Уровни локальности (от лучшего к худшему):
| Уровень | Данные находятся | Задержка |
|---|---|---|
PROCESS_LOCAL |
В памяти того же JVM (кеш) | ~0 |
NODE_LOCAL |
На диске того же узла (HDFS block) | Низкая |
RACK_LOCAL |
На другом узле в том же rack | Средняя |
ANY |
На любом узле кластера | Сетевая |
Spark ждёт spark.locality.wait (3 сек) на каждом уровне, прежде чем опуститься ниже.
На S3 Data Locality не работает: объектное хранилище не привязано к узлам кластера, все данные читаются по сети. Поэтому для S3-workload'ов нет смысла задерживаться - можно установить spark.locality.wait=0.
Каковы основные ограничения (недостатки) Apache Spark?¶
Понимание слабых сторон Spark важно, чтобы выбирать инструмент осознанно.
| Ограничение | Подробнее |
|---|---|
| Нет собственного хранилища | Spark - только вычислительный движок. Требует HDFS, S3, Iceberg и т.д. |
| Высокое потребление памяти | In-memory модель требует больших heap'ов; OOM-ошибки сложнее диагностировать, чем disk-based подходы |
| Не масштабируется для CPU-intensive задач | Оптимизирован под I/O-heavy workload'ы; для ML-инференса на GPU лучше TensorFlow/PyTorch |
| Latency выше, чем у OLTP | Минимальная задержка - секунды (micro-batch); для sub-millisecond нужны Flink/Redis |
| Сложная настройка памяти | Executor memory, overhead, off-heap, G1GC - требует экспертизы для production-тюнинга |
| Small files problem | Без компактификации накапливаются миллионы мелких Parquet-файлов, замедляющих listing |
| MLlib уступает специализированным фреймворкам | Нет ряда алгоритмов (например, Tanimoto distance); для глубокого обучения - PyTorch/TF лучше |
Когда Spark - правильный выбор:
- Объём данных > нескольких GB (не помещается в pandas)
- Batch ETL: трансформации, агрегации, joins
- Стриминг с задержкой >1 сек (Structured Streaming)
- Единое API для batch + stream в одном пайплайне
Когда лучше использовать другое:
- Sub-second latency → Apache Flink
- Маленькие датасеты (<1 GB) → pandas/DuckDB
- GPU-inference → TensorFlow Serving / Triton
- OLTP запросы → PostgreSQL / ClickHouse