Перейти к содержанию
Learning Platform

Главная ментальная модель

Когда вы пишете df.filter(col("age") > 30).select("name"), Spark не выполняет этот код буквально. Он не вызывает ваш filter, а потом ваш select. Вместо этого Spark рассматривает написанное как декларацию намерения — описание того, какой результат вы хотите получить, — и оставляет за собой полную свободу решать, как его получить. Между вашим кодом и реальной работой на executor’ах стоят два независимых слоя: Catalyst решает, что вычислять (планирование), а Tungsten решает, как это вычислять эффективно на железе (исполнение).

Это разделение — ключ ко всему. Catalyst оперирует деревьями: он берёт ваш запрос как дерево операторов и многократно переписывает это дерево, пока не получит дешёвый план. Tungsten оперирует байтами и CPU: он берёт финальный план и генерирует под него низкоуровневый Java-код, который работает с данными в плотном бинарном формате мимо JVM-объектов. Если держать эти два слоя в голове по отдельности, .explain(true) перестаёт быть стеной текста и становится читаемой картой.

Дальше — самодостаточный разбор. Если хотите трогать всё руками, у нас есть бесплатный курс Spark Internals с runnable-песочницей в браузере, но статья работает и без него.

Catalyst: четыре дерева, одна труба

Catalyst — это оптимизатор запросов, построенный вокруг одной идеи: запрос — это дерево, оптимизация — это применение правил, переписывающих дерево. Каждое правило получает на вход дерево и возвращает новое дерево; Catalyst прогоняет наборы правил до тех пор, пока дерево не перестанет меняться (fixed point) или не исчерпается лимит итераций.

Запрос проходит через четыре последовательных представления. Понимание границ между ними и есть навык чтения планов.

ФазаЧто на входеЧто делаетРезультат
Parsed / Unresolved Logical PlanSQL-строка или DataFrame APIстроит дерево операторов, имена пока не провереныплан с 'unresolvedColumn
Analyzed Logical Planunresolved-план + Catalogрезолвит таблицы/колонки/типы, проверяет существованиетипизированный логический план
Optimized Logical Plananalyzed-планприменяет rule-based правила (RBO) и cost-based (CBO)дешёвый логический план
Physical Plan(s)optimized-планвыбирает конкретные алгоритмы, считает стоимостьодин выбранный SparkPlan

Unresolved: дерево без значений

Парсер превращает текст в дерево, но ничего о нём не знает. В выражении SELECT name FROM users WHERE age > 30 слова name, users, age — это просто строки. Catalyst помечает их как unresolved: он ещё не знает, существует ли таблица users, есть ли в ней колонка age и какого она типа. На этой фазе ловятся только синтаксические ошибки.

Analyzed: подключаем Catalog

Analyzer берёт unresolved-дерево и Catalog и разрешает каждую ссылку. users превращается в конкретное отношение со схемой, age — в колонку типа int под определённым внутренним идентификатором (exprId). Здесь же выводятся типы, вставляются неявные касты (age > 30 где age это bigint30 поднимется до bigint), разворачиваются * и проверяется, что GROUP BY согласован с агрегациями. Если колонки нет — именно тут вылетит AnalysisException. После этой фазы план семантически корректен, но ещё наивен: он буквально повторяет то, что вы написали.

Optimized: где происходит магия

Это сердце Catalyst. Оптимизатор прогоняет десятки правил, которые переписывают логический план, гарантированно сохраняя его результат, но удешевляя исполнение. Два правила окупают изучение всей фазы.

Predicate pushdown (проталкивание фильтров). Фильтр, написанный в любом месте запроса, Catalyst старается сдвинуть как можно ближе к источнику данных — в идеале внутрь самого чтения. Если источник это Parquet или JDBC, фильтр age > 30 уезжает в источник: Parquet пропустит целые row group по статистике min/max, а JDBC добавит WHERE в SQL и не потащит лишние строки по сети. Меньше прочитано — меньше всё последующее.

Projection pruning (отсечение колонок). Spark читает только те колонки, что реально нужны на выходе. В SELECT name FROM users WHERE age > 30 физически читаются ровно name и age; остальные сто колонок таблицы не трогаются вообще. Для колоночного формата это не оптимизация «на полях», а кратное сокращение I/O.

Сюда же входят constant folding (вычислить 1 + 2 в плане, а не на каждой строке), null propagation, boolean simplification, схлопывание соседних Filter и Project, выталкивание LIMIT вниз. Все эти правила — rule-based optimization (RBO): они детерминированы, не смотрят на данные и применяются всегда.

RBO против CBO

RBO опирается только на структуру запроса. Этого достаточно для «всегда хороших» переписываний, но не отвечает на вопросы, ответ на которые зависит от данных: в каком порядке джойнить три таблицы? использовать broadcast join или sort-merge? Здесь подключается cost-based optimization (CBO).

CBO смотрит на собранную статистику — число строк, размер в байтах, число уникальных значений, гистограммы — и численно оценивает стоимость альтернативных планов, выбирая дешёвый. Классический выигрыш: переупорядочивание джойнов, чтобы сначала схлопнулись самые селективные, и автоматический выбор broadcast join, когда CBO видит, что одна сторона помещается в память. Критическая оговорка: CBO бесполезен без актуальной статистики. Статистика собирается командой ANALYZE TABLE ... COMPUTE STATISTICS и сама не обновляется. Поэтому в реальных пайплайнах основную тяжесть несёт AQE — он корректирует план по фактическим размерам данных прямо во время выполнения, не полагаясь на заранее собранную статистику.

Physical plan: от «что» к «как»

Оптимизированный логический план говорит «соедини A и B по ключу k». Это всё ещё абстракция: «соединить» — не алгоритм. SparkPlanner применяет стратегии и порождает один или несколько физических планов, где каждый логический оператор заменён конкретным алгоритмом. Логический Join становится BroadcastHashJoinExec, SortMergeJoinExec или ShuffledHashJoinExec — в зависимости от размеров и конфигурации. Среди кандидатов выбирается план с наименьшей оценочной стоимостью, и именно он уходит в исполнение. На этой фазе появляются Exchange (это shuffle), Sort, HashAggregate — операторы, которых не было в логическом плане, потому что они описывают физику, а не семантику.

Tungsten: исполнение близко к железу

Catalyst отдаёт физический план. Теперь надо его выполнить, и здесь вступает Tungsten — движок исполнения, переписанный, чтобы Spark упирался в CPU и память так же эффективно, как нативный код, а не как типичное JVM-приложение, тонущее в указателях и сборке мусора. Три столпа.

Бинарный формат вне кучи

Обычный JVM-объект для строки данных — это набор разбросанных по куче объектов с заголовками, выравниванием и указателями: одно поле int в боксированном Integer стоит десятки байт и порождает работу для GC. Tungsten хранит строку как UnsafeRow — плоский непрерывный участок памяти: блок фиксированной длины со смещениями полей, за ним блок переменной длины для строк и массивов. Такая строка живёт off-heap, читается прямыми обращениями по смещению без разыменования объектов и почти не нагружает GC. Бонус: сравнение и хеширование можно делать прямо над байтами, не материализуя объекты — это разгоняет сортировку и агрегацию.

Cache-aware вычисления

Tungsten проектирует структуры данных под иерархию кэшей CPU. Алгоритмы сортировки и агрегации работают с компактными массивами, где рядом лежат ключ-указатель и префикс ключа, так что сравнение часто решается, не выходя за пределы L1/L2-кэша и не прыгая по случайным адресам в куче. Для процессора, который простаивает сотни тактов в ожидании промаха кэша, плотная раскладка иногда важнее самого алгоритма.

Whole-stage codegen

Самый заметный рычаг. Классические движки используют Volcano model: каждый оператор — итератор, строки тянутся через next() вверх по цепочке. На миллиард строк и пять операторов это пять миллиардов виртуальных вызовов, которые JVM не инлайнит, плюс boxing и промахи branch predictor.

Whole-stage codegen (WSCG) схлопывает целую цепочку операторов (scan, filter, project, часть join) в один сгенерированный Java-метод — один плотный цикл по строкам без межоператорных вызовов. Код компилируется на лету компилятором Janino, и JIT инлайнит его в нативный машинный код. Это превращает интерпретацию плана в исполнение, неотличимое от написанного вручную цикла.

spark.range(0, 1_000_000) \
    .selectExpr("id", "id % 7 AS bucket") \
    .filter("bucket > 2") \
    .groupBy("bucket").count() \
    .explain(True)

В выводе плана операторы, попавшие в один codegen-блок, помечены звёздочкой и общим номером, например *(1). Это и есть граница whole-stage: всё под одной звёздочкой исполняется как один скомпилированный метод. Exchange (shuffle) разрывает codegen — после него начинается новый блок *(2), потому что данные физически перетасовываются по сети между стадиями.

Как читать .explain(true)

.explain(true) печатает все четыре дерева подряд — ровно фазы Catalyst. Читать их надо снизу вверх (источник внизу, результат сверху) и по такому чек-листу:

== Parsed Logical Plan ==        # 'unresolvedColumn — имена ещё строки
== Analyzed Logical Plan ==      # типы и exprId проставлены
== Optimized Logical Plan ==     # фильтры уехали вниз, колонки отсечены
== Physical Plan ==              # алгоритмы выбраны, появились звёздочки

Что искать в Physical Plan — самой полезной секции:

  • PushedFilters: [GreaterThan(age,30)] в строке FileScan/Scan — predicate pushdown сработал, фильтр ушёл в источник. Если фильтр там, где вы его ждёте в источнике, нет — он применяется уже после чтения, и вы тащите лишнее.
  • ReadSchema: struct<name:string,age:int> — projection pruning: читаются только перечисленные колонки. Видите там колонки, которые не используете, — что-то мешает отсечению.
  • *(1) HashAggregate со звёздочкой — оператор внутри whole-stage codegen. Оператор без звёздочки кодген не получил (частая причина — UDF на Python, разрывающий pipeline); это место для подозрений на медленность.
  • Exchange hashpartitioning(...) — это shuffle, граница стадии и разрыв codegen. Каждый Exchange — материализация и сеть; их число — первое, что стоит сокращать.
  • BroadcastHashJoin против SortMergeJoin — какую стратегию выбрал планировщик. Неожиданный SortMergeJoin на маленькой таблице часто означает, что у Spark нет статистики и стоит подсказать broadcast.
  • AdaptiveSparkPlan isFinalPlan=false — план будет уточнён AQE во время выполнения; финальную форму смотрите в Spark UI после завершения, а не в статическом explain.

Практический цикл оптимизации простой: сначала по логическому плану убедиться, что фильтры протолкнулись, а колонки отсеклись; затем по физическому плану посчитать Exchange и проверить, что тяжёлые операторы попали под звёздочку codegen. Это покрывает большинство «почему медленно» ещё до того, как вы тронете executor.memory.

Что забрать с собой

Catalyst и Tungsten — это разделение «что» и «как». Catalyst переписывает дерево запроса через RBO (всегда-полезные правила вроде predicate pushdown и projection pruning) и CBO (выбор по статистике/AQE), проводя его через стадии unresolved, analyzed, optimized, physical. Tungsten исполняет финальный план близко к железу: бинарный UnsafeRow вне кучи, cache-aware алгоритмы и whole-stage codegen, схлопывающий операторы в один JIT-инлайнящийся цикл. .explain(true) показывает обе половины напрямую — научитесь его читать, и большинство проблем производительности станут видны на плане, а не в логах после часа ожидания.

Хотите прогнать эти планы на живом Spark и потрогать каждую фазу — бесплатный курс Spark Internals с песочницей в браузере ведёт от RDD до кодогенерации, а соседние материалы по направлению Data Engineering дают остальную картину стека.

Ещё в направлении · Data Engineering

Все материалы направления →