Skew Detection: как увидеть skew в Task Duration и shuffle read
Анатомия data skew: горячие ключи, straggler tasks, shuffle read spill. Как читать Spark UI, писать диагностические запросы и готовить данные к лечению.
Анатомия data skew: почему один task топит весь stage¶
Spark - по природе распределённая система. Её производительность напрямую зависит от равномерного распределения работы между исполнителями. Когда распределение нарушается, возникает data skew - перекос данных.
Представьте конвейер из 200 tasks в одном stage. 199 задач обрабатывают по 10 MB каждая и завершаются за 8 секунд. Одна задача получает 4 GB и работает 40 минут. Весь stage не завершится, пока не закончится последняя задача. 199 исполнителей будут простаивать 40 минут, ожидая единственного straggler.
Это и есть проклятие «99%»: джоб добегает до почти полного завершения, а потом замирает. Пользователь видит Stage 5: 199/200 tasks completed и бесконечно ждёт.
Почему Spark не может сам перераспределить работу¶
В классическом Spark (без AQE) задачи назначаются исполнителям в момент запуска stage, после shuffle. Если shuffle записал 4 GB в один partition key, executor, обрабатывающий этот partition, обречён. Другие executors не могут "помочь" - каждый partition закреплён за одним task.
AQE (Adaptive Query Execution) частично решает эту проблему автоматически. Но для этого сначала нужно понять как skew проявляется, как его обнаружить и в каких случаях AQE недостаточно.
Виды skew в Spark¶
Skew возникает в разных точках пайплайна и имеет разные причины.
Data skew - самый распространённый тип¶
Возникает когда значения ключа join/groupBy/window распределены неравномерно. Несколько значений ("горячие ключи", hot keys) присутствуют в миллионах строк, остальные - в единицах или тысячах.
Примеры горячих ключей:
- Транзакции маркетплейса: Wildberries, Ozon, Яндекс Маркет как продавцы генерируют на порядки больше заказов, чем тысячи мелких продавцов
- Пользовательские события: боты, парсеры, тестовые аккаунты создают аномально много событий
- Sentinel-значения:
NULL,"unknown","N/A",0,""- дефолтные значения, которые встречаются в миллионах строк там, где данных нет - CDC-события: таблица с высокочастотными обновлениями одной строки даёт тысячи событий для одного ключа
Partition skew - неравные входные файлы¶
Даже до shuffle партиции могут быть неравными: один S3-файл весит 2 GB, остальные по 100 MB. Spark создаёт один task на каждый файл (если нет auto-splitting), и первый task будет в 20 раз тяжелее.
Explode skew - взрыв из массивов¶
После explode(items) строки с массивом из 1 элемента и строки с массивом из 10 000 элементов оказываются в одной партиции. Partition, содержавший "тяжёлые" строки, после explode становится огромным.
Window skew - giant partitions¶
partitionBy в window functions работает как shuffle по ключу. Если один ключ (например, country = "RU") содержит 80% данных, весь этот раздел обрабатывается одним executor без возможности параллелизации.
Как Spark распределяет данные при shuffle¶
Чтобы понять skew, нужно понять механизм shuffle.
После операции groupBy("key") или join(..., on="key"):
- Map side (shuffle write): каждый executor вычисляет для каждой строки номер partition назначения:
partition_id = hash(key) % numPartitions. Результаты записываются в shuffle files на локальный диск - Shuffle: каждый reduce-executor читает из всех map-executors строки для своих partition_id
- Reduce side (shuffle read): executor получает все строки для своих partition_id и выполняет агрегацию/join
Проблема: функция hash(key) % numPartitions детерминирована. Если ключ "unknown" встречается в 50 миллионах строк, все 50 миллионов строк будут отправлены в один и тот же partition (один и тот же hash("unknown") % 200). Один executor получит 50 млн строк вместо ожидаемых ~250 тысяч (50M / 200).
Spark UI: где искать skew¶
Вкладка Stages: главное место диагностики¶
Откройте Spark UI → Jobs → нужный Job → выберите подозрительный Stage. Или сразу через вкладку Stages и найдите стадию с большим временем выполнения.
Ключевая секция на странице Stage - Summary Metrics for Completed Tasks. Это таблица квантилей по всем задачам.
Summary Metrics for Completed Tasks
────────────────────────────────────────────────────────────────────────
Metric | Min | 25th %ile | Median | 75th %ile | Max
────────────────────────────────────────────────────────────────────────
Duration | 1.2 s | 3.4 s | 4.1 s | 5.0 s | 42 min
GC Time | 0.1 s | 0.3 s | 0.4 s | 0.5 s | 8 min
Shuffle Read Size | 1.2 MB | 3.8 MB | 4.5 MB | 5.2 MB | 3.8 GB
Shuffle Read Records| 12K | 38K | 44K | 51K | 38.2M
Input Size | 0 B | 0 B | 0 B | 0 B | 0 B
Spill (Memory) | 0 B | 0 B | 0 B | 0 B | 12.4 GB
Spill (Disk) | 0 B | 0 B | 0 B | 0 B | 2.1 GB
────────────────────────────────────────────────────────────────────────
Этот вывод - классический диагноз data skew. Разберём каждую строку.
Метрика Duration: стопроцентный индикатор¶
Медиана 4.1 секунды, максимум 42 минуты - это соотношение max/median ≈ 615. Ни один нормальный вариационный разброс не даёт такой аномалии. Это однозначно skew.
Практические пороги:
max / median < 3- нормально, незначительный дисбалансmax / median 3–10- умеренный skew, стоит исследоватьmax / median > 10- явный skew, требует решенияmax / median > 100- критический skew, джоб практически не работает
Посмотрите на 75th percentile: 5.0 секунд. Это значит что 75% tasks завершились за 5 секунд, и только несколько задач тянули весь stage. Если бы скью был равномерным, 75th percentile тоже был бы сильно выше медианы.
Метрика Shuffle Read: находим виновника¶
Медиана 4.5 MB, максимум 3.8 GB - один task прочитал в 844 раза больше данных, чем обычный.
Shuffle Read Records: медиана 44K, максимум 38.2M - один task обработал в 868 раз больше строк. Это подтверждает: дело не в размере строк, а в количестве. Горячий ключ встречался в 38 миллионах строк.
Shuffle Read - самая важная метрика для диагностики join/groupBy skew. Именно здесь видно: какой объём данных "притянул" к себе один executor через сеть.
Метрика Spill: executor не справился¶
Spill (Memory): 12.4 GB, Spill (Disk): 2.1 GB - executor, получивший перекошенную партицию, не смог поместить 3.8 GB данных в выделенную Execution Memory. Spark начал сбрасывать сериализованные данные на локальный диск воркера.
Spill to disk - катастрофа для производительности:
- Disk I/O во много раз медленнее RAM
- Данные сериализуются перед записью (CPU overhead)
- При десериализации снова тратится CPU
- Если диск заполнен - task падает с
java.io.IOException: No space left on device
Spill (Memory) > 0 - всегда повод для расследования. Это значит что executor работает "за пределами" памяти.
GC Time: вторичный симптом¶
GC Time max: 8 минут при медиане 0.4 секунды. Executor с перекошенной партицией постоянно аллоцирует объекты и давит на GC. Длинные GC паузы останавливают обработку на всё время сборки мусора - "stop-the-world" паузы по несколько секунд.
Event Timeline: визуальная диагностика¶
На странице Stage нажмите кнопку Event Timeline. Вы увидите горизонтальную диаграмму Ганта: по оси X - время, по оси Y - executors. Каждая задача - полоска.
При skew картина выглядит так: большинство полосок коротки и кластеризованы в начале оси времени. Одна или несколько полосок тянутся далеко вправо - это straggler tasks. Остальные executors к этому моменту уже простаивают.
Вкладка SQL и DAG-граф¶
На вкладке SQL Spark UI показывает physical plan с metrics для каждого оператора. Это позволяет точно локализовать узкое место.
Как найти проблемный оператор¶
- Откройте Spark UI → SQL
- Найдите SQL-запрос с большим временем выполнения
- Нажмите на запрос - откроется DAG-граф физического плана
- Ищите узлы
Exchange(shuffle) - они разделяют стадии - Над каждым Exchange будет оператор агрегации или join. Смотрите метрики этих операторов
Каждый оператор в DAG-графе показывает:
- number of output rows: сколько строк оператор произвёл
- spill size: сколько данных было сброшено на диск
- peak memory usage: пиковое потребление памяти
Если SortMergeJoin показывает spill size = 2.1 GB - это ваш виновник.
Связь Stage ID с кодом¶
В Spark UI каждый Stage имеет ID и содержит ссылку на строки кода, которые создали его. Нажмите на Stage → посмотрите секцию Associated Job IDs и Details - там будет стектрейс с именами методов и номерами строк PySpark-кода.
Это позволяет однозначно ответить: "Skew возник из-за df.join(other, on='user_id', how='left') в строке 47 файла pipeline.py".
Программная диагностика: находим горячие ключи¶
Spark UI показывает симптомы. PySpark-запросы позволяют найти причину - конкретные ключи, которые создают skew.
Шаг 1: Создаём синтетический skewed датасет¶
from pyspark.sql import SparkSession
from pyspark.sql.functions import (
col, count, desc, lit, when, rand, concat_ws,
percentile_approx, sum as spark_sum, round as spark_round,
spark_partition_id, countDistinct
)
import pyspark.sql.functions as F
spark = SparkSession.builder \
.appName("skew-detection") \
.config("spark.sql.shuffle.partitions", "50") \
.getOrCreate()
# Создаём датасет с искусственным skew:
# user_id = "bot_001" встречается в 40% строк
# user_id = "guest" встречается в 20% строк
# Остальные 40% - нормальные уникальные пользователи
from pyspark.sql.types import *
n_rows = 2_000_000
# Генерируем строки с перекосом
clickstream = spark.range(n_rows).select(
col("id").alias("event_id"),
when(rand() < 0.40, lit("bot_001")) # 40% - один бот
.when(rand() < 0.60, lit("guest")) # 20% (0.60 - 0.40) - гость
.when(rand() < 0.70, lit("unknown")) # 10% - unknown
.otherwise(concat_ws("_", lit("user"), (col("id") % 50000).cast("string")))
.alias("user_id"),
(rand() * 100).alias("page_id"),
F.current_timestamp().alias("event_ts"),
)
clickstream.cache()
clickstream.count() # материализуем
Шаг 2: Профилирование ключей - находим hot keys¶
Первый диагностический запрос - распределение значений по ключу join/groupBy:
# Топ-20 значений ключа по количеству строк
key_distribution = clickstream \
.groupBy("user_id") \
.count() \
.withColumn("pct_of_total", spark_round(col("count") / n_rows * 100, 2)) \
.orderBy(desc("count"))
key_distribution.show(20, truncate=False)
+---------+--------+------------+
| user_id| count|pct_of_total|
+---------+--------+------------+
| bot_001| 800234| 40.01|
| guest| 399876| 19.99|
| unknown| 199843| 9.99|
| user_0 | 47| 0.00|
| user_1 | 42| 0.00|
| user_2 | 38| 0.00|
| ... | ...| ...|
+---------+--------+------------+
Картина очевидна: три ключа (bot_001, guest, unknown) содержат 70% всех строк. Если делать join или groupBy по user_id, executor, обрабатывающий bot_001, получит в 17 000 раз больше строк, чем executor с user_0.
Шаг 3: Статистика распределения¶
Для более полного понимания - статистические метрики:
# Агрегат: сколько уникальных ключей и как распределены строки
key_stats = clickstream \
.groupBy("user_id") \
.count() \
.agg(
F.count("user_id").alias("distinct_keys"),
F.min("count").alias("min_rows_per_key"),
F.max("count").alias("max_rows_per_key"),
F.avg("count").alias("avg_rows_per_key"),
percentile_approx("count", 0.5).alias("median_rows"),
percentile_approx("count", 0.95).alias("p95_rows"),
percentile_approx("count", 0.99).alias("p99_rows"),
)
key_stats.show()
+------------+----------------+----------------+----------------+-----------+---------+---------+
|distinct_keys|min_rows_per_key|max_rows_per_key|avg_rows_per_key|median_rows|p95_rows |p99_rows |
+------------+----------------+----------------+----------------+-----------+---------+---------+
| 50003| 1| 800234| 39.9| 40| 47| 200|
+------------+----------------+----------------+----------------+-----------+---------+---------+
Интерпретация: 50 003 уникальных ключа. Медиана - 40 строк на ключ. Максимум - 800 234 строк. Отношение max/median = 20 006. Это катастрофический skew.
Шаг 4: Анализ NULL-аномалий¶
NULL - особый случай. В Spark NULL != NULL, поэтому при join по ключу с NULL: строки с NULL не будут матчиться (в INNER join) или создадут отдельный overloaded bucket (в некоторых реализациях). Нужно оценить их долю:
from pyspark.sql.functions import isnull
null_analysis = clickstream.agg(
count("*").alias("total_rows"),
count(when(isnull("user_id"), True)).alias("null_user_id"),
count(when(col("user_id") == "unknown", True)).alias("unknown_user_id"),
count(when(col("user_id") == "guest", True)).alias("guest_user_id"),
).withColumn(
"null_pct",
spark_round(col("null_user_id") / col("total_rows") * 100, 2)
).withColumn(
"unknown_pct",
spark_round(col("unknown_user_id") / col("total_rows") * 100, 2)
)
null_analysis.show()
+----------+------------+---------------+-------------+--------+-----------+
|total_rows|null_user_id|unknown_user_id|guest_user_id|null_pct|unknown_pct|
+----------+------------+---------------+-------------+--------+-----------+
| 2000000| 0| 199843| 399876| 0.0| 9.99|
+----------+------------+---------------+-------------+--------+-----------+
Шаг 5: Оценка влияния hot keys на shuffle объём¶
Важно понять: если мы сделаем join по user_id, сколько байт будет переведено через shuffle для каждого ключа?
# Приблизительный объём shuffle по ключу
shuffle_estimate = clickstream \
.groupBy("user_id") \
.agg(
count("*").alias("row_count"),
# Оцениваем ~100 байт на строку (event_id + user_id + page_id + ts)
(count("*") * 100).alias("estimated_shuffle_bytes"),
) \
.orderBy(desc("row_count")) \
.limit(10)
shuffle_estimate.withColumn(
"estimated_shuffle_mb",
spark_round(col("estimated_shuffle_bytes") / 1024 / 1024, 1)
).select("user_id", "row_count", "estimated_shuffle_mb").show()
+---------+--------+--------------------+
| user_id|row_count|estimated_shuffle_mb|
+---------+--------+--------------------+
| bot_001| 800234| 76.3|
| guest| 399876| 38.1|
| unknown| 199843| 19.1|
| user_0 | 47| 0.0|
+---------+--------+--------------------+
Только три ключа создадут 133 MB из ~190 MB общего shuffle трафика. Остальные 50 000 ключей - менее 57 MB суммарно.
Наблюдение skew в action: join с перекосом¶
Создадим join на skewed датасете и посмотрим что происходит:
# Таблица профилей пользователей
user_profiles = spark.range(50000).select(
concat_ws("_", lit("user"), col("id")).alias("user_id"),
(rand() * 10).cast("int").alias("age_group"),
F.element_at(F.array(lit("RU"), lit("US"), lit("DE")),
((col("id") % 3) + 1).cast("int")).alias("country")
)
# Добавляем "горячих" пользователей
hot_users = spark.createDataFrame([
("bot_001", 99, "XX"),
("guest", 0, "XX"),
("unknown", 0, "XX"),
], ["user_id", "age_group", "country"])
all_profiles = user_profiles.union(hot_users)
# JOIN: clickstream × profiles
# Это МЕДЛЕННЫЙ join - будет skew
spark.conf.set("spark.sql.adaptive.enabled", "false") # отключаем AQE для наглядности
result = clickstream.join(
all_profiles,
on="user_id",
how="left"
).groupBy("country").count()
result.show()
После выполнения этого кода откройте Spark UI: stage с join будет иметь один task, работающий значительно дольше остальных.
Визуализация partition size до и после shuffle¶
# Смотрим распределение ПОСЛЕ shuffle (partition distribution post-join)
result_with_pid = clickstream.join(
all_profiles.hint("skew", "user_id"),
on="user_id",
how="left"
).select(spark_partition_id().alias("pid"), "user_id") \
partition_sizes = result_with_pid \
.groupBy("pid") \
.agg(
count("*").alias("rows_in_partition"),
countDistinct("user_id").alias("distinct_keys"),
) \
.orderBy(desc("rows_in_partition"))
# Топ-5 самых больших партиций
partition_sizes.show(10)
+----+------------------+-------------+
| pid| rows_in_partition|distinct_keys|
+----+------------------+-------------+
| 7 | 800234| 1| ← bot_001
| 31 | 399876| 1| ← guest
| 42 | 199843| 1| ← unknown
| 0 | 1245| 31|
| 1 | 1198| 30|
| 2 | 1267| 32|
| 3 | 1231| 31|
+----+------------------+-------------+
Картина предельно ясна: partition 7 содержит 800 234 строки и только один уникальный ключ (bot_001). Это та самая перекошенная партиция, которую один executor будет обрабатывать пока остальные 49 ждут.
Диагностика skew в aggreration и window functions¶
Skew не ограничивается join. Проверяем groupBy:
# Skew в groupBy
skewed_agg = clickstream \
.groupBy("user_id") \
.agg(
count("*").alias("event_count"),
F.max("event_ts").alias("last_event"),
)
# Измеряем размеры партиций в результате агрегации
skewed_agg.select(spark_partition_id().alias("pid"), "user_id") \
.groupBy("pid") \
.agg(count("*").alias("keys_in_partition")) \
.orderBy(desc("keys_in_partition")) \
.show(5)
Skew в window functions¶
Window functions используют partitionBy - по сути тот же shuffle:
from pyspark.sql.window import Window
from pyspark.sql.functions import row_number, rank
# ОПАСНО: partitionBy("user_id") при skewed user_id
window_spec = Window.partitionBy("user_id").orderBy("event_ts")
# Весь bot_001 - 800K строк - попадёт в одну партицию
# Executor должен держать в памяти весь отсортированный массив для user_id
result_with_rank = clickstream.withColumn(
"session_rank",
row_number().over(window_spec)
)
# Этот код выполнится, но будет очень медленным и, вероятно, вызовет spill
Количественные критерии для классификации skew¶
На основе метрик из Spark UI можно формализовать оценку:
def assess_skew(df, key_col: str, sample_fraction: float = 1.0) -> None:
"""
Профилирует распределение ключа и оценивает severity skew.
При большом датасете можно использовать sample_fraction < 1.0.
"""
if sample_fraction < 1.0:
sample_df = df.sample(fraction=sample_fraction, seed=42)
else:
sample_df = df
key_counts = sample_df.groupBy(key_col).count()
stats = key_counts.agg(
count(key_col).alias("distinct_keys"),
F.min("count").alias("min_count"),
F.max("count").alias("max_count"),
F.avg("count").alias("avg_count"),
percentile_approx("count", 0.5).alias("p50"),
percentile_approx("count", 0.9).alias("p90"),
percentile_approx("count", 0.99).alias("p99"),
).collect()[0]
max_to_median = stats["max_count"] / max(stats["p50"], 1)
p99_to_p50 = stats["p99"] / max(stats["p50"], 1)
print(f"\n=== Skew Assessment for column: {key_col} ===")
print(f" Distinct keys: {stats['distinct_keys']:,}")
print(f" Min rows per key: {stats['min_count']:,}")
print(f" Median rows: {stats['p50']:,}")
print(f" P99 rows: {stats['p99']:,}")
print(f" Max rows per key: {stats['max_count']:,}")
print(f" Max/Median ratio: {max_to_median:.0f}x")
print(f" P99/Median ratio: {p99_to_p50:.0f}x")
if max_to_median > 100:
severity = "CRITICAL - немедленное вмешательство"
elif max_to_median > 20:
severity = "HIGH - требует исправления"
elif max_to_median > 5:
severity = "MEDIUM - стоит исследовать"
else:
severity = "LOW - приемлемо"
print(f"\n Severity: {severity}")
# Топ hot keys
print(f"\n Top-10 hot keys:")
top_keys = key_counts.orderBy(desc("count")).limit(10).collect()
total = sample_df.count() if sample_fraction < 1.0 else df.count()
for row in top_keys:
pct = row["count"] / total * 100
bar = "█" * int(pct / 2)
print(f" {str(row[key_col]):20s} {row['count']:>10,} rows ({pct:5.1f}%) {bar}")
# Применяем
assess_skew(clickstream, "user_id")
=== Skew Assessment for column: user_id ===
Distinct keys: 50,003
Min rows per key: 1
Median rows: 40
P99 rows: 200
Max rows per key: 800,234
Max/Median ratio: 20006x
Severity: CRITICAL - немедленное вмешательство
Top-10 hot keys:
bot_001 800,234 rows (40.0%) ████████████████████
guest 399,876 rows (20.0%) ██████████
unknown 199,843 rows (10.0%) █████
user_12345 47 rows ( 0.0%)
user_28901 46 rows ( 0.0%)
Быстрый sampling для больших датасетов¶
Полный scan 100+ GB датасета только для профилирования ключей - расточительно. Используйте sampling:
# Для 1 TB датасета - сэмплируем 1%, результат статистически надёжен
sample_fraction = 0.01
assess_skew(huge_df, "join_key", sample_fraction=sample_fraction)
# Или approximate distinct count через HLL
from pyspark.sql.functions import approx_count_distinct
quick_stats = clickstream.agg(
count("*").alias("total"),
approx_count_distinct("user_id", rsd=0.05).alias("approx_distinct_users"),
).show()
approx_count_distinct использует алгоритм HyperLogLog - точность 5% (rsd=0.05), скорость в 10-20 раз выше точного countDistinct.
Skew и object storage: усиленный эффект¶
На S3 / MinIO skew имеет дополнительный эффект, которого нет при работе с HDFS.
S3 ограничивает пропускную способность на уровне prefix partition: 5500 GET/s и 3500 PUT/s на один prefix (на один "хот" ключ в контексте S3-путей). При shuffle executor с перекошенной партицией делает тысячи мелких HTTP GET-запросов к S3 для чтения shuffle-файлов - и может упереться в лимиты S3 API.
Кроме того, при spill-to-disk executor пишет временные файлы. На S3 нет "временных файлов" - всё идёт через HTTP. Spill на S3 - это медленнее, чем spill на локальный SSD в 10-50 раз.
Вывод: на object storage объём spill нужно поддерживать строго равным нулю. Это ещё один аргумент для борьбы со skew.
Мониторинг skew в production¶
Spark History Server¶
После завершения джобов их метрики сохраняются в Spark History Server (если настроен spark.eventLog.enabled=true). Это позволяет проводить post-mortem анализ неудавшихся или медленных джобов.
# В spark-defaults.conf или при создании SparkSession
spark = SparkSession.builder \
.config("spark.eventLog.enabled", "true") \
.config("spark.eventLog.dir", "s3a://bucket/spark-logs/") \
.getOrCreate()
В History Server ищите паттерны:
- Stage с duration > ожидаемого в 3+ раза
- Tasks с max duration >> median duration
- Non-zero Spill (Disk) в любом stage
Автоматические алерты¶
В production-окружениях (Databricks, EMR, GCP Dataproc) можно настроить алерты на метрики:
- Task duration 95th percentile > threshold
- Spill > 0 MB в финансово-критичных джобах
- Stage duration > SLA
Логирование метрик из кода¶
from pyspark.sql.functions import count, max as spark_max, min as spark_min
def log_partition_stats(df, label: str):
"""Логирует распределение партиций для мониторинга."""
stats = df.select(spark_partition_id().alias("pid")) \
.groupBy("pid").count() \
.agg(
spark_min("count").alias("min_rows"),
spark_max("count").alias("max_rows"),
F.avg("count").alias("avg_rows"),
count("pid").alias("num_partitions"),
).collect()[0]
skew_ratio = stats["max_rows"] / max(stats["avg_rows"], 1)
print(f"[{label}] partitions={stats['num_partitions']}, "
f"min={stats['min_rows']}, avg={stats['avg_rows']:.0f}, "
f"max={stats['max_rows']}, skew_ratio={skew_ratio:.1f}x")
if skew_ratio > 10:
print(f" ⚠ WARNING: high skew detected (ratio={skew_ratio:.0f}x)")
# Используем перед тяжёлыми операциями
log_partition_stats(clickstream.repartition(50, "user_id"), "before_join")
Лабораторная работа: полная диагностика аварии¶
Сценарий¶
У вас есть пайплайн, который падает через 90 минут. Финальный stage застрял на 1/200 tasks remaining. Spark UI недоступен напрямую - вы получаете метрики через History Server.
Данные из History Server:
Stage 5: "SortMergeJoin [user_id]"
Duration: 1h 38m
Tasks: 200/200 completed
Task Duration:
Min: 0.8s | 25%: 2.1s | Median: 3.4s | 75%: 4.2s | Max: 1h 37m
Shuffle Read Size:
Min: 0.5MB | 25%: 1.8MB | Median: 2.3MB | 75%: 2.9MB | Max: 6.2GB
Shuffle Read Records:
Min: 5K | 25%: 18K | Median: 23K | 75%: 29K | Max: 61.8M
Spill (Memory): Max = 28.4GB
Spill (Disk): Max = 4.7GB
Associated query: events.join(users, on="user_id", how="left")
Диагностический аудит¶
Проведём анализ по методологии Senior Data Engineer:
# 1. Подтверждаем наличие skew по квантилям
max_duration_sec = 97 * 60 # 1h 37m в секундах
median_duration_sec = 3.4
skew_ratio_time = max_duration_sec / median_duration_sec
print(f"Task Duration Skew: {skew_ratio_time:.0f}x") # → 1706x
max_shuffle_bytes = 6.2 * 1024 # MB
median_shuffle_bytes = 2.3
skew_ratio_shuffle = max_shuffle_bytes / median_shuffle_bytes
print(f"Shuffle Read Skew: {skew_ratio_shuffle:.0f}x") # → 2762x
max_records = 61_800_000
median_records = 23_000
skew_ratio_records = max_records / median_records
print(f"Records Skew: {skew_ratio_records:.0f}x") # → 2687x
# 2. Оцениваем объём spill
print(f"\nSpill to disk: 4.7 GB - критично для производительности")
print(f"Spill to memory: 28.4 GB - executor работает за пределами памяти")
# 3. Делаем вывод
print("\n=== Заключение ===")
print("Диагноз: Data Skew в SortMergeJoin [user_id]")
print(f" Severity: CRITICAL (time ratio: {skew_ratio_time:.0f}x)")
print(f" Проблемный ключ: 1 значение user_id содержит 61.8M строк")
print(f" vs медиана 23K строк - это 2687x превышение нормы")
print(f" Следствие: Executor вышел за пределы памяти:")
print(f" Spill(Disk) = 4.7 GB, Spill(Memory) = 28.4 GB")
print(f" Рекомендация: проверить распределение user_id,")
print(f" изолировать NULL/sentinel значения,")
print(f" применить AQE skewJoin или salting")
Задание для самостоятельной работы¶
- Запустить код из раздела "Создаём синтетический skewed датасет" с параметрами:
- 5 000 000 строк
hot_key_1- 35% строкhot_key_2- 25% строк-
Остальные 40% - 100 000 уникальных ключей
-
Выполнить join с таблицей profiles (100 003 строки: по одной на каждый ключ)
-
Открыть Spark UI → Stage с join → сделать скриншот Task Duration metrics
-
Выполнить функцию
assess_skewи зафиксировать: - max/median ratio для duration
- max/median ratio для shuffle read
-
наличие spill
-
Написать аудиторское заключение: подтвердить skew, указать горячие ключи, оценить объём spill
Anti-patterns¶
Anti-pattern 1: Игнорировать skew "потому что работает"¶
# Джоб работает за 2 часа - "нормально", никто не исследует
result = huge_events.join(profiles, on="user_id")
result.write.parquet(...)
# На самом деле: медиана задач - 3 секунды, один task - 1.5 часа
# При росте данных в 2x время вырастет не в 2x, а в 10-20x
# Skew усиливается нелинейно с ростом данных
Anti-pattern 2: Слепой repartition без анализа ключа¶
# "Перераспределяем" 1000 партиций - не помогает при hot keys
df.repartition(1000).join(other, on="hot_key")
# Если hash("bot_001") % 1000 = 42, все 800K строк всё равно пойдут
# в partition 42. Количество партиций не решает проблему hot keys.
Anti-pattern 3: Запускать тяжёлый join без предварительного профилирования¶
# НЕПРАВИЛЬНО: сразу запускаем join на терабайтном датасете
result = tb_scale_events.join(profiles, on="user_id") # упадёт через час
# ПРАВИЛЬНО: сначала профилируем ключ на 1% sample
assess_skew(tb_scale_events.sample(0.01), "user_id")
# Потом принимаем архитектурное решение: AQE / broadcast / salting
Anti-pattern 4: Не смотреть на spill в Spark UI¶
# "Джоб завершился, всё ок"
# На самом деле: Spill(Disk) = 50 GB, время в 5x дольше оптимального
# Spill - это молчаливый убийца производительности
Что дальше: методы лечения¶
Обнаружить skew - половина работы. Вторая половина - устранить его. Основные инструменты:
AQE Skew Join (spark.sql.adaptive.skewJoin.enabled=true) - Spark 3.x автоматически обнаруживает перекошенные партиции после shuffle и разбивает их на несколько подзадач. Работает без изменения кода, но требует включения AQE и не справляется с экстремальным skew.
Broadcast Join - для небольшой таблицы (< spark.sql.autoBroadcastJoinThreshold, по умолчанию 10 MB) Spark транслирует её на все executors. Никакого shuffle, никакого skew. Работает только если малая сторона join помещается в память executor.
Salting - добавление случайного суффикса к горячим ключам ("bot_001_3", "bot_001_7" и т.д.) для распределения их по нескольким партициям. Требует изменения логики join и дополнительной агрегации.
Pre-filtering / изоляция sentinel-значений - отдельная обработка NULL, "unknown", "guest" вне основного join-пайплайна.
Эти техники разбираются в следующем уроке. Первый шаг всегда один: найти и измерить skew - чем мы и занялись в этом уроке.
Чеклист¶
- [ ] Понимаю что такое data skew и почему один straggler task блокирует весь stage
- [ ] Знаю четыре вида skew: data skew, partition skew, explode skew, window skew
- [ ] Умею читать Summary Metrics for Completed Tasks в Spark UI: Duration, Shuffle Read, Spill
- [ ] Знаю пороги severity: max/median > 10 - явный skew, > 100 - критический
- [ ] Понимаю что означает non-zero Spill (Disk) и почему это критично
- [ ] Умею находить горячие ключи через
groupBy(key).count().orderBy(desc("count")) - [ ] Знаю функцию
assess_skew- могу написать аналогичную диагностику самостоятельно - [ ] Понимаю особую опасность sentinel-значений (NULL, "unknown", "guest") как hot keys
- [ ] Умею использовать sampling для профилирования больших датасетов без полного scan
- [ ] Знаю как найти проблемный оператор через SQL tab и DAG-граф в Spark UI
- [ ] Понимаю почему slind repartition не спасает от hot key skew
- [ ] Могу написать аудиторское заключение по метрикам из History Server