AQE: adaptive coalescing - автоматический подбор числа партиций

Как Adaptive Query Execution устраняет проблему фиксированного spark.sql.shuffle.partitions, собирая реальную статистику после Shuffle Write и динамически схлопывая тысячи мелких партиций в оптимальные крупные таски.

optimization

Проблема фиксированного параллелизма после Shuffle

Каждый опытный Spark-инженер рано или поздно сталкивается с одним и тем же архитектурным тупиком: нужно задать единое число shuffle-партиций для пайплайна, который обрабатывает принципиально разные объёмы данных. В понедельник приходит 500 ГБ данных о продажах, в субботу - 3 ГБ. Параметр spark.sql.shuffle.partitions один. Какое значение поставить?

Это не вопрос вкуса - это вопрос архитектурного компромисса, который при любом фиксированном значении будет давать сбои. AQE (Adaptive Query Execution) устраняет этот компромисс, заменяя статическое решение динамическим.

Архитектурный тупик spark.sql.shuffle.partitions

Параметр spark.sql.shuffle.partitions определяет число партиций, которые Spark создаёт на выходе каждой широкой трансформации - groupBy, join, distinct, repartition без явного числа партиций. Значение по умолчанию - 200, что исторически подходило для кластеров начала 2010-х.

Проблема в том, что этот параметр задаётся до начала выполнения, когда Spark ещё ничего не знает о реальном объёме данных. Он фиксирован на весь lifetime сессии и не меняется от запроса к запросу.

Представьте конвейер, который агрегирует логи по дням:

  • Понедельник: 800 ГБ входных данных → нужно ≈ 1600 партиций по 500 МБ каждая
  • Суббота: 4 ГБ входных данных → нужно ≈ 8 партиций по 500 МБ каждая

Если поставить spark.sql.shuffle.partitions=1600, то в субботу Spark создаст 1600 партиций по 2.5 МБ каждая. Если поставить spark.sql.shuffle.partitions=8, то в понедельник Spark создаст 8 огромных партиций по 100 ГБ каждая - гарантированные OOM и spill на диск.

Именно поэтому принятая в индустрии практика - ставить значение «с запасом» (1000, 2000, даже 5000), принося в жертву эффективность маленьких джобов ради устойчивости больших. AQE ликвидирует эту жертву.

Эффект Over-partitioning: что происходит, когда партиций слишком много

Рассмотрим сценарий: spark.sql.shuffle.partitions=2000 при объёме данных 2 ГБ после shuffle.

Каждая из 2000 партиций получает в среднем 1 МБ данных. Spark запускает 2000 тасок на стадии Reduce.

Каждая такса - это не просто «прочитать данные». Это целый жизненный цикл:

  1. Планирование на Драйвере: TaskScheduler сериализует объект задачи, сохраняет его в памяти
  2. Сериализация и отправка: TaskDescription (≈50-200 КБ) отправляется на Executor по сети
  3. Десериализация на Executor: создание объектов Java/Scala, загрузка closure
  4. Инициализация JVM-структур: выделение BufferedReader, буферов shuffle fetch
  5. Само вычисление: чтение 1 МБ данных, агрегация - занимает 20-50 мс
  6. Сбор метрик: AccumulatorV2, отправка heartbeat обратно на Драйвер
  7. Подтверждение завершения: StatusUpdate от Executor к Драйверу

При 2000 тасках overhead на шаги 1-4 и 6-7 суммарно превышает время реального вычисления на шаге 5. Executor тратит больше времени на административную работу, чем на обработку данных.

Помимо CPU-расходов, есть и память: Драйвер хранит метаданные для каждой активной и завершённой задачи. При 2000 пустых тасок он держит в памяти тысячи объектов TaskInfo, замедляя GC и затрудняя работу TaskScheduler для других запросов.

И наконец - файловая система. Если после shuffle идёт запись в Parquet на S3/MinIO, каждая такса создаёт отдельный файл. 2000 тасок × (возможно, ещё N партиций по дате) = миллионы мелких файлов. Это проблема «Small Files Disease», которую мы разбираем в отдельном уроке, но AQE помогает предотвратить её у самого источника.

Диаграмма: Тупик фиксированного параллелизма

Схема ниже показывает, как один и тот же spark.sql.shuffle.partitions=1000 ведёт себя по-разному при разных объёмах входных данных. Оба сценария проблемные - но в противоположные стороны.

Оба сценария требуют разного значения параметра, но параметр один и фиксирован. AQE ломает эту дилемму: вы задаёте заведомо завышенное значение (например, 2000), а Spark после реального shuffle сам определяет, сколько из этих 2000 партиций нужно схлопнуть в одну.


Анатомия работы Adaptive Coalescing под капотом

Чтобы понять, как AQE решает задачу динамического определения числа партиций, нужно разобраться в том, как устроено выполнение Spark-запроса на уровне стадий (Stages).

Граница стадий и Shuffle Map Stage

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

Стадия, которая записывает shuffle-данные, называется Shuffle Map Stage (или Map Stage). Её задачи называются map-тасками, и каждая из них записывает данные в локальные shuffle-файлы на дисках Executor'ов, разбивая их на spark.sql.shuffle.partitions партиций по ключу.

После завершения Map Stage данные физически существуют на дисках Executor'ов. Вот здесь у AQE появляется уникальная возможность: все данные уже записаны и их размер известен.

Сбор точной статистики MapOutput

После того как Map Stage завершается, Shuffle Map Output Tracker (компонент Spark на Драйвере) собирает со всех Executor'ов точные сведения о размере каждого shuffle-вывода.

Для каждого из N map-тасок и каждой из P партиций Драйвер знает точный размер (в байтах). В итоге формируется матрица N × P, где каждая ячейка - это количество байт, которые map-таска i отправила в reduce-партицию j.

Суммируя по столбцам (суммируя по всем map-таскам для каждой партиции), AQE получает точный размер каждой reduce-партиции в байтах. Не оценку, не статистику по семплу — а точное значение, потому что данные уже записаны.

Именно эта информация - золото. Теперь AQE знает: «партиция 0 весит 45 МБ, партиция 1 - 2 МБ, партиция 2 - 800 КБ, партиция 3 - 3 МБ...» и так далее для всех P партиций.

Алгоритм объединения партиций

Имея точные размеры всех партиций, AQE запускает алгоритм жадного объединения соседних партиций (greedy consecutive coalescing). Он работает следующим образом:

  1. Берём первую партицию, начинаем накапливать размер в текущий «bucket»
  2. Добавляем следующую партицию к bucket'у
  3. Если суммарный размер bucket'а не превышает advisoryPartitionSizeInBytes (целевой размер) - продолжаем
  4. Если превышает - закрываем текущий bucket (он станет одной reduce-таской) и начинаем новый
  5. Повторяем до конца списка партиций

Ключевое слово - соседних. AQE может объединять только непрерывные партиции. Это обусловлено физической природой shuffle: данные в файлах упорядочены по partition ID, и читать их в произвольном порядке нельзя без значительного I/O overhead. Поэтому AQE не может взять партицию 5, партицию 100 и партицию 300 и объединить их в одну таску. Только 5+6+7+...+N.

После завершения алгоритма Spark вставляет в физический план специальный оператор — CustomShuffleReader. Он настроен читать несколько исходных shuffle-партиций как одну reduce-таску. Вместо 2000 тасок может быть запущено, например, 14 - и каждая из них прочитает ~140 исходных партиций.

Полная диаграмма жизненного цикла AQE Coalescing

Обратите внимание на разрыв между Map Stage и Reduce Stage - именно здесь AQE вставляет свою оптимизацию. Map Stage выполняется по «старому» плану (2000 partition ID), а Reduce Stage - уже по адаптированному (14 тасок). Это и называется runtime plan rewrite: план переписывается прямо во время выполнения, после получения реальных данных.


Ключевые конфигурации и рычаги управления

Поведение AQE контролируется набором взаимосвязанных параметров. Разберём каждый из них подробно, потому что знание только их названий без понимания их взаимодействия часто приводит к неожиданным результатам.

spark.sql.adaptive.enabled

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

Глобальный тумблер AQE. В Apache Spark 3.2+ включён по умолчанию. Если он выключен, никакая адаптивная оптимизация не работает - ни coalescing, ни skew join handling, ни динамическое переключение join-стратегий.

Отключать его стоит только в двух случаях:

  • Debugging: нужно воспроизвести поведение Spark 2.x или сравнить планы без адаптации
  • Deterministic тестирование: план должен быть строго воспроизводим между запусками

spark.sql.adaptive.coalescePartitions.enabled

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

Включает именно механизм adaptive coalescing - схлопывание мелких партиций. Когда AQE включён, но этот параметр выключен, AQE будет работать, но только для других оптимизаций (skew join, runtime broadcast). По умолчанию true.

spark.sql.adaptive.advisoryPartitionSizeInBytes

spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "128m")

Целевой (совещательный) размер партиции - главный рычаг управления coalescing'ом. Именно к этому размеру AQE стремится при объединении партиций. Слово «advisory» (совещательный) означает, что это рекомендация, а не жёсткий лимит: если одна партиция уже превышает это значение, AQE всё равно оставит её как есть, не разбивая (разбиение - задача skew join optimization).

Дефолтное значение: 64 МБ (в Spark 3.x), 128 МБ (в ряде дистрибутивов).

Когда поднимать до 128-256 МБ:

  • Мощные Executor'ы (16+ ядер, 64+ ГБ RAM): большая таска лучше использует кэш CPU и JIT
  • Запись в Parquet/ORC: целевой размер файла 128-256 МБ снижает количество файлов в lake
  • CPU-интенсивные агрегации: накладные расходы на запуск таски становятся пренебрежимо малы

Когда оставлять 64 МБ или снижать до 32 МБ:

  • Маленькие Executor'ы (4 ядра, 16 ГБ): большие партиции могут вызвать GC-паузы
  • Streaming micro-batch: меньший размер → более быстрые итерации
  • Операции с высокой кардинальностью (distinct, count distinct): нужно больше параллелизма

spark.sql.adaptive.coalescePartitions.minPartitionNum

spark.conf.set("spark.sql.adaptive.coalescePartitions.minPartitionNum", "1")

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

Это критически важный параметр для кластеров с высоким параллелизмом. Если у вас 200 воркер-ядер, и AQE схлопнул всё в 5 партиций, то 195 ядер простаивают. По умолчанию в Spark 3.x значение равно CoalescePartitionsUtil.DEFAULT_MIN_PARTITION_NUM, что обычно соответствует числу ядер кластера.

Явно задавать имеет смысл, когда дефолтное значение слишком агрессивно или слишком консервативно.

spark.sql.adaptive.coalescePartitions.parallelismFirst

spark.conf.set("spark.sql.adaptive.coalescePartitions.parallelismFirst", "true")

Баланс между параллелизмом и размером партиций. Когда true (дефолт в Spark 3.x), AQE в первую очередь стремится сохранить параллелизм кластера, и только потом смотрит на целевой размер. Фактически при parallelismFirst=true параметр advisoryPartitionSizeInBytes игнорируется, если его соблюдение означало бы слишком мало партиций.

Когда false, AQE строго следует advisoryPartitionSizeInBytes, даже если это означает меньше партиций, чем ядер кластера. Это полезно для сценариев записи (write-heavy workloads), где важнее размер выходных файлов, чем максимальный параллелизм чтения.

Совет: начните с дефолтного true, переключайте на false только если боретесь с проблемой мелких файлов на выходе.

spark.sql.shuffle.partitions - исходная точка

spark.conf.set("spark.sql.shuffle.partitions", "2000")

При включённом AQE этот параметр становится верхней границей числа партиций, а не целевым значением. Задавайте его «с запасом», опираясь на пиковый объём данных. AQE при необходимости схлопнет их.

Правило большого пальца: spark.sql.shuffle.partitions = (пиковый объём данных в ГБ) × 4. Для пика 500 ГБ → 2000 партиций. Для пика 2 ТБ → 8000 партиций.


Как читать адаптивный план в Spark UI

Понимание того, как AQE отображается в Spark UI и в explain(), - это практический навык, без которого сложно проверить, работает ли оптимизация так, как вы ожидаете.

Explain и AdaptiveSparkPlanExec

При вызове df.explain("formatted") до выполнения запроса вы увидите статический план — тот, который Catalyst построил без знания реальных данных. AQE-операторы в нём не видны: план ещё не выполнялся и статистики нет.

# Статический explain - план до выполнения
df.groupBy("region").agg(sum("revenue")).explain("formatted")
== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=false
+- HashAggregate(keys=[region#12], functions=[sum(revenue#34)])
   +- Exchange hashpartitioning(region#12, 2000), ENSURE_REQUIREMENTS, [id=#89]
      +- ...

Строка AdaptiveSparkPlan isFinalPlan=false говорит: «план будет адаптирован в рантайме». isFinalPlan=false означает, что оптимизация ещё не произошла.

После выполнения запроса (например, после df.write.parquet(...) или df.show()), можно вызвать explain() на том же объекте или посмотреть план в Spark UI. Теперь увидите:

== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=true
+- == Final Plan ==
   HashAggregate(keys=[region#12], functions=[sum(revenue#34)])
   +- AQEShuffleRead coalesced
      +- ShuffleQueryStage 0
         +- Exchange hashpartitioning(region#12, 2000), ENSURE_REQUIREMENTS, [id=#89]
            +- ...

Ключевые операторы в финальном плане:

  • ShuffleQueryStage: точка, где Spark «заморозил» выполнение и собрал MapOutput статистику. Номер после него (0, 1, 2...) - это порядковый номер стадии в рамках всего запроса.

  • AQEShuffleRead coalesced: CustomShuffleReader, который читает несколько исходных partition ID как одну reduce-таску. Слово coalesced показывает, что произошло именно схлопывание. Если бы план был адаптирован по причине skew join - здесь было бы skewed.

Визуальные маркеры в SQL-вкладке Spark UI

В Spark UI → вкладка SQL → найдите нужный запрос → нажмите на него → откроется DAG-граф.

В DAG при включённом AQE вы увидите:

  1. Серые узлы с надписью AQE - маркеры стадий, которые были оптимизированы в рантайме
  2. Оператор Exchange с числом справа - сколько партиций было создано на Map Stage
  3. Оператор CustomShuffleReader под Exchange - сколько партиций реально прочитала Reduce Stage

Ищите в описании CustomShuffleReader строку вида:

coalesced from 2000 to 14 partitions

Это и есть доказательство работы AQE Coalescing: из 2000 запланированных партиций Spark реально создал 14 reduce-тасок.

Метрики тасок для анализа эффекта

В Spark UI → Stages → найдите Reduce Stage. Ключевые метрики для анализа:

Метрика Без AQE (2000 партиций) С AQE (14 партиций)
Number of Tasks 2000 14
Median Task Duration 12 мс (пустые) 8.5 с
Max Task Duration 45 мс 12 с
Scheduler Delay 15-25% от Duration 2-5% от Duration
Task Deserialization Time 5-8 мс 5-8 мс
Shuffle Read 1 МБ / task 140 МБ / task

Колонка «Scheduler Delay» - самый показательный индикатор over-partitioning. Если она составляет 15-25% от времени таски, значит executor тратит четверть времени на административную работу, а не на реальные вычисления. AQE снижает этот процент до нормальных 2-5%.


Практика: сравнение поведения с AQE и без

Подготовка данных и сессии

Подготовим среду для экспериментов. Создадим два варианта SparkSession - с AQE и без — и датасет, который имитирует реальный ETL с сильной вариативностью объёма данных.

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

def create_spark(aqe_enabled: bool, shuffle_partitions: int = 2000) -> SparkSession:
    """
    Создаёт SparkSession с заданными настройками AQE.
    Параметр aqe_enabled позволяет сравнить поведение до и после включения AQE.
    shuffle_partitions=2000 симулирует настройку "с запасом" для пиковой нагрузки.
    """
    builder = (
        SparkSession.builder
        .appName(f"AQE-Demo-{'enabled' if aqe_enabled else 'disabled'}")
        .master("local[4]")

        # Фиксируем высокое число shuffle партиций, как на продакшн-кластере
        # для пиковой нагрузки. Без AQE это будет катастрофой на малых данных.
        .config("spark.sql.shuffle.partitions", str(shuffle_partitions))

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

        # Adaptive Coalescing: целевой размер партиции 64 МБ
        .config("spark.sql.adaptive.coalescePartitions.enabled", "true")
        .config("spark.sql.adaptive.advisoryPartitionSizeInBytes", "64m")

        # Минимум партиций - не ниже числа доступных ядер (4 в данном случае)
        .config("spark.sql.adaptive.coalescePartitions.minPartitionNum", "4")

        # Включаем логирование адаптивного плана для анализа
        .config("spark.sql.adaptive.logLevel", "WARN")
    )
    return builder.getOrCreate()


def generate_small_dataset(spark: SparkSession, size_mb: float = 50.0):
    """
    Генерирует датасет ~50 МБ - типичный «выходной день» для нашего ETL.
    При spark.sql.shuffle.partitions=2000 это создаст партиции по ~25 КБ,
    что является экстремальным примером over-partitioning.
    """
    # Оцениваем число строк для нужного объёма данных.
    # Строка с 5 колонками ≈ 100 байт → 50 МБ ≈ 500_000 строк
    num_rows = int(size_mb * 1024 * 1024 / 100)

    df = spark.range(num_rows).select(
        F.col("id"),
        # Регион: 10 уникальных значений - типичный GROUP BY ключ
        (F.col("id") % 10).cast("string").alias("region"),
        # Категория: 50 значений - второй уровень группировки
        (F.col("id") % 50).cast("string").alias("category"),
        # Выручка: случайное значение от 1 до 1000
        (F.rand(seed=42) * 999 + 1).alias("revenue"),
        # Дата: последние 30 дней
        (F.current_date() - (F.col("id") % 30).cast("int")).alias("sale_date"),
    )
    return df.cache()


def run_aggregation(spark: SparkSession, df, label: str):
    """
    Выполняет типичную GROUP BY агрегацию и замеряет время.
    Именно после GROUP BY происходит Shuffle с adaptivePartitions или без.
    """
    print(f"\n{'='*60}")
    print(f"  {label}")
    print(f"{'='*60}")

    # Запрос с двухуровневой агрегацией - two-stage shuffle
    query = (
        df
        .groupBy("region", "category")
        .agg(
            F.sum("revenue").alias("total_revenue"),
            F.count("*").alias("order_count"),
            F.avg("revenue").alias("avg_revenue"),
            F.countDistinct("sale_date").alias("active_days"),
        )
        .orderBy("total_revenue", ascending=False)
    )

    start = time.time()
    result = query.collect()  # Действие, которое реально запускает выполнение
    elapsed = time.time() - start

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

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

# Создаём сессию БЕЗ AQE
spark_no_aqe = create_spark(aqe_enabled=False, shuffle_partitions=2000)

# Генерируем 50 МБ данных - «выходной» объём
df_small = generate_small_dataset(spark_no_aqe, size_mb=50.0)
df_small.count()  # Прогреваем кэш

# Смотрим статический план - видим 2000 партиций
print("=== СТАТИЧЕСКИЙ ПЛАН (без AQE) ===")
df_small.groupBy("region").agg(F.sum("revenue")).explain("formatted")

В выводе explain увидите:

== Physical Plan ==
*(2) HashAggregate(keys=[region#12], functions=[sum(revenue#34)])
+- Exchange hashpartitioning(region#12, 2000), ENSURE_REQUIREMENTS, [id=#89]
   +- *(1) HashAggregate(keys=[region#12], ...

Цифра 2000 после hashpartitioning - это число shuffle партиций. Без AQE оно не изменится ни при каких обстоятельствах. Spark запустит ровно 2000 reduce-тасок, подавляющее большинство которых прочитает по 25-50 КБ данных.

# Запускаем и замеряем время
time_no_aqe = run_aggregation(
    spark_no_aqe,
    df_small,
    "Без AQE: spark.sql.shuffle.partitions=2000, данные=50МБ"
)

# После выполнения откройте Spark UI (обычно localhost:4040)
# Перейдите в Stages → найдите стадию с 2000 тасками
# Обратите внимание на Scheduler Delay в сводной таблице тасок

Что вы увидите в Spark UI:

  • 2000 завершённых тасок в Reduce Stage
  • Медианное время таски: 8-20 мс (из которых 5-8 мс - сериализация и scheduling overhead)
  • Scheduler Delay: 4-10 мс на каждую таску - это чистые потери на административную работу
  • Shuffle Read Size: 20-50 КБ на таску - мизерный объём полезной работы

Шаг 2: Запуск с AQE - наблюдаем адаптацию

# Создаём новую сессию С AQE
spark_aqe = create_spark(aqe_enabled=True, shuffle_partitions=2000)

# Используем те же самые данные того же объёма
df_small_aqe = generate_small_dataset(spark_aqe, size_mb=50.0)
df_small_aqe.count()

# Статический план всё ещё покажет 2000 - AQE ещё не выполнялся
print("=== СТАТИЧЕСКИЙ ПЛАН (с AQE, до выполнения) ===")
df_small_aqe.groupBy("region").agg(F.sum("revenue")).explain("formatted")
== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=false   ← Маркер: план будет адаптирован
+- HashAggregate(keys=[region#12], ...)
   +- Exchange hashpartitioning(region#12, 2000), ...

Статический план идентичен. AQE ещё не знает о реальных размерах партиций — этот план существует только на бумаге.

# Запускаем и замеряем время
time_aqe = run_aggregation(
    spark_aqe,
    df_small_aqe,
    "С AQE: advisoryPartitionSizeInBytes=64MB, данные=50МБ"
)

# Теперь смотрим финальный (адаптированный) план
print("=== ФИНАЛЬНЫЙ ПЛАН (с AQE, после выполнения) ===")
adapted_df = df_small_aqe.groupBy("region").agg(F.sum("revenue"))
adapted_df.explain("formatted")

После выполнения план изменится:

== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=true
+- == Final Plan ==
   *(2) HashAggregate(...)
   +- AQEShuffleRead coalesced          ← Вот он!
      +- ShuffleQueryStage 0, Statistics(sizeInBytes=48.0 MiB, ...)
         +- Exchange hashpartitioning(region#12, 2000), ...

== Initial Plan ==                       ← Исходный план для сравнения
   HashAggregate(...)
   +- Exchange hashpartitioning(region#12, 2000), ...

Оператор AQEShuffleRead coalesced говорит: «Я прочитал данные из 2000 shuffle-партиций, но объединил их в несколько крупных тасок». Сколько именно - видно в Spark UI.

Шаг 3: Анализ результатов и метрик

def analyze_stage_metrics(spark: SparkSession, has_aqe: bool):
    """
    Выводит ключевые метрики для понимания эффекта AQE.
    В реальном коде эти данные смотрим в Spark UI → Stages.
    Здесь используем SparkContext.statusTracker для программного доступа.
    """
    sc = spark.sparkContext
    # Получаем информацию о всех завершённых стадиях
    stage_ids = sc.statusTracker().getJobIdsForGroup(None)

    aqe_label = "С AQE" if has_aqe else "Без AQE"
    print(f"\n{'='*60}")
    print(f"  Анализ метрик: {aqe_label}")
    print(f"{'='*60}")

    # Ключевые наблюдения для интерпретации результата
    if has_aqe:
        print("""
  Ожидаемые метрики в Spark UI (Stages → Reduce Stage):
  ─────────────────────────────────────────────────────
  • Number of Tasks:     ~4 (после coalescing с minPartitionNum=4)
  • Median Task Duration: 0.5-2с (больше полезной работы)
  • Scheduler Delay:     2-5% от Duration (норма)
  • Shuffle Read/task:   ~12 МБ (вместо 25 КБ)
  • Files written:       4-8 файлов (вместо 2000)

  В DAG-графе (SQL tab):
  • Оператор: AQEShuffleRead coalesced
  • Надпись: "coalesced from 2000 to 4 partitions"
        """)
    else:
        print("""
  Ожидаемые метрики в Spark UI (Stages → Reduce Stage):
  ─────────────────────────────────────────────────────
  • Number of Tasks:     2000 (все запущены!)
  • Median Task Duration: 8-20 мс (в основном overhead)
  • Scheduler Delay:     15-30% от Duration (красный флаг!)
  • Shuffle Read/task:   ~25 КБ (мизерная полезная работа)
  • Files written:       2000 мелких файлов

  Это признаки over-partitioning:
  ─────────────────────────────────────────────────────
  ⚠ Scheduler Delay > 10% → executor тратит время на orchestration
  ⚠ Task Duration < 100мс → задача тривиальная, overhead доминирует
  ⚠ Shuffle Read < 1MB → партиция почти пустая
        """)


# Запускаем анализ для обоих вариантов
analyze_stage_metrics(spark_no_aqe, has_aqe=False)
analyze_stage_metrics(spark_aqe, has_aqe=True)

# Итоговое сравнение
print(f"\n{'='*60}")
print(f"  ИТОГОВОЕ СРАВНЕНИЕ")
print(f"{'='*60}")
print(f"  Без AQE:   {time_no_aqe:.2f}с")
print(f"  С AQE:     {time_aqe:.2f}с")
print(f"  Ускорение: {time_no_aqe/time_aqe:.1f}x")

Шаг 4: Тонкая настройка advisoryPartitionSizeInBytes

def benchmark_advisory_sizes(spark_base: SparkSession, df, sizes_mb: list):
    """
    Сравниваем влияние разных значений advisoryPartitionSizeInBytes.
    Понимание этого помогает найти баланс между параллелизмом и размером файлов.
    """
    results = {}

    for size_mb in sizes_mb:
        size_bytes = size_mb * 1024 * 1024

        # Создаём отдельную сессию для чистоты эксперимента
        spark = (
            SparkSession.builder
            .appName(f"AQE-Advisory-{size_mb}MB")
            .master("local[4]")
            .config("spark.sql.shuffle.partitions", "2000")
            .config("spark.sql.adaptive.enabled", "true")
            .config("spark.sql.adaptive.coalescePartitions.enabled", "true")
            .config("spark.sql.adaptive.advisoryPartitionSizeInBytes", str(int(size_bytes)))
            .config("spark.sql.adaptive.coalescePartitions.minPartitionNum", "1")
            # parallelismFirst=false → строго соблюдаем advisoryPartitionSizeInBytes
            .config("spark.sql.adaptive.coalescePartitions.parallelismFirst", "false")
            .getOrCreate()
        )

        df_local = df.repartition(2000)  # Принудительно 2000 исходных партиций

        start = time.time()
        query_result = (
            df_local
            .groupBy("region", "category")
            .agg(F.sum("revenue").alias("total"))
            .collect()
        )
        elapsed = time.time() - start

        results[size_mb] = elapsed
        print(f"  advisoryPartitionSizeInBytes={size_mb}MB: {elapsed:.2f}с")
        spark.stop()

    return results


# Тестируем разные целевые размеры партиций
print("\n=== Benchmark: advisoryPartitionSizeInBytes ===")
print("Данные: 50 МБ, shuffle.partitions=2000, parallelismFirst=false")
print()
# В реальном эксперименте запускаем с нашим датасетом
# benchmark_results = benchmark_advisory_sizes(spark_aqe, df_small_aqe, [16, 32, 64, 128, 256])

Шаг 5: Симуляция инкрементального ETL с переменным объёмом

def simulate_incremental_etl(spark: SparkSession, day_type: str):
    """
    Симулирует реальный ETL-пайплайн с разным объёмом данных:
    - 'weekday': типичный будний день (500 МБ)
    - 'weekend': выходной день (5 МБ)
    - 'peak': пиковый день (5 ГБ, конец месяца)

    Именно для этого сценария AQE наиболее ценен:
    один и тот же код и конфигурация работают оптимально на всех трёх объёмах.
    """
    sizes = {
        "weekday": 500,    # 500 МБ → ~250 строк × 2 МБ блоки
        "weekend": 5,      # 5 МБ → мизерный объём, 2000 партиций = катастрофа без AQE
        "peak":    5000,   # 5 ГБ → нужен высокий параллелизм
    }

    size_mb = sizes[day_type]
    # Масштабируем число строк к нужному размеру
    num_rows = int(size_mb * 1024 * 1024 / 100)

    raw = spark.range(num_rows).select(
        F.col("id"),
        (F.col("id") % 100).cast("string").alias("product_id"),
        (F.col("id") % 50).cast("string").alias("store_id"),
        (F.rand(seed=42) * 999 + 1).alias("amount"),
        (F.current_date() - (F.col("id") % 7).cast("int")).alias("txn_date"),
    )

    # Типичный ETL: агрегация по продуктам и магазинам
    silver = (
        raw
        .groupBy("product_id", "store_id", "txn_date")
        .agg(
            F.sum("amount").alias("daily_revenue"),
            F.count("*").alias("txn_count"),
        )
    )

    # Вторая агрегация - Gold layer
    gold = (
        silver
        .groupBy("product_id", "txn_date")
        .agg(
            F.sum("daily_revenue").alias("total_revenue"),
            F.sum("txn_count").alias("total_txns"),
            F.countDistinct("store_id").alias("active_stores"),
        )
    )

    start = time.time()
    row_count = gold.count()
    elapsed = time.time() - start

    print(f"  Day type: {day_type:8s} | Size: {size_mb:5d} МБ | "
          f"Rows: {row_count:6d} | Time: {elapsed:.2f}с")
    return elapsed


print("\n=== Симуляция инкрементального ETL (с AQE) ===")
print("Конфигурация: shuffle.partitions=2000, advisoryPartitionSizeInBytes=64MB")
print()
for day in ["weekend", "weekday", "peak"]:
    simulate_incremental_etl(spark_aqe, day)

Взаимодействие с другими механизмами AQE

Adaptive Coalescing - лишь одна из трёх ключевых оптимизаций, входящих в AQE. Они работают вместе и взаимно усиливают друг друга.

Почему это важно: Coalescing и Skew Optimization - противоположные операции. Coalescing объединяет мелкие партиции в крупные. Skew Optimization разбивает крупные партиции на мелкие. Обе задачи решаются в рамках одного ShuffleQueryStage: AQE анализирует статистику, и для каждой партиции независимо решает, что с ней сделать.

Взаимодействие с Dynamic Allocation

AQE и Dynamic Allocation (динамическое выделение Executor'ов) хорошо работают вместе, но управляют разными ресурсами:

  • AQE оптимизирует количество тасок (партиций) в рамках заданного числа Executor'ов
  • Dynamic Allocation оптимизирует количество Executor'ов под текущую нагрузку

В связке они дают двойной эффект: AQE схлопывает 2000 партиций до 14 тасок, а Dynamic Allocation освобождает лишние Executor'ы, которые теперь не нужны.

Ограничения: что AQE не может сделать

AQE - мощный инструмент, но у него есть принципиальные ограничения:

  1. Coalescing не решает проблему skew. Если одна партиция весит 50 ГБ, AQE не схлопнет её — он может только оставить её как есть. Для skew нужен отдельный механизм (Skew Join Optimization).

  2. AQE работает только на границах shuffle. Оптимизации внутри стадии (без shuffle) AQE не делает. Если у вас медленный сканер или тяжёлая UDF - AQE не поможет.

  3. Streaming ограничен. В Structured Streaming AQE работает в отдельных micro-batch'ах, но не может адаптировать план между ними. Stateful операторы обходят AQE.

  4. Кэшированные датасеты. После df.cache() данные материализованы, и AQE не может перепланировать доступ к ним. Адаптация происходит только при первом выполнении.


Практический чеклист настройки AQE в production

Ниже - рецепт, который работает для большинства batch-пайплайнов. Подстраивайте под свои данные.

def production_aqe_config(
    spark_builder,
    peak_data_gb: float,
    executor_memory_gb: int,
    num_executors: int,
    target_file_size_mb: int = 128,
):
    """
    Задаёт production-готовые настройки AQE на основе характеристик кластера.

    Параметры:
        peak_data_gb:       Пиковый объём данных после shuffle в ГБ
        executor_memory_gb: RAM одного Executor'а
        num_executors:      Количество Executor'ов
        target_file_size_mb: Целевой размер Parquet-файла на выходе

    Логика расчёта:
        - shuffle_partitions: правило 4x от пикового объёма (с запасом под AQE)
        - advisoryPartitionSizeInBytes: целевой размер файла - хотим, чтобы
          каждая AQE-таска создала ровно один файл нужного размера
        - minPartitionNum: число ядер кластера - чтобы кластер не простаивал
    """
    total_cores = num_executors * 4  # Предполагаем 4 ядра на Executor

    # Число shuffle партиций с запасом (AQE схлопнет лишние)
    shuffle_partitions = max(200, int(peak_data_gb * 1024 / 64 * 4))

    # Целевой размер = целевой размер файла (для хорошего layout в Lakehouse)
    advisory_size_bytes = target_file_size_mb * 1024 * 1024

    # Нижний лимит = число ядер (чтобы все заняты)
    min_partitions = total_cores

    return (
        spark_builder
        # Включаем AQE
        .config("spark.sql.adaptive.enabled", "true")
        .config("spark.sql.adaptive.coalescePartitions.enabled", "true")

        # Ключевые параметры coalescing
        .config("spark.sql.shuffle.partitions", str(shuffle_partitions))
        .config("spark.sql.adaptive.advisoryPartitionSizeInBytes",
                str(advisory_size_bytes))
        .config("spark.sql.adaptive.coalescePartitions.minPartitionNum",
                str(min_partitions))

        # parallelismFirst=false: соблюдаем целевой размер строго
        # Это важно для Lakehouse: хотим файлы нужного размера, не максимальный параллелизм
        .config("spark.sql.adaptive.coalescePartitions.parallelismFirst", "false")

        # Skew Join - включаем вместе с coalescing (работают дополняя друг друга)
        .config("spark.sql.adaptive.skewJoin.enabled", "true")
        .config("spark.sql.adaptive.skewJoin.skewedPartitionFactor", "5")
        .config("spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes", "256m")
    )


# Пример использования для кластера: 20 Executor'ов, 16 ГБ каждый, пик 2 ТБ данных
builder = SparkSession.builder.appName("Production-ETL")
configured_builder = production_aqe_config(
    spark_builder=builder,
    peak_data_gb=2048,      # 2 ТБ пиковый объём после shuffle
    executor_memory_gb=16,  # 16 ГБ на Executor
    num_executors=20,       # 20 Executor'ов
    target_file_size_mb=128 # Хотим Parquet-файлы по 128 МБ
)

# Расчётные значения для данного примера:
# shuffle_partitions = max(200, 2048 * 1024 / 64 * 4) = max(200, 131072) = 131072
# advisoryPartitionSizeInBytes = 128 МБ
# minPartitionNum = 20 × 4 = 80 ядер
# Итог: AQE может схлопнуть 131072 исходных партиций до ~(2TB/128MB) = ~16384 тасок

Типичные антипаттерны и как их избегать

Антипаттерн 1: Отключать AQE ради «предсказуемости»

# НЕПРАВИЛЬНО: отключение AQE для детерминированного числа файлов
spark.conf.set("spark.sql.adaptive.enabled", "false")
spark.conf.set("spark.sql.shuffle.partitions", "500")
df.write.parquet("s3://bucket/output/")
# Результат: 500 файлов независимо от объёма данных.
# При 5 МБ данных - 500 файлов по 10 КБ (катастрофа для S3 listing).
# При 5 ТБ данных - 500 файлов по 10 ГБ (Parquet reader будет страдать).

# ПРАВИЛЬНО: доверяйте AQE + задайте advisoryPartitionSizeInBytes
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.shuffle.partitions", "10000")
spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "128m")
spark.conf.set("spark.sql.adaptive.coalescePartitions.parallelismFirst", "false")
df.write.parquet("s3://bucket/output/")
# Результат: ~(total_data_size / 128MB) файлов.
# При 5 МБ → 1 файл. При 5 ТБ → ~40000 файлов.

Антипаттерн 2: Ставить минимальный shuffle_partitions при включённом AQE

# НЕПРАВИЛЬНО: слишком мало исходных партиций
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.shuffle.partitions", "50")  # Мало!
# Если данных окажется 10 ТБ → 50 партиций по 200 ГБ → огромные таски, spill, OOM.
# AQE может только СХЛОПЫВАТЬ, но не РАСШИРЯТЬ число партиций выше заданного значения.

# ПРАВИЛЬНО: задавать с запасом на пиковый объём
spark.conf.set("spark.sql.shuffle.partitions", "8000")  # Для пика 2 ТБ
# При 10 МБ данных AQE схлопнет до 1-4 партиций - всё нормально.
# При 2 ТБ данных AQE оставит ~16000 партиций по 128 МБ - всё нормально.

Антипаттерн 3: Игнорировать parallelismFirst

# Ситуация: хотим строго контролировать размер выходных файлов в Iceberg
spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "128m")

# Но если оставить parallelismFirst=true (дефолт), Spark может создать
# больше партиций, чем нужно, чтобы «занять» все ядра кластера.
# В итоге файлы будут мельче 128 МБ - цель не достигнута.

# ПРАВИЛЬНО для write-heavy пайплайнов с Lakehouse:
spark.conf.set("spark.sql.adaptive.coalescePartitions.parallelismFirst", "false")
# Теперь Spark строго соблюдает advisoryPartitionSizeInBytes,
# даже если часть ядер кластера будет простаивать.

Итоги урока

AQE Adaptive Coalescing решает фундаментальное противоречие статического планирования: невозможность задать одно число shuffle-партиций, которое будет оптимальным для всех объёмов данных.

Вместо того чтобы пытаться угадать правильное значение до запуска, AQE переносит решение в момент, когда информация уже известна: после завершения Map Stage, когда каждая исходная партиция имеет точно измеренный размер. Алгоритм жадного объединения соседних партиций прост, но эффективен: он превращает тысячи мелких пустых тасок в десятки полноценных, устраняя основной источник scheduling overhead в over-partitioned запросах.

На практике рецепт прост:

  • Задайте spark.sql.shuffle.partitions с запасом под пиковый объём данных
  • Установите advisoryPartitionSizeInBytes равным целевому размеру файла в Lakehouse (128 МБ - хороший старт)
  • Доверяйте AQE: он знает о реальных данных больше, чем любая статическая конфигурация

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

Вам предоставлен следующий датасет и пайплайн:

# Датасет: транзакции e-commerce за квартал
# Объём: ~1.2 ТБ (после shuffle - около 800 ГБ после агрегации)
# Кластер: 40 Executor'ов, 8 ядер каждый, 32 ГБ RAM

raw_txns = spark.read.parquet("s3://data-lake/bronze/transactions/2024-Q4/")

# Ваш Gold-layer пайплайн:
gold = (
    raw_txns
    .groupBy("merchant_id", "product_category", "txn_date")
    .agg(
        F.sum("amount").alias("daily_gmv"),
        F.count("*").alias("txn_count"),
        F.approx_count_distinct("user_id").alias("unique_buyers"),
    )
    .groupBy("merchant_id", "product_category")
    .agg(
        F.sum("daily_gmv").alias("quarterly_gmv"),
        F.sum("txn_count").alias("quarterly_txns"),
        F.avg("unique_buyers").alias("avg_daily_buyers"),
    )
)

gold.write.format("parquet").mode("overwrite").save("s3://data-lake/gold/merchant_summary/")

Задание:

  1. Математически рассчитайте оптимальный advisoryPartitionSizeInBytes для этого пайплайна, если целевой размер Parquet-файла в Gold-витрине - 256 МБ, а объём данных после второй агрегации ожидается ~50 ГБ.

  2. Рассчитайте spark.sql.shuffle.partitions с достаточным запасом.

  3. Запустите пайплайн с рассчитанными параметрами и AQE включённым.

  4. Зафиксируйте в Spark UI вкладку Stages для Reduce Stage второй агрегации:

  5. Фактическое число тасок после coalescing
  6. Scheduler Delay (должен быть < 5% от Task Duration)
  7. Shuffle Read per task (должен быть близок к 256 МБ)

  8. Проверьте в финальном explain("formatted") наличие оператора AQEShuffleRead coalesced и убедитесь, что число итоговых партиций соответствует формуле (50 GB / 256 MB).

Ожидаемый результат: в выходной директории должно быть около 200 файлов по ~256 МБ, а в Spark UI - 0 тасок с Task Duration < 100 мс.