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