AQE Skew Join: автоматическое обнаружение и разбиение hot partitions

Глубокий разбор AQE Skew Join в Spark 3.x: анатомия Data Skew и эффект Straggler Task, физика Shuffle при джойне, исторические паттерны salting и ручной фильтрации, алгоритм детекции горячих партиций, механика разбиения, конфигурационные триггеры, ограничения и мониторинг через Spark UI и .explain().

optimization

1. Анатомия Data Skew: почему одна таска парализует весь кластер

Data Skew - это неравномерное распределение данных по партициям в Shuffle-операциях. Это не исключение в production Data Lake, это норма. Понять почему - значит понять физику распределённых вычислений.

Фундаментальный закон Spark Stages

Spark выполняет запрос в виде Stages - групп задач, разделённых Shuffle-барьерами. Критически важное правило: Stage считается завершённым только когда завершена ПОСЛЕДНЯЯ Task. Все остальные Task могут выполниться мгновенно - Stage всё равно будет ждать.

Схема наглядно показывает катастрофу: без скоса 200 тасок выполняются параллельно за 2.6 минуты. Со скосом 199 тасок завершаются за 3 минуты, но Stage продолжает работать ещё 28 часов - ждёт одну задачу. Все 199 Executor'ов в это время потребляют память и CPU кластера, ничего не производя.

Физические причины Data Skew в Data Lake

Бизнес-логика порождает неравномерность. В B2B-системах крупные клиенты (Apple, Amazon, Сбербанк) могут иметь в сотни раз больше транзакций, чем средний клиент. Клиент customer_id=1 может иметь 500 миллионов заказов, пока остальные имеют по 1000.

Технические дефолтные значения. В логах событий «незалогиненные» пользователи часто имеют user_id = NULL, user_id = 0 или user_id = -1. Если 80% трафика - анонимные пользователи, эти ключи будут катастрофически перегружены.

Временные паттерны. Данные за «пиковый» день (например, Чёрная пятница) или «пиковый» час создают огромные партиции при GROUP BY date/hour.

Иерархическая вложенность. В телеком- и IoT-данных один устройство может генерировать миллионы событий, пока тысячи других - по десятку.

Физика Shuffle при JOIN: как данные попадают в один Executor

Чтобы понять почему skew - это проблема Shuffle, нужно разобраться как Spark распределяет данные при JOIN:

Схема объясняет фундаментальную причину: хэш-функция детерминирована. Все строки с customer_id=NULL всегда дают одинаковый хэш → всегда попадают в одну партицию → всегда обрабатываются одним Executor'ом. Это не баг хэш-функции - это её свойство. И именно это свойство создаёт проблему при неравномерных данных.

Эффект Straggler Task: реальная стоимость для бизнеса

Straggler Task (отстающая задача) - Task, которая выполняется значительно дольше медианного времени в Stage. В условиях Data Skew Straggler Task почти всегда вызвана перегруженной Shuffle Partition.

Реальные последствия:

  • Airflow DAG задерживается → нарушение SLA по доставке отчётов
  • YARN/K8s выделяет ресурсы на весь период ожидания → прямые финансовые потери
  • При достаточно большом skew Executor падает с OOM → пересчёт всего Stage → многократное увеличение времени
  • В multi-tenant кластере другие jobs ждут освобождения ресурсов

2. Как инженеры боролись со Skew до эпохи AQE

До появления AQE Skew Join в Spark 3.x инженеры использовали несколько обходных путей. Все они требовали ручной работы, знания данных заранее и создавали собственные проблемы.

Паттерн 1: ручная фильтрация горячих ключей

Самый прямолинейный подход - обработать «горячие» ключи и «нормальные» ключи отдельно:

from pyspark.sql import SparkSession, functions as F

spark = SparkSession.builder.getOrCreate()

orders = spark.table("bronze.orders")
customers = spark.table("silver.customers")

# ── Шаг 1: Разделяем данные на две группы ────────────────────────────
# «Горячие» ключи: NULL и дефолтные значения
hot_keys = [None, 0, -1, "UNKNOWN"]

orders_hot = orders.filter(F.col("customer_id").isNull() |
                           F.col("customer_id").isin(0, -1))
orders_normal = orders.filter(F.col("customer_id").isNotNull() &
                              ~F.col("customer_id").isin(0, -1))

# ── Шаг 2: Разные стратегии JOIN ─────────────────────────────────────
# Для «горячих»: broadcast маленькую таблицу customers (она небольшая)
# Даже если orders_hot огромный, customers маленький → BHJ без shuffle
result_hot = orders_hot.join(
    F.broadcast(customers),
    "customer_id",
    "left"
)

# Для «нормальных»: обычный SortMergeJoin (данные равномерны)
result_normal = orders_normal.join(customers, "customer_id", "left")

# ── Шаг 3: Объединяем результаты ─────────────────────────────────────
final_result = result_hot.union(result_normal)

Проблемы этого подхода:

  • Нужно знать список «горячих» ключей заранее. Что если их характер меняется со временем?
  • Два отдельных scan исходных данных - двойная нагрузка на HDFS/S3
  • Код становится сложным, хрупким, требует поддержки
  • UNION ALL добавляет лишний Shuffle Stage

Паттерн 2: Key Salting (сальтирование ключей)

Идея: «размыть» горячий ключ по нескольким партициям, добавив случайный суффикс:

import random

SALT_FACTOR = 100  # число «копий» горячего ключа

# ── Шаг 1: Добавляем случайную «соль» к ключам левой таблицы ──────────
# Каждая строка получает случайный суффикс от 0 до SALT_FACTOR-1
# NULL_0, NULL_1, NULL_2, ..., NULL_99
# customer_123_0, customer_123_1, ..., customer_123_99
orders_salted = orders.withColumn(
    "salt",
    (F.rand() * SALT_FACTOR).cast("int")
).withColumn(
    "customer_id_salted",
    F.concat(
        F.coalesce(F.col("customer_id").cast("string"), F.lit("NULL")),
        F.lit("_"),
        F.col("salt").cast("string")
    )
)
# Теперь 15M NULL строк → 15M строк с ключами NULL_0..NULL_99
# Hash(NULL_0), Hash(NULL_1), ..., Hash(NULL_99) → 100 разных партиций!

# ── Шаг 2: «Размножаем» правую таблицу ───────────────────────────────
# Каждая строка customers должна существовать в SALT_FACTOR копиях
# чтобы JOIN с salted-ключами дал правильный результат
# customers: customer_id=NULL → NULL_0, NULL_1, ..., NULL_99

from pyspark.sql.types import IntegerType

salt_df = spark.range(SALT_FACTOR).select(
    F.col("id").cast(IntegerType()).alias("salt")
)

# CrossJoin customers × [0..99] - получаем 100x больше строк customers
customers_expanded = customers.crossJoin(salt_df).withColumn(
    "customer_id_salted",
    F.concat(
        F.coalesce(F.col("customer_id").cast("string"), F.lit("NULL")),
        F.lit("_"),
        F.col("salt").cast("string")
    )
)

# ── Шаг 3: JOIN по salted ключу ──────────────────────────────────────
result = orders_salted.join(
    customers_expanded,
    "customer_id_salted",
    "left"
)

Проблемы salting:

  • Размер таблицы customers_expanded вырастает в SALT_FACTOR раз → память Executor'ов
  • Нужно знать SALT_FACTOR заранее: слишком маленький → skew сохраняется, слишком большой → OOM на expanded таблице
  • Бизнес-логика перемешивается с инфраструктурным кодом
  • При изменении данных скрипт salting нужно пересматривать

Почему эти костыли неприемлемы в 2026 году

Ключевая проблема обоих паттернов: они требуют знания данных заранее и ручной поддержки. В production Data Lake данные постоянно меняются: появляются новые горячие ключи, старые охлаждаются, объём растёт. Код salting превращается в техдолг, который регулярно ломается при изменениях в данных.

Именно это побудило команду Spark разработать автоматическое решение.


3. Алгоритм работы AQE Skew Join под капотом Catalyst

AQE Skew Join - это полностью автоматический механизм, работающий на уровне Spark Runtime без изменений в пользовательском коде.

Стадия детектирования: MapOutputStatistics

После завершения Map-стадии (Shuffle Write) каждый Executor отправляет Driver'у компактный отчёт MapOutputStatistics. Этот отчёт содержит размер каждой Shuffle Partition в байтах.

Driver собирает эти отчёты и строит полную картину распределения данных:

Важный момент: AQE анализирует суммарный размер каждой партиции по всем Executor'ам. Если Partition 0 получила по 600-800 MB от каждого Executor'а → суммарный размер огромный.

Математика определения перекоса: два условия одновременно

AQE считает партицию скошенной если выполнены оба условия:

Условие 1: partition_size > skewedPartitionFactor × median_partition_size
Условие 2: partition_size > skewedPartitionThresholdInBytes

Почему нужны оба? Рассмотрим:

  • Только условие 1 (relative): если все партиции по 1 KB и одна 10 KB - formально она в 10 раз больше медианы. Но нет смысла разбивать 10 KB партицию.
  • Только условие 2 (absolute): если все партиции по 500 MB и threshold = 256 MB - все считаются скошенными, что нелепо.

Комбинация: «партиция значительно больше медианы И абсолютно большая» - это и есть настоящий skew.

# Пример вычисления для нашего кейса:
# Partition sizes: [2.8, 2.1, 1.9, 2.3, ..., 1460.0] MB
# Sorted: [1.9, 2.1, 2.3, 2.8, ..., 1460.0] MB
# Median = 2.1 MB

skew_factor = 5      # spark.sql.adaptive.skewJoin.skewedPartitionFactor
skew_threshold = 256 # MB, spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes

median = 2.1
p0_size = 1460.0

condition_1 = p0_size > skew_factor * median  # 1460 > 5 * 2.1 = 10.5 → True
condition_2 = p0_size > skew_threshold         # 1460 > 256 → True
is_skewed = condition_1 and condition_2         # True!

# Сколько частей нужно создать?
advisory_size = 64   # MB, spark.sql.adaptive.advisoryPartitionSizeInBytes
n_splits = max(2, int(p0_size / advisory_size))  # max(2, int(1460/64)) = 22 части!

Механика разбиения: Partition Splitting

После детекции скошенной партиции AQE выполняет физическое разбиение Shuffle данных:

Координация противоположной стороны JOIN

Это самая сложная часть алгоритма. Когда левая сторона JOIN разбита на 22 части, каждая из них должна джойниться с соответствующими строками правой стороны.

AQE решает эту задачу через дублирование: для каждого Split левой таблицы создаётся копия соответствующей партиции правой таблицы.

Правая партиция (2 MB) дублируется 22 раза - это дешевле, чем иметь одну задачу с 1460 MB. Каждая Task работает параллельно и независимо.

Корректность результата гарантирована: каждый Split левой таблицы получает все строки правой таблицы для соответствующего ключа (NULL). Union всех результатов даёт тот же результат, что и полный JOIN.


4. Конфигурационные триггеры и параметры детекции Skew

Главный выключатель и иерархия зависимостей

from pyspark.sql import SparkSession

spark = SparkSession.builder \

    # Главный тумблер AQE - должен быть включён для Skew Join
    # По умолчанию true в Spark 3.2+
    .config("spark.sql.adaptive.enabled", "true") \

    # Тумблер Skew Join - по умолчанию true при enabled AQE
    .config("spark.sql.adaptive.skewJoin.enabled", "true") \

    # ── Порог кратности (относительный) ──────────────────────────────
    # Партиция скошена если: partition_size > factor × median_size
    # Дефолт: 5 (партиция в 5x больше медианы → skew)
    # Снизьте до 3 для агрессивного обнаружения
    # Поднимите до 10 для консервативного (меньше false positive)
    .config("spark.sql.adaptive.skewJoin.skewedPartitionFactor", "5") \

    # ── Минимальный абсолютный размер для skew ───────────────────────
    # Партиция скошена если: partition_size > threshold
    # Дефолт: 256 MB. Защита от false positive на маленьких данных.
    # Снизьте до 64 MB если работаете с датасетами < 10 GB
    # и хотите увидеть AQE в действии на примерах
    .config("spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes", "256MB") \

    # ── Целевой размер Split ─────────────────────────────────────────
    # На сколько нарезать скошенную партицию
    # Дефолт: 64 MB. Связан с advisoryPartitionSizeInBytes (Coalesce)
    # Рекомендация: совпадает с размером блока HDFS (128 MB для HDFS)
    .config("spark.sql.adaptive.advisoryPartitionSizeInBytes", "64MB") \

    .getOrCreate()

Расчёт числа сплитов

Число сплитов вычисляется автоматически:

# Внутренняя логика Spark (псевдокод)
def calculate_splits(skewed_partition_size_bytes: int,
                     advisory_size_bytes: int) -> int:
    """
    Вычисляет оптимальное число частей для скошенной партиции.
    Минимум 2 части (иначе нет смысла разбивать).
    """
    n_splits = max(2, math.ceil(skewed_partition_size_bytes / advisory_size_bytes))
    return n_splits

# Примеры:
# Partition 1460 MB, advisory = 64 MB → ceil(1460/64) = 23 splits
# Partition 500 MB,  advisory = 64 MB → ceil(500/64)  = 8 splits
# Partition 300 MB,  advisory = 64 MB → ceil(300/64)  = 5 splits
# Partition 128 MB,  advisory = 64 MB → max(2, 2)     = 2 splits

# Итоговое время выполнения:
# 1460 MB / 64 MB/split ≈ 23 задачи × (64 MB / disk_throughput) = параллельно!

Таблица: когда изменять дефолтные параметры

Сценарий Рекомендация Почему
Датасеты < 10 GB skewedPartitionThresholdInBytes = 32MB Дефолт 256 MB слишком большой
Экстремальный skew (1 ключ = 90%+ данных) skewedPartitionFactor = 3 Aggressive detection
Нагруженный кластер (много concurrent jobs) advisoryPartitionSizeInBytes = 128MB Меньше задач = меньше Scheduler Overhead
Нестабильные данные с частым изменением распределения skewedPartitionFactor = 3 Раньше обнаруживает новые hot keys
Маленький кластер (< 10 Executors) advisoryPartitionSizeInBytes = 256MB Слишком много мелких tasks невыгодны

5. Ограничения и слепые зоны AQE Skew Join

AQE Skew Join - мощный инструмент, но не универсальное решение. Понимание его ограничений критически важно для корректного применения.

Ограничение 1: только SortMergeJoin

AQE Skew Join работает только тогда, когда Spark выбирает SortMergeJoin (SMJ). BroadcastHashJoin не имеет Shuffle → нет MapOutputStatistics → AQE не может анализировать распределение данных.

# Если таблица попала под broadcast threshold → BHJ → нет skew optimization
# В этом случае skew уже не проблема (broadcast = весь датасет в памяти каждого Executor)

# Проблема возникает только при SMJ:
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "10MB")
# Если правая таблица > 10 MB → SMJ → AQE Skew может помочь

# Если правая таблица < 10 MB → BHJ → никакого skew, никакой оптимизации не нужно

Ограничение 2: Full Outer Join

Full Outer Join несовместим с AQE Skew Join из-за сложности генерации NULL-строк при отсутствии совпадений. При полном внешнем соединении AQE Skew не применяется.

# Проблема: Full Outer Join
result = left.join(right, "id", "full")
# AQE Skew Join НЕ применяется

# Решение: разбить на два join'а
inner_result = left.join(right, "id", "inner")
left_only = left.join(right, "id", "left_anti")  # строки только из left
right_only = right.join(left, "id", "right_anti")  # строки только из right
# Затем union: inner_result + left_only + right_only = эквивалент full outer

Ограничение 3: экстремальный skew с OOM

Если скошенная партиция настолько велика, что даже один Split превышает доступную память Executor'а, Spark спишет данные на диск (spill). Если spill tоже невозможен - OOM.

Пример: партиция 500 GB, advisory = 64 MB → 7813 splits. Каждый split 64 MB - теоретически управляемо. Но если Executor имеет только 32 GB RAM и пытается загрузить в memory-intensive join 64 MB данных + структуры JOIN + overhead → спилл.

# Признаки что AQE Skew недостаточен:
# 1. В Spark UI: Executor #N имеет значительный "Spill (Memory)" и "Spill (Disk)"
# 2. Логи: "Container killed by YARN for exceeding memory limits"
# 3. Stage выполняется медленно несмотря на skew optimization

# В этом случае нужна дополнительная ручная оптимизация:
# a) Увеличить Executor memory: spark.executor.memory = 16g
# b) Уменьшить advisoryPartitionSizeInBytes: 32MB вместо 64MB
# c) Применить salting для экстремальных ключей (NULL, дефолтные)
# d) Фильтровать NULL-ключи до JOIN

Ограничение 4: skew до Shuffle Stage

AQE Skew Join работает только в Shuffle Phase. Если данные уже неравномерны в исходных файлах (например, все файлы одной даты партиции находятся на одном DataNode), AQE не может это исправить.

Исходные файлы (HDFS Input Stage):
- date=2024-01-15/part-000.parquet: 200 MB  → Task 1: 2 min
- date=2024-01-15/part-001.parquet: 200 MB  → Task 2: 2 min
...
- date=2024-01-15/part-999.parquet: 200 MB  → Task 1000: 2 min

Это не skew в Shuffle, это просто большой датасет.
AQE Coalesce поможет (если файлы маленькие), Skew Join - нет.

6. Мониторинг и чтение планов выполнения

.explain(): ключевые маркеры AQE Skew Join

from pyspark.sql import SparkSession, functions as F

spark = SparkSession.builder \
    .master("local[8]") \
    .config("spark.sql.adaptive.enabled", "true") \
    .config("spark.sql.adaptive.skewJoin.enabled", "true") \
    .config("spark.sql.adaptive.skewJoin.skewedPartitionFactor", "3") \
    .config("spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes", "10MB") \
    .config("spark.sql.adaptive.advisoryPartitionSizeInBytes", "10MB") \
    .getOrCreate()

# Создаём скошенный датасет
orders = spark.range(1_000_000).select(
    F.col("id").alias("order_id"),
    # 90% заказов - user_id=0 (незалогиненные)
    F.when(F.rand() < 0.9, F.lit(0))
     .otherwise((F.rand() * 99 + 1).cast("int"))
     .alias("user_id"),
    (F.rand() * 1000).alias("amount")
)

users = spark.range(100).select(
    F.col("id").alias("user_id"),
    F.concat(F.lit("User_"), F.col("id")).alias("username")
)

result = orders.join(users, "user_id")

# Смотрим план ДО выполнения
print("=== ПЛАН ДО ВЫПОЛНЕНИЯ ===")
result.explain("formatted")
# Увидим: AdaptiveSparkPlan isFinalPlan=false
# Current plan: SortMergeJoin (AQE не знает о skew заранее)

# Выполняем
result.collect()

# Смотрим план ПОСЛЕ выполнения
print("=== ПЛАН ПОСЛЕ ВЫПОЛНЕНИЯ ===")
result.explain("formatted")

Пример вывода после выполнения (фрагмент):

== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=true
+- == Final Plan ==
   Project [order_id#0, user_id#1, amount#2, username#12]
   +- SortMergeJoin [user_id#1], [user_id#7], Inner
      :- Sort [user_id#1 ASC NULLS FIRST], false, 0
      :  +- AQEShuffleRead coalesced                      ← AQE обработал shuffle read
      :     +- ShuffleQueryStage 0
      :        +- Exchange hashpartitioning(user_id#1, 200)
      :           +- Range (0, 1000000, step=1, splits=8)
      +- Sort [user_id#7 ASC NULLS FIRST], false, 0
         +- AQEShuffleRead skewed                         ← skewed партиции обнаружены!
            +- ShuffleQueryStage 1
               +- Exchange hashpartitioning(user_id#7, 200)
                  +- Range (0, 100, step=1, splits=8)

== Initial Plan ==          ← До AQE: без skew обработки
   SortMergeJoin ...

Ключевые строки:

  • AQEShuffleRead skewed - подтверждение применения Skew Join оптимизации
  • isFinalPlan=true - финальный адаптированный план
  • Если бы не было skew - написало бы AQEShuffleRead coalesced (только coalesce)

Spark UI: что искать для верификации Skew Join

В Spark UI вкладка SQL / DataFrame показывает граф выполнения. При AQE Skew Join ищите:

В узле SortMergeJoin (кликаем на него):

  • number of skewed partitions: N - сколько партиций были признаны скошенными
  • number of skewed partition splits: M - на сколько частей они были разбиты

В вкладке Stages для Stage с JOIN:

  • Общее число Tasks = обычные_tasks + skew_split_tasks (больше чем ожидаете)
  • Равномерное распределение времени Task'ов в Task Distribution (нет длинного «хвоста»)
  • Все Task'и завершились примерно одновременно

В вкладке Stages → конкретный Stage → Tasks:

  • Ищите Tasks с суффиксами в ID: task_0001, task_0001_1, task_0001_2 - split tasks от скошенной партиции

Программный мониторинг через SparkListener

from pyspark import SparkContext
from pyspark.listener import SparkListener

class AQESkewMonitor:
    """
    Собирает метрики о применении AQE Skew Join.
    Полезно для production мониторинга и алертинга.
    """

    def __init__(self, spark: SparkSession):
        self.spark = spark
        self.skewed_stages = []

    def check_last_job_for_skew(self) -> dict:
        """
        Проверяет последний завершённый job на наличие skew оптимизации.
        Использует SparkContext StatusTracker.
        """
        sc = self.spark.sparkContext
        completed_jobs = sc.statusTracker().getJobIdsForGroup(None)

        skew_info = {
            "skew_detected": False,
            "stages_with_skew": [],
            "total_skewed_partitions": 0,
        }

        # В production используйте SparkListener для real-time мониторинга
        # Здесь упрощённая версия через SQL metrics

        # Получаем метрики из последнего SQL плана
        last_sql = self.spark.sql("""
            SELECT
                executionId,
                description,
                metrics
            FROM spark_catalog.information_schema.sql_operations
            LIMIT 1
        """)
        # (Реальный доступ через REST API History Server)

        return skew_info

7. Практика: симуляция Data Skew и сравнение результатов

Создание реалистичного скошенного датасета

from pyspark.sql import SparkSession, functions as F
from pyspark.sql.types import StructType, StructField, LongType, StringType, DoubleType
import time

spark = SparkSession.builder \
    .master("local[8]") \
    .appName("skew-join-benchmark") \
    .getOrCreate()

spark.sparkContext.setLogLevel("ERROR")


def create_skewed_orders(spark: SparkSession, total_rows: int = 2_000_000):
    """
    Создаёт реалистичный датасет заказов с Data Skew.

    Распределение user_id:
    - user_id=0 (гость/незалогиненный): 80% строк → горячий ключ
    - user_id=1 (крупный B2B клиент):  10% строк → второй горячий ключ
    - user_id=2..99:                   10% строк → нормальные клиенты

    Это имитирует реальный e-commerce: большинство трафика анонимные пользователи
    и один или несколько крупных корпоративных клиентов.
    """
    n_guest = int(total_rows * 0.80)   # 80% гости
    n_b2b = int(total_rows * 0.10)     # 10% крупный B2B
    n_normal = total_rows - n_guest - n_b2b  # 10% нормальные

    # Гостевые заказы: user_id=0
    guest_orders = spark.range(n_guest).select(
        F.col("id").alias("order_id"),
        F.lit(0).alias("user_id"),
        (F.rand() * 1000).alias("amount"),
        F.lit("guest").alias("order_type")
    )

    # B2B заказы: user_id=1
    b2b_orders = spark.range(n_b2b).select(
        (F.col("id") + n_guest).alias("order_id"),
        F.lit(1).alias("user_id"),
        (F.rand() * 50000).alias("amount"),
        F.lit("b2b").alias("order_type")
    )

    # Нормальные заказы: user_id=2..99
    normal_orders = spark.range(n_normal).select(
        (F.col("id") + n_guest + n_b2b).alias("order_id"),
        (F.col("id") % 98 + 2).cast("long").alias("user_id"),
        (F.rand() * 500).alias("amount"),
        F.lit("retail").alias("order_type")
    )

    return guest_orders.union(b2b_orders).union(normal_orders)


def create_users_table(spark: SparkSession):
    """Маленький справочник пользователей."""
    return spark.range(100).select(
        F.col("id").alias("user_id"),
        F.concat(F.lit("User_"), F.col("id")).alias("username"),
        F.array(F.lit("Bronze"), F.lit("Silver"), F.lit("Gold")).getItem(
            (F.col("id") % 3).cast("int")
        ).alias("tier")
    )


def run_join_benchmark(
    spark: SparkSession,
    orders,
    users,
    use_aqe: bool,
    skew_threshold_mb: int = 256,
    skew_factor: int = 5,
    label: str = ""
) -> tuple[float, int]:
    """
    Запускает JOIN и возвращает (время в секундах, число задач).
    """
    if use_aqe:
        spark.conf.set("spark.sql.adaptive.enabled", "true")
        spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
        spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionFactor",
                       str(skew_factor))
        spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes",
                       f"{skew_threshold_mb}MB")
        spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "16MB")
    else:
        spark.conf.set("spark.sql.adaptive.enabled", "false")
        spark.conf.set("spark.sql.shuffle.partitions", "200")

    result = orders.join(users, "user_id") \
        .groupBy("username", "tier") \
        .agg(
            F.sum("amount").alias("total_amount"),
            F.count("*").alias("order_count"),
            F.avg("amount").alias("avg_amount")
        )

    print(f"\n{'='*50}")
    print(f"Запуск: {label}")
    print(f"AQE: {'ВКЛЮЧЁН' if use_aqe else 'ВЫКЛЮЧЕН'}")

    t0 = time.time()
    collected = result.collect()
    elapsed = time.time() - t0

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

    return elapsed, len(collected)


# ── Основной бенчмарк ─────────────────────────────────────────────────────
if __name__ == "__main__":
    print("Создаём датасеты...")
    orders = create_skewed_orders(spark, total_rows=2_000_000)
    users = create_users_table(spark)

    # Кешируем исходные данные
    orders.cache()
    users.cache()
    orders.count()
    users.count()

    # Показываем распределение (чтобы видеть skew)
    print("\nРаспределение user_id в orders:")
    orders.groupBy("user_id") \
        .count() \
        .orderBy(F.desc("count")) \
        .show(10)

    # ── Тест 1: без AQE (имитация Spark 2.x или выключенного AQE) ────────
    t_no_aqe, rows_no_aqe = run_join_benchmark(
        spark, orders, users,
        use_aqe=False,
        label="Без AQE (статический Catalyst)"
    )

    # ── Тест 2: с AQE Skew Join (настройки под наш датасет) ──────────────
    t_aqe, rows_aqe = run_join_benchmark(
        spark, orders, users,
        use_aqe=True,
        skew_threshold_mb=10,  # Уменьшен для демонстрации (датасет ~200 MB)
        skew_factor=3,          # Агрессивная детекция
        label="С AQE Skew Join"
    )

    # ── Итоговое сравнение ────────────────────────────────────────────────
    print("\n" + "="*60)
    print("ИТОГОВОЕ СРАВНЕНИЕ")
    print("="*60)
    print(f"Без AQE:   {t_no_aqe:.1f}с  ({rows_no_aqe} строк в результате)")
    print(f"С AQE:     {t_aqe:.1f}с  ({rows_aqe} строк в результате)")
    if t_no_aqe > 0:
        print(f"Ускорение: {t_no_aqe/t_aqe:.2f}x")
    print()
    print("Проверьте Spark UI (http://localhost:4040):")
    print("  → Вкладка SQL: найдите SortMergeJoin")
    print("  → Кликните на него: ищите 'skewed partitions'")
    print("  → Вкладка Stages: сравните Task Distribution")

    spark.stop()

Верификация: сравниваем планы вручную

from pyspark.sql import SparkSession, functions as F

spark = SparkSession.builder \
    .master("local[4]") \
    .config("spark.sql.adaptive.skewJoin.skewedPartitionFactor", "3") \
    .config("spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes", "5MB") \
    .config("spark.sql.adaptive.advisoryPartitionSizeInBytes", "5MB") \
    .getOrCreate()

orders = spark.range(500_000).select(
    F.col("id"),
    F.when(F.rand() < 0.85, F.lit(0))
     .otherwise((F.rand() * 10 + 1).cast("int"))
     .alias("user_id"),
    F.rand().alias("amount")
)

users = spark.range(11).select(
    F.col("id").alias("user_id"),
    F.concat(F.lit("U"), F.col("id")).alias("name")
)

result = orders.join(users, "user_id").groupBy("name").sum("amount")

# ── Без AQE ──────────────────────────────────────────────────────────────
spark.conf.set("spark.sql.adaptive.enabled", "false")
print("=== БЕЗ AQE: статический план ===")
result.explain("formatted")

# ── С AQE: до выполнения ─────────────────────────────────────────────────
spark.conf.set("spark.sql.adaptive.enabled", "true")
print("\n=== С AQE: план ДО выполнения ===")
result.explain("formatted")
# isFinalPlan=false - план ещё не финальный

# ── С AQE: после выполнения ──────────────────────────────────────────────
result.collect()  # выполняем
print("\n=== С AQE: план ПОСЛЕ выполнения ===")
result.explain("formatted")
# isFinalPlan=true
# Ищем: AQEShuffleRead skewed
# Ищем: SkewJoin (если Spark вывел этот оператор)

Диагностический скрипт: анализ распределения без выполнения JOIN

Перед запуском тяжёлого JOIN полезно проверить распределение данных:

def diagnose_skew_potential(
    spark: SparkSession,
    df,
    join_key: str,
    n_shuffle_partitions: int = 200,
    top_n: int = 10
) -> dict:
    """
    Анализирует потенциальный skew ДО выполнения JOIN.
    Помогает понять нужен ли AQE и насколько агрессивными
    должны быть параметры детекции.

    Возвращает рекомендуемые параметры AQE.
    """
    print(f"Анализ распределения '{join_key}'...")

    # Подсчёт частоты каждого значения ключа
    key_counts = df.groupBy(join_key).count() \
        .orderBy(F.desc("count"))

    total_rows = df.count()
    top_keys = key_counts.limit(top_n).collect()

    print(f"\nТоп-{top_n} значений '{join_key}':")
    for row in top_keys:
        pct = row["count"] / total_rows * 100
        print(f"  {str(row[join_key]):20s}: {row['count']:>10,} строк ({pct:.1f}%)")

    # Вычисляем медиану по хэш-партициям (приближённо)
    # Реальная медиана доступна только после Shuffle, но можно оценить
    key_counts_df = key_counts.withColumn(
        "partition",
        (F.hash(F.col(join_key)) % n_shuffle_partitions).cast("int")
    )

    partition_sizes = key_counts_df.groupBy("partition").sum("count")
    stats = partition_sizes.agg(
        F.percentile_approx("sum(count)", 0.5).alias("median"),
        F.max("sum(count)").alias("max"),
        F.min("sum(count)").alias("min"),
        F.avg("sum(count)").alias("avg"),
    ).first()

    print(f"\nСтатистика по {n_shuffle_partitions} партициям (оценка):")
    print(f"  Медиана партиции: {stats.median:,} строк")
    print(f"  Максимум:         {stats.max:,} строк")
    print(f"  Минимум:          {stats.min:,} строк")
    print(f"  Отношение Max/Median: {stats.max/stats.median:.1f}x")

    # Рекомендации
    skew_ratio = stats.max / stats.median
    print(f"\nРекомендации AQE:")
    if skew_ratio > 20:
        print(f"  ⚠️  КРИТИЧЕСКИЙ SKEW ({skew_ratio:.0f}x)!")
        print(f"  → skewedPartitionFactor = 3")
        print(f"  → Рассмотрите ручное salting для top-1 ключей")
    elif skew_ratio > 5:
        print(f"  ⚠️  ЗНАЧИТЕЛЬНЫЙ SKEW ({skew_ratio:.0f}x)")
        print(f"  → skewedPartitionFactor = 5 (дефолт)")
        print(f"  → AQE Skew Join должен справиться")
    else:
        print(f"  ✅ Skew умеренный ({skew_ratio:.1f}x)")
        print(f"  → AQE Skew Join с дефолтными параметрами")

    return {
        "skew_ratio": skew_ratio,
        "recommended_factor": 3 if skew_ratio > 20 else 5,
        "hot_keys": [row[join_key] for row in top_keys[:3]],
    }


# Использование перед запуском JOIN:
orders = spark.range(1_000_000).select(
    F.col("id"),
    F.when(F.rand() < 0.85, F.lit(0))
     .otherwise((F.rand() * 99 + 1).cast("int"))
     .alias("user_id"),
    F.rand().alias("amount")
)

diagnosis = diagnose_skew_potential(spark, orders, "user_id")
print(f"\nВыявленные горячие ключи: {diagnosis['hot_keys']}")

Когда AQE Skew Join достаточен, а когда нужны дополнительные меры

Ключевой вывод: AQE Skew Join - это не магия, а инструментальная оптимизация с конкретными условиями применимости. Он отлично справляется с умеренным skew в Inner/Outer JOIN на SortMergeJoin без изменений в пользовательском коде. Для экстремальных случаев или специальных типов JOIN необходимы дополнительные ручные техники.

AQE Skew Join не устраняет причину data skew - он устраняет её последствия. Данные остаются неравномерными, но Spark умеет с этим работать. Если бизнес-логика позволяет - правильнее нормализовать данные или изменить ключ партиционирования. AQE - это страховка, а не замена правильному проектированию данных.