AQE: автоматическое переключение join-стратегий в рантайме

Как AQE динамически переключает SortMergeJoin в BroadcastHashJoin уже во время выполнения запроса, собирая реальные размеры таблиц после фильтрации и устраняя дорогостоящий shuffle.

optimization

Проблема «слепого» планирования джоинов

Join - самая дорогая операция в аналитических workloads. Неправильный выбор стратегии соединения способен превратить запрос на несколько секунд в многоминутный shuffle-ад. Но ещё хуже то, что Catalyst Optimizer принимает это решение вслепую: до начала выполнения, основываясь на статистике, которая может быть устаревшей, неточной или вовсе отсутствовать.

AQE решает эту проблему радикально: откладывает выбор стратегии join'а до момента, когда реальные данные уже частично обработаны и их точный размер известен.

Ограничения статического Catalyst Optimizer

Catalyst строит физический план выполнения на основе статистики таблиц: числа строк, распределения значений, размера данных. Эту статистику он получает из метастора (Hive Metastore, Iceberg catalog) и обновляет командой ANALYZE TABLE ... COMPUTE STATISTICS.

На практике статистика редко актуальна по нескольким причинам:

Команда ANALYZE запускается нечасто. Для Bronze-слоя, куда непрерывно поступают новые данные, статистика актуальна только сразу после её сбора. Через день-два после загрузки новой порции данных Catalyst снова работает с устаревшими цифрами. Многие команды вообще не настраивают автоматический ANALYZE, полагаясь на дефолтное поведение.

Статистика не учитывает WHERE-условия. Это принципиальное ограничение. Предположим, таблица categories весит 100 МБ. Catalyst знает это из метастора и планирует SortMergeJoin, потому что 100 МБ > autoBroadcastJoinThreshold (10 МБ по умолчанию). Но в запросе есть фильтр WHERE season = 'winter' AND is_active = true, который оставит от этой таблицы лишь 2 МБ. Catalyst на этапе планирования не может точно знать, сколько строк пройдёт фильтр. Он использует эвристики (column statistics, histogram-based estimation), которые часто ошибаются в 10-100 раз.

Фильтры создают «эффект схлопывания». Строгий предикат по конкретной категории или дате может уменьшить 100-гигабайтную таблицу до 500 мегабайт. И в этот момент оптимальная стратегия полностью меняется: вместо дорогого SortMergeJoin можно сделать быстрый BroadcastHashJoin. Но Catalyst, планируя запрос статически, ещё не знает о таком схлопывании.

Сложные подзапросы и CTE ломают оценки. Когда одна часть запроса является результатом другой (WITH filtered AS (SELECT ... WHERE ...)) или когда данные приходят из оконной функции, Catalyst практически лишён возможности точно оценить кардинальность. Оценки становятся произведением нескольких неточных оценок - ошибка множится.

Три стратегии join'а и их стоимость

Прежде чем разбирать, как AQE переключает стратегии, нужно понять, что именно он переключает и почему это важно. Spark поддерживает несколько стратегий join'а, кардинально отличающихся по стоимости.

BroadcastHashJoin (BHJ) - самая быстрая стратегия для случаев, когда одна из таблиц помещается в память. Spark целиком загружает меньшую таблицу в память каждого Executor'а (broadcast), строит из неё хэш-таблицу, и затем для каждой строки большой таблицы делает мгновенный lookup в этой хэш-таблице. Никакого сетевого shuffle для большой таблицы - данные остаются там, где они есть. Никакой сортировки. Никакого merge.

SortMergeJoin (SMJ) - стандартная стратегия для случаев, когда обе таблицы большие. Spark делает shuffle обеих таблиц по join-ключу (данные с одним ключом оказываются на одном Executor'е), затем сортирует каждую сторону по ключу, и наконец сливает два отсортированных потока. Это надёжно и масштабируется на любые объёмы, но стоит дорого: два полных shuffle по сети + двойная сортировка - это гигабайты сетевого трафика и минуты CPU на сортировку.

ShuffleHashJoin (SHJ) - промежуточный вариант. Shuffle обеих таблиц по ключу происходит, как в SMJ, но вместо сортировки + merge одна сторона (меньшая) строит хэш-таблицу, а вторая итерируется по ней. Нет сортировки - быстрее SMJ, но требует, чтобы меньшая сторона после shuffle поместилась в память одного Executor'а.

Разница в производительности между BHJ и SMJ колоссальная: BHJ на практике в 5-20 раз быстрее SMJ для типичных star-schema запросов, потому что устраняет shuffle большой таблицы - самую дорогую операцию.

Диаграмма: стоимость join-стратегий

Схема показывает, что делает Spark при каждой стратегии join'а. Стрелки от большой таблицы наглядно демонстрируют, где возникает или не возникает дорогостоящий сетевой shuffle.

Ключевое отличие: в BHJ большая таблица (1 ТБ) вообще не участвует в shuffle. Данные читаются локально на каждом Executor'е и немедленно джоинятся с уже загруженной хэш-таблицей. В SMJ обе таблицы перегоняются по сети - это и есть главный источник потерь.


Механика динамической смены стратегий

AQE превращает момент окончания Map Stage в «точку пересмотра плана». Именно здесь, когда данные физически записаны на диски Executor'ов и их реальный размер известен с точностью до байта, Spark решает: продолжать ли по исходному плану или переключиться на лучшую стратегию.

Концепция ShuffleQueryStage как точки принятия решения

В предыдущем уроке мы говорили о ShuffleQueryStage в контексте adaptive coalescing. В контексте join optimization этот же механизм работает иначе.

Когда Spark выполняет запрос с JOIN, он разбивает его на стадии. Первые стадии читают и фильтруют данные из обеих сторон джоина, записывая результаты в shuffle-файлы (Map Stage). После завершения Map Stage для обеих сторон AQE оказывается в уникальной позиции: он точно знает, сколько байт данных содержит каждая сторона после фильтрации.

Это и есть «точка невозврата» в понимании оптимального join'а. Теперь, а не на этапе планирования, Spark может принять взвешенное решение.

Шаг за шагом: как SMJ превращается в BHJ

Рассмотрим конкретный сценарий: JOIN таблицы кликов (500 ГБ) с таблицей категорий (100 МБ), при этом запрос фильтрует категории по условию WHERE season = 'winter'.

SELECT c.user_id, cat.name, count(*) as click_count
FROM clicks c
JOIN categories cat ON c.category_id = cat.id
WHERE cat.season = 'winter'
GROUP BY c.user_id, cat.name

Статический план (без AQE или с AQE до выполнения):

Catalyst видит: clicks = 500 ГБ, categories = 100 МБ > 10 МБ порог broadcast. Оба dataset'а большие → планирует SortMergeJoin для обоих.

Что происходит с AQE:

  1. Spark запускает Map Stage для таблицы categories (стадия чтения и фильтрации). Каждая map-таска читает свой чанк из Parquet, применяет предикат season = 'winter', записывает отфильтрованные данные в shuffle-файлы.

  2. Map Stage завершается. MapOutputTracker сообщает Драйверу: реальный объём отфильтрованных данных categories - 1.8 МБ (из 100 МБ осталось только 1.8 МБ зимних категорий).

  3. AQE Optimizer проверяет: 1.8 МБ ≤ spark.sql.adaptive.autoBroadcastJoinThreshold (10 МБ по умолчанию)? Да! Принимается решение переключить стратегию.

  4. Spark отменяет запланированный Shuffle для таблицы clicks. Map Stage clicks, если ещё не начался - не начнётся с shuffle-write. Если начался - результаты будут прочитаны локально.

  5. Данные categories (1.8 МБ) broadcast'ятся на все Executor'ы. Каждый Executor строит из них локальную хэш-таблицу.

  6. Map Stage для clicks читает данные локально и немедленно джоинит с хэш-таблицей категорий. Никакого сетевого трафика для 500 ГБ кликов. Никакой сортировки.

Результат: вместо двух полных shuffle (100 ГБ + 500 ГБ по сети) происходит broadcast 1.8 МБ × N Executor'ов - экономия в сотни раз.

Диаграмма: жизненный цикл runtime join rewrite

Обратите внимание на ключевой момент: AQE Decision Point происходит между стадиями. Это принципиально - Spark не прерывает выполнение на середине стадии. Он ждёт завершения Map Stage, смотрит на реальные размеры и принимает решение о следующей стадии. Именно поэтому этот механизм называется inter-stage optimization, а не intra-stage.

LocalShuffleReader: устранение сетевого трафика

При обычном SortMergeJoin shuffle работает так: каждый Executor записывает свои данные в файлы, разбитые по partition ID. Потом каждый Executor читает со всех остальных Executor'ов нужные ему партиции. Это и есть shuffle: данные летят по сети в обоих направлениях.

Когда AQE переключается на BroadcastHashJoin в рантайме, возникает вопрос: что делать с уже записанными shuffle-файлами большой таблицы (clicks)? В них данные уже «перемешаны» (shuffled), но BHJ не нуждается в этом перемешивании - он читает данные локально.

Здесь вступает LocalShuffleReader (управляется параметром spark.sql.adaptive.localShuffleReader.enabled). Вместо того чтобы читать shuffle-файлы в порядке partition ID (что требует сетевого обращения к другим Executor'ам), LocalShuffleReader читает только локальные shuffle-файлы на своём диске. Каждый Executor читает только то, что сам же и записал.

Это значит: нулевой сетевой трафик для большой таблицы. Вся сеть используется только для broadcast малой таблицы (1.8 МБ × число Executor'ов).

В Spark UI это видно как Shuffle Read = 0 (или очень малое значение) для стадии, где срабатывает LocalShuffleReader.


Настройка параметров динамического выбора join-стратегии

Поведение AQE-оптимизации join'ов контролируется несколькими параметрами, которые важно понимать не изолированно, а в их взаимодействии.

spark.sql.adaptive.autoBroadcastJoinThreshold

spark.conf.set("spark.sql.adaptive.autoBroadcastJoinThreshold", "10m")

Это AQE-специфичный порог broadcast'а, который используется для принятия рантайм-решений. Он работает в паре со статическим spark.sql.autoBroadcastJoinThreshold, но применяется уже после того, как Spark увидел реальный размер данных после Map Stage.

Важное различие между двумя параметрами:

  • spark.sql.autoBroadcastJoinThreshold - статический порог. Catalyst проверяет его до выполнения, сравнивая с оценкой из статистики метастора. Если оценка говорит «таблица больше порога» - Catalyst планирует SMJ. Если меньше - BHJ сразу в исходном плане.

  • spark.sql.adaptive.autoBroadcastJoinThreshold - AQE-порог. Применяется в рантайме, когда AQE сравнивает его с реальным измеренным размером таблицы после Map Stage.

Когда spark.sql.adaptive.autoBroadcastJoinThreshold не задан явно, AQE использует значение spark.sql.autoBroadcastJoinThreshold. Но задавая их независимо, можно получить тонкий контроль: статически планировать SMJ (консервативно), но разрешить AQE переключить на BHJ в рантайме для таблиц до, например, 50 МБ.

# Статический планировщик: broadcast только до 10 МБ (консервативно)
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "10m")

# AQE в рантайме: broadcast до 50 МБ (агрессивнее, т.к. знаем реальный размер)
spark.conf.set("spark.sql.adaptive.autoBroadcastJoinThreshold", "50m")

Почему AQE-порог можно ставить выше? Потому что статический планировщик работает с оценками, которые могут ошибаться. Если выставить статический порог 50 МБ, Catalyst может попытаться broadcast'ить таблицу, которую он оценил в 40 МБ, а реально она весит 800 МБ - это OOM. AQE же работает с реальными измеренными данными, поэтому риск ошибки минимален.

spark.sql.adaptive.localShuffleReader.enabled

spark.conf.set("spark.sql.adaptive.localShuffleReader.enabled", "true")

По умолчанию: true. Включает LocalShuffleReader при динамическом переключении на BHJ. Отключать имеет смысл только для диагностики - чтобы понять, насколько LocalShuffleReader реально снижает сетевой трафик в вашем конкретном случае.

spark.sql.autoBroadcastJoinThreshold = -1

Иногда полезный приём: отключить статический broadcast, оставив только AQE-динамический.

# Отключаем статический broadcast
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")

# AQE будет делать broadcast только на основе реальных размеров
spark.conf.set("spark.sql.adaptive.autoBroadcastJoinThreshold", "64m")
spark.conf.set("spark.sql.adaptive.enabled", "true")

Это полезно, когда статистика в метасторе ненадёжна и вы боитесь, что Catalyst попытается broadcast'ить большую таблицу по ошибке ещё на этапе планирования.

spark.sql.adaptive.nonEmptyPartitionRatioForBroadcast

spark.conf.set("spark.sql.adaptive.nonEmptyPartitionRatioForBroadcast", "0.2")

Дополнительная эвристика для AQE. Если менее 20% партиций одной из сторон join'а непустые (значит, данные очень сильно разрежены), AQE воздержится от принятия решения о broadcast'е даже если суммарный размер данных меньше порога. Это защита от edge-case: крохотные данные, разбитые по миллионам почти пустых партиций, могут давать нулевой суммарный размер, но при этом быть трудны для broadcast'а.

Взаимодействие со стратегией ShuffleHashJoin

AQE может переключать не только в BHJ, но и в ShuffleHashJoin (SHJ) - промежуточную стратегию, которую Catalyst часто игнорирует в пользу более консервативного SortMergeJoin.

Переключение SMJ → SHJ происходит, когда:

  • Одна сторона join'а достаточно мала, чтобы поместиться в памяти одного Executor'а
  • Но слишком велика для полного broadcast'а на все Executor'ы

В этом случае Spark делает shuffle обоих датасетов по ключу (как в SMJ), но затем не сортирует, а строит хэш-таблицу из меньшей стороны на каждом Executor'е. Экономия: нет сортировки (O(n log n) → O(n)), что существенно для больших датасетов.

Параметр spark.sql.adaptive.maxShuffledHashJoinLocalMapThreshold контролирует, при каком соотношении размеров сторон AQE рассмотрит SHJ как кандидата.


Как читать адаптивные join'ы в Spark UI и explain()

Умение читать адаптированный план - обязательный навык для диагностики. Только по плану можно убедиться, что AQE действительно переключил стратегию, а не сохранил изначальный SMJ.

explain() до и после выполнения

Статический план (до выполнения - isFinalPlan=false):

query = clicks.join(categories, "category_id").groupBy("user_id").count()
query.explain("formatted")
== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=false
+- HashAggregate(keys=[user_id#12], ...)
   +- Exchange hashpartitioning(user_id#12, 200)
      +- HashAggregate(...)
         +- SortMergeJoin [category_id#5], [id#78]    ← SMJ в статическом плане
            :- Sort [category_id#5 ASC]
            :  +- Exchange hashpartitioning(category_id#5, 200)
            :     +- Scan parquet clicks
            +- Sort [id#78 ASC]
               +- Exchange hashpartitioning(id#78, 200)
                  +- Scan parquet categories

Видим SortMergeJoin - именно он был запланирован статически, потому что categories оценена как 100 МБ > 10 МБ порога broadcast.

Финальный план (после выполнения - isFinalPlan=true):

# После query.collect() или query.show()
query.explain("formatted")
== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=true
+- == Final Plan ==
   HashAggregate(keys=[user_id#12], ...)
   +- HashAggregate(...)
      +- BroadcastHashJoin [category_id#5], [id#78], Inner, BuildRight   ← BHJ!
         :- AQEShuffleRead local                                          ← LocalShuffleReader!
         :  +- ShuffleQueryStage 0, Statistics(sizeInBytes=487.3 GiB, ...)
         :     +- Exchange hashpartitioning(category_id#5, 200)
         :        +- Scan parquet clicks
         +- BroadcastQueryStage 1, Statistics(sizeInBytes=1.8 MiB, ...)  ← 1.8 МБ после фильтра!
            +- BroadcastExchange HashedRelationBroadcastMode
               +- Filter (season#89 = winter)
                  +- Scan parquet categories

Разберём ключевые операторы в финальном плане:

  • BroadcastHashJoin ... BuildRight: SMJ успешно заменён на BHJ. BuildRight означает, что правая сторона (categories) использована для построения broadcast хэш-таблицы.

  • AQEShuffleRead local: LocalShuffleReader читает shuffle-файлы clicks локально, без сетевого трафика. Это ключевое слово local - именно оно означает нулевой shuffle по сети.

  • ShuffleQueryStage 0, Statistics(sizeInBytes=487.3 GiB): размер данных clicks после фильтрации (большая таблица, её shuffle-файлы читаются локально).

  • BroadcastQueryStage 1, Statistics(sizeInBytes=1.8 MiB): реальный размер categories после фильтрации. Именно это значение сравнивалось с порогом broadcast. 1.8 МБ - вот откуда AQE узнал, что можно переключить на BHJ.

  • == Initial Plan == (в полном explain): раздел с исходным планом, который был до адаптации. Useful для сравнения «до и после».

Поиск маркеров в Spark UI

В Spark UI → SQL tab → найдите нужный запрос по Job ID. В DAG-графе ищите:

  1. Серые «bubbles» с надписью AQE: маркеры стадий, переплановонных в рантайме
  2. Оператор BroadcastExchange вместо ShuffleExchange для малой таблицы
  3. Метрику Shuffle Read (local) в деталях стадии clicks: должна быть ненулевой, тогда как Shuffle Read (remote) - нулевой или минимальной

В метриках стадии clicks при успешном LocalShuffleReader:

Метрика Ожидаемое значение
Shuffle Read (remote) 0 Б (!)
Shuffle Read (local) ~487 ГБ
Task deserialization time низкий
Scheduler delay низкий

Нулевой Shuffle Read (remote) при большом Shuffle Read (local) - это «подпись» LocalShuffleReader. Данные читались с локального диска, а не гонялись по сети.


Практика: спасение пайплайна от лишнего Shuffle

Бизнес-кейс и подготовка данных

Смоделируем реальный сценарий из аналитики e-commerce: нужно найти суммарную выручку по пользователям для зимних категорий. Таблица кликов - 500 ГБ, таблица категорий - 100 МБ, но после фильтра season = 'winter' остаётся лишь 2% от категорий.

from pyspark.sql import SparkSession
from pyspark.sql import functions as F
import time

def create_spark(aqe_enabled: bool) -> SparkSession:
    """
    Создаёт SparkSession с заданными настройками AQE.

    Намеренно отключаем статический broadcast (autoBroadcastJoinThreshold=-1),
    чтобы без AQE Spark гарантированно планировал SMJ, - так эффект AQE виден чисто.
    """
    return (
        SparkSession.builder
        .appName(f"AQE-Join-{'enabled' if aqe_enabled else 'disabled'}")
        .master("local[4]")

        # Отключаем статический broadcast - хотим наблюдать только AQE-эффект
        .config("spark.sql.autoBroadcastJoinThreshold", "-1")

        # AQE: глобальный тумблер
        .config("spark.sql.adaptive.enabled", str(aqe_enabled).lower())

        # AQE: порог для рантайм-broadcast (применяется только при AQE enabled)
        .config("spark.sql.adaptive.autoBroadcastJoinThreshold", "10m")

        # LocalShuffleReader: читать shuffle локально после переключения на BHJ
        .config("spark.sql.adaptive.localShuffleReader.enabled", "true")

        # Достаточно партиций для демонстрации
        .config("spark.sql.shuffle.partitions", "200")
        .getOrCreate()
    )


def generate_datasets(spark: SparkSession):
    """
    Генерирует два датасета, имитирующих структуру e-commerce:
    - clicks: таблица фактов (большая), ~300 тысяч строк ≈ имитация большой таблицы
    - categories: справочник (меньше), 10000 строк, из которых только 200 «зимних»

    В продакшн-версии clicks был бы терабайтным, но для демо масштабируем вниз,
    сохраняя правильное соотношение: clicks в 30x больше categories.
    """
    # Таблица кликов: большая таблица фактов
    clicks = (
        spark.range(300_000)
        .select(
            F.col("id").alias("click_id"),
            (F.col("id") % 1000).cast("long").alias("user_id"),
            # category_id: 0..999, но только 20 из них - «зимние»
            (F.col("id") % 1000).cast("long").alias("category_id"),
            (F.rand(seed=42) * 999 + 1).alias("revenue"),
            F.current_timestamp().alias("click_ts"),
        )
        .cache()
    )

    # Таблица категорий: справочник
    # Из 10000 категорий только 200 помечены как season='winter'
    # Это 2% от данных - после фильтра таблица схлопнется до маленького размера
    categories = (
        spark.range(10_000)
        .select(
            F.col("id").alias("cat_id"),
            (F.concat(F.lit("Category_"), F.col("id"))).alias("name"),
            # Только каждая 50-я категория - зимняя (2% от всех)
            F.when(F.col("id") % 50 == 0, F.lit("winter"))
             .otherwise(F.lit("other"))
             .alias("season"),
            F.lit(True).alias("is_active"),
        )
        .cache()
    )

    # Прогреваем кэш
    print(f"  Clicks: {clicks.count():,} строк")
    print(f"  Categories: {categories.count():,} строк (из них зимних: "
          f"{categories.filter('season = \"winter\"').count()})")

    return clicks, categories


def run_winter_analysis(spark: SparkSession, clicks, categories, label: str):
    """
    Запрос: суммарная выручка по пользователям для зимних категорий.

    Ключевой момент: фильтр WHERE season='winter' оставляет от categories
    только 2% строк. Без AQE это не влияет на план - Spark сделает SMJ.
    С AQE он увидит реальный размер после фильтра и переключится на BHJ.
    """
    query = (
        clicks
        .join(
            categories.filter(F.col("season") == "winter"),
            clicks.category_id == categories.cat_id,
            "inner"
        )
        .groupBy("user_id", "name")
        .agg(
            F.sum("revenue").alias("total_revenue"),
            F.count("*").alias("click_count"),
        )
        .orderBy(F.desc("total_revenue"))
    )

    print(f"\n{'='*65}")
    print(f"  {label}")
    print(f"{'='*65}")

    # Показываем статический план (до выполнения)
    print("\n--- Статический план (isFinalPlan=false) ---")
    query.explain(mode="simple")

    # Замеряем реальное время
    start = time.time()
    result = query.collect()
    elapsed = time.time() - start

    print(f"\n  Результат: {len(result)} строк")
    print(f"  Время выполнения: {elapsed:.2f}с")

    # Показываем финальный план (после выполнения)
    print("\n--- Финальный план (isFinalPlan=true) ---")
    query.explain(mode="simple")

    return elapsed

Шаг 1: Запуск без AQE - наблюдаем SortMergeJoin

print("=== ЭКСПЕРИМЕНТ 1: БЕЗ AQE ===")
spark_no_aqe = create_spark(aqe_enabled=False)
clicks_no, cats_no = generate_datasets(spark_no_aqe)

time_no_aqe = run_winter_analysis(
    spark_no_aqe, clicks_no, cats_no,
    "Без AQE: autoBroadcastJoinThreshold=-1 → гарантированный SMJ"
)

Ожидаемый вывод статического и финального плана (они совпадают - без AQE план не меняется):

== Physical Plan ==
*(5) HashAggregate(...)
+- Exchange hashpartitioning(user_id#5, 200)
   +- *(4) HashAggregate(...)
      +- *(4) SortMergeJoin [category_id#3], [cat_id#21], Inner
         :- *(2) Sort [category_id#3 ASC]
         :  +- Exchange hashpartitioning(category_id#3, 200)   ← shuffle clicks!
         :     +- *(1) Scan InMemoryRelation (clicks)
         +- *(3) Sort [cat_id#21 ASC]
            +- Exchange hashpartitioning(cat_id#21, 200)       ← shuffle cats!
               +- *(1) Filter (season='winter')
                  +- Scan InMemoryRelation (categories)

Что здесь происходит: два Exchange hashpartitioning - это два полных shuffle. Clicks (все 300 тысяч строк) перемешиваются по сети. Categories тоже перемешиваются, пусть даже 98% строк отфильтровано - данные сначала фильтруются на map-стороне, а потом оставшиеся 2% всё равно записываются в shuffle-файлы и читаются по сети. В продакшне с 500 ГБ clicks это означало бы 500 ГБ сетевого трафика.

В Spark UI вы увидите:

  • Стадия 0: scan + filter clicks (300 тыс строк) → shuffle write 200 файлов
  • Стадия 1: scan + filter categories → shuffle write (пусть маленький, но есть)
  • Стадия 2: shuffle read с обеих сторон → SMJ → sort → merge

Шаг 2: Запуск с AQE - наблюдаем SMJ → BHJ

print("\n\n=== ЭКСПЕРИМЕНТ 2: С AQE ===")
spark_aqe = create_spark(aqe_enabled=True)
clicks_aqe, cats_aqe = generate_datasets(spark_aqe)

time_aqe = run_winter_analysis(
    spark_aqe, clicks_aqe, cats_aqe,
    "С AQE: autoBroadcastJoinThreshold=10MB → динамический BHJ"
)

Статический план идентичен примеру без AQE - AQE ещё не выполнялся. Но финальный план после collect() будет кардинально другим:

== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=true
+- == Final Plan ==
   *(4) HashAggregate(...)
   +- *(4) BroadcastHashJoin [category_id#3], [cat_id#21], Inner, BuildRight
      :- AQEShuffleRead local                          ← LocalShuffleReader для clicks!
      :  +- ShuffleQueryStage 0, Statistics(sizeInBytes=...)
      :     +- Exchange hashpartitioning(category_id#3, 200)
      :        +- *(1) Scan InMemoryRelation (clicks)
      +- BroadcastQueryStage 1, Statistics(sizeInBytes=X.X KiB) ← реальный размер!
         +- BroadcastExchange HashedRelationBroadcastMode
            +- *(2) Filter (season='winter')
               +- *(2) Scan InMemoryRelation (categories)

Ключевые отличия от плана без AQE:

  • SortMergeJoin заменён на BroadcastHashJoin, BuildRight - categories используется для broadcast, clicks обрабатывается локально

  • AQEShuffleRead local - clicks читаются с локального диска, не по сети

  • BroadcastQueryStage 1 - AQE материализовал categories как broadcast-таблицу, и в Statistics видим реальный размер после фильтрации

  • Нет Sort операторов - сортировка для SMJ больше не нужна

Шаг 3: Сравнение результатов и диагностика

print(f"\n{'='*65}")
print(f"  ИТОГОВОЕ СРАВНЕНИЕ")
print(f"{'='*65}")
print(f"  Без AQE (SortMergeJoin):     {time_no_aqe:.2f}с")
print(f"  С AQE (BroadcastHashJoin):   {time_aqe:.2f}с")
if time_no_aqe > 0:
    print(f"  Ускорение:                   {time_no_aqe/time_aqe:.1f}x")

print("""
  Что изменилось (ожидаемые наблюдения):
  ─────────────────────────────────────────────────────────────
  Shuffle:
    Без AQE: 2 Exchange (clicks + cats) → сетевой трафик обеих сторон
    С AQE:   1 LocalShuffleRead (clicks, локально) + broadcast 1 таблицы

  Сортировка:
    Без AQE: 2 Sort оператора (обе стороны) - CPU-intensive
    С AQE:   0 Sort операторов - BHJ не требует сортировки

  Стадии в Spark UI:
    Без AQE: scan/shuffle write × 2 + shuffle read/sort/merge
    С AQE:   scan/shuffle write (clicks) + broadcast + local read/join
""")

Шаг 4: Тонкая настройка порога broadcast для сложных запросов

def analyze_join_thresholds(spark_base, clicks, categories):
    """
    Исследуем влияние разных значений autoBroadcastJoinThreshold на поведение AQE.

    Это помогает понять: при каком пороге AQE будет агрессивнее переключать на BHJ,
    и когда это может привести к проблемам (OOM из-за слишком большого broadcast).
    """
    thresholds = ["1m", "5m", "10m", "50m", "100m"]

    print("\n=== Анализ влияния autoBroadcastJoinThreshold ===")
    print(f"{'Порог':>10} | {'Время':>10} | {'Ожидаемая стратегия'}")
    print(f"{'-'*10}-+-{'-'*10}-+-{'-'*40}")

    for threshold in thresholds:
        spark = (
            SparkSession.builder
            .appName(f"AQE-Join-threshold-{threshold}")
            .master("local[4]")
            .config("spark.sql.autoBroadcastJoinThreshold", "-1")
            .config("spark.sql.adaptive.enabled", "true")
            .config("spark.sql.adaptive.autoBroadcastJoinThreshold", threshold)
            .config("spark.sql.shuffle.partitions", "200")
            .getOrCreate()
        )

        query = (
            clicks.join(
                categories.filter(F.col("season") == "winter"),
                clicks.category_id == categories.cat_id,
            )
            .groupBy("user_id").agg(F.sum("revenue").alias("total"))
        )

        # Размер categories после фильтра - примерно X КБ в нашем датасете
        # Реальный размер определяется после выполнения
        start = time.time()
        query.count()
        elapsed = time.time() - start

        # Оцениваем стратегию по ожидаемым размерам (в реальности смотрим explain)
        # filtered cats ≈ 2% × 10000 строк × ~50 байт/строка ≈ 10 КБ
        strategy = "BHJ (broadcast)" if threshold not in ["1m"] else "SMJ (слишком мало?)"
        print(f"{threshold:>10} | {elapsed:>8.2f}с | {strategy}")

        spark.stop()

Шаг 5: Демонстрация с тремя последовательными join'ами

def multi_join_demo(spark: SparkSession):
    """
    Демонстрирует, как AQE обрабатывает цепочку join'ов.

    Сценарий: star-schema с таблицей фактов и тремя измерениями.
    Без AQE: все три join'а - SortMergeJoin (всё большое).
    С AQE: часть переключается на BHJ в зависимости от реального размера после фильтров.
    """
    # Таблица фактов: большая
    facts = (
        spark.range(100_000)
        .select(
            F.col("id").alias("fact_id"),
            (F.col("id") % 500).alias("user_id"),
            (F.col("id") % 200).alias("product_id"),
            (F.col("id") % 100).alias("store_id"),
            (F.rand(seed=1) * 1000).alias("amount"),
        )
    )

    # Измерение 1: Users (большое, нет фильтра → будет SMJ)
    users = (
        spark.range(500)
        .select(
            F.col("id").alias("user_id"),
            F.concat(F.lit("User_"), F.col("id")).alias("user_name"),
            F.lit("RU").alias("country"),
        )
    )

    # Измерение 2: Products (среднее, фильтр оставит ~10% → BHJ вероятен)
    products = (
        spark.range(200)
        .select(
            F.col("id").alias("product_id"),
            F.concat(F.lit("Product_"), F.col("id")).alias("product_name"),
            F.when(F.col("id") % 10 == 0, F.lit("electronics"))
             .otherwise(F.lit("other"))
             .alias("category"),
        )
    )

    # Измерение 3: Stores (очень маленькое после фильтра → точно BHJ)
    stores = (
        spark.range(100)
        .select(
            F.col("id").alias("store_id"),
            F.concat(F.lit("Store_"), F.col("id")).alias("store_name"),
            F.when(F.col("id") % 20 == 0, F.lit("flagship"))
             .otherwise(F.lit("standard"))
             .alias("store_type"),
        )
    )

    # Запрос с тремя join'ами и разными условиями фильтрации
    result = (
        facts
        .join(users, "user_id")                                    # join 1: users (нет фильтра)
        .join(
            products.filter(F.col("category") == "electronics"),   # join 2: фильтр 10%
            "product_id"
        )
        .join(
            stores.filter(F.col("store_type") == "flagship"),      # join 3: фильтр 5%
            "store_id"
        )
        .groupBy("user_name", "product_name", "store_name")
        .agg(F.sum("amount").alias("total_sales"))
    )

    print("\n=== Многоуровневый join (3 измерения) ===")
    print("Статический план:")
    result.explain(mode="simple")

    result.collect()

    print("\nФинальный адаптированный план:")
    result.explain(mode="simple")

    # Интерпретация: смотрим, сколько BHJ появилось в финальном плане
    final_plan = result._sc._jvm.PythonSQLUtils.explainString(
        result._jdf.queryExecution(), "simple"
    )
    bhj_count = final_plan.count("BroadcastHashJoin")
    smj_count = final_plan.count("SortMergeJoin")
    print(f"\n  BroadcastHashJoin в финальном плане: {bhj_count}")
    print(f"  SortMergeJoin в финальном плане:     {smj_count}")

Best practices и подводные камни

Опасность OOM при агрессивном пороге broadcast

AQE работает с реальными размерами данных - это безопаснее статического планировщика. Но риски всё равно есть, и их нужно понимать.

Когда AQE решает broadcast'ить таблицу, он:

  1. Собирает данные со всех map-тасок на Драйвер (BroadcastExchange → Драйвер)
  2. Сериализует их в объект Broadcast на Драйвере
  3. Отправляет копию на каждый Executor

Это значит: таблица размером X МБ на диске после десериализации в JVM-памяти Драйвера займёт 4-10x больше места из-за объектного overhead'а Java. Таблица в 50 МБ в Parquet (columnar, compressed) может занять 200-500 МБ в JVM heap.

# ОПАСНО: порог 500 МБ без учёта распаковки в RAM
spark.conf.set("spark.sql.adaptive.autoBroadcastJoinThreshold", "500m")
# Реальное потребление памяти Драйвера: 500 МБ × 5 = 2.5 ГБ → OOM если spark.driver.memory=2g

# БЕЗОПАСНО: не превышать 1/4 от spark.driver.memoryOverheadFactor × driver memory
# При spark.driver.memory=4g → порог не более 200-300 МБ
spark.conf.set("spark.sql.adaptive.autoBroadcastJoinThreshold", "100m")

Также важен параметр spark.sql.broadcastTimeout (default 300 секунд). Если broadcast большой таблицы занимает более 5 минут (медленная сеть или большой датасет), стадия упадёт с org.apache.spark.SparkException: Could not execute broadcast in time.

# Для больших broadcast'ов или медленных сетей
spark.conf.set("spark.sql.broadcastTimeout", "600")  # 10 минут

Взаимодействие с ручными JOIN HINT'ами

AQE и join hints (broadcast(), shuffle_merge(), shuffle_hash()) работают по-разному. Понимание их приоритетов важно для предсказуемого поведения.

from pyspark.sql.functions import broadcast

# Вариант 1: явный hint BROADCAST - всегда переопределяет AQE
# Spark попытается broadcast'ить независимо от размера и настроек AQE
result_forced = clicks.join(broadcast(categories), "category_id")

# Вариант 2: нет hint'а - AQE сам решает на основе реального размера
result_adaptive = clicks.join(categories.filter(...), "category_id")

# Вариант 3: hint MERGE - принудительно SMJ, AQE не переключит на BHJ
result_merge = clicks.join(categories.hint("merge"), "category_id")

Рекомендация: в первую очередь доверяйте AQE (вариант 2). Используйте broadcast() hint только когда уверены в размере таблицы независимо от фильтров. Используйте merge hint при диагностике или когда знаете о проблемах с broadcast (медленная сеть, OOM).

Когда AQE не переключает стратегию

AQE откажется переключать на BHJ в следующих случаях:

  • Реальный размер таблицы после Map Stage превышает порог autoBroadcastJoinThreshold
  • Join использует оператор OUTER JOIN или FULL OUTER JOIN (broadcast доступен не для всех типов)
  • Таблица находится на стороне outer в LEFT OUTER JOIN (broadcast возможен только для правой стороны)
  • Включён hint SHUFFLE_MERGE или SHUFFLE_HASH, явно запрещающий BHJ
  • spark.sql.adaptive.enabled=false
# Пример: LEFT OUTER JOIN - broadcast возможен только для правой (inner) стороны
result = large_table.join(small_table, "key", "left_outer")
# AQE может broadcast'ить small_table (правая сторона - она inner)
# Но НЕ может broadcast'ить large_table (левая сторона - она outer)

result = small_table.join(large_table, "key", "left_outer")
# Здесь AQE НЕ сможет broadcast'ить large_table - она на inner стороне,
# но слишком велика. И не может broadcast'ить small_table - она outer.
# Вывод: переформулируйте как RIGHT OUTER или внимательно думайте о порядке таблиц.

Синергия с другими оптимизациями Spark

Dynamic Join Selection не работает изолированно - он взаимодействует с другими механизмами, усиливая их эффект.

DPP + AQE Join Selection: Dynamic Partition Pruning (разбирается в следующем уроке) применяет результаты агрегации размерной таблицы к фильтрации партиций таблицы фактов ещё до shuffle. Это дополнительно уменьшает объём данных, который видит AQE при принятии решения о стратегии join'а.

Coalescing + Join Selection: если в цепочке запросов несколько join'ов, каждый ShuffleQueryStage становится точкой как для Adaptive Coalescing (мелкие партиции → крупные), так и для Join Selection (проверка порога broadcast). Они дополняют друг друга.


Итоги урока

AQE runtime join optimization устраняет фундаментальное ограничение статического планировщика: невозможность знать реальный размер данных после фильтрации до начала выполнения.

Переключение SortMergeJoin → BroadcastHashJoin в рантайме - это не просто техническая деталь. Для star-schema запросов с сильными фильтрами по размерным таблицам это разница между «гоняю терабайты по сети» и «broadcast'лю мегабайты на все узлы». В продакшне с реальными данными это ускорение в 5-20 раз на join-heavy пайплайнах.

Ключевые выводы для практики:

  • Доверяйте AQE: выставьте разумный autoBroadcastJoinThreshold (10-50 МБ) и дайте Spark принимать решения на основе реальных данных

  • Не злоупотребляйте broadcast() hint'ами - жёсткие подсказки устаревают при изменении данных

  • Следите за OOM: autoBroadcastJoinThreshold × 5-10 не должно превышать heap Драйвера
  • Используйте explain("formatted") после выполнения, а не до: только финальный план показывает, что AQE реально сделал

Домашнее задание

Дан следующий запрос с тремя последовательными join'ами:

spark = SparkSession.builder.master("local[4]") \
    .config("spark.sql.autoBroadcastJoinThreshold", "-1") \
    .config("spark.sql.adaptive.enabled", "true") \
    .getOrCreate()

orders = spark.read.parquet("data/orders/")          # 200 ГБ
customers = spark.read.parquet("data/customers/")    # 80 МБ, нет фильтра
products = spark.read.parquet("data/products/")      # 500 МБ, фильтр по is_active=True (10%)
promotions = spark.read.parquet("data/promotions/")  # 2 ГБ, фильтр по month='2024-12' (1%)

result = (
    orders
    .join(customers, "customer_id")
    .join(products.filter("is_active = true"), "product_id")
    .join(promotions.filter("month = '2024-12'"), "promo_id", "left_outer")
    .groupBy("customer_id", "product_id")
    .agg(F.sum("amount").alias("total"))
)

Задание:

  1. Определите, какие из трёх join'ов AQE скорее всего переключит на BHJ (исходя из ожидаемых размеров после фильтрации), а какой оставит как SMJ и почему.

  2. Настройте параметры AQE так, чтобы принудительно добиться переключения двух join'ов на BHJ без изменения spark.sql.autoBroadcastJoinThreshold (используйте только spark.sql.adaptive.autoBroadcastJoinThreshold).

  3. Объясните, почему третий join (promotions, LEFT OUTER) ведёт себя иначе и какие ограничения AQE накладывает на OUTER JOIN.

  4. Приложите вывод result.explain("formatted") после выполнения запроса с указанием строк, где видны операторы BroadcastHashJoin и AQEShuffleRead local.