Skew Detection: как увидеть skew в Task Duration и shuffle read

Анатомия data skew: горячие ключи, straggler tasks, shuffle read spill. Как читать Spark UI, писать диагностические запросы и готовить данные к лечению.

optimization

Анатомия 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"):

  1. Map side (shuffle write): каждый executor вычисляет для каждой строки номер partition назначения: partition_id = hash(key) % numPartitions. Результаты записываются в shuffle files на локальный диск
  2. Shuffle: каждый reduce-executor читает из всех map-executors строки для своих partition_id
  3. 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 для каждого оператора. Это позволяет точно локализовать узкое место.

Как найти проблемный оператор

  1. Откройте Spark UI → SQL
  2. Найдите SQL-запрос с большим временем выполнения
  3. Нажмите на запрос - откроется DAG-граф физического плана
  4. Ищите узлы Exchange (shuffle) - они разделяют стадии
  5. Над каждым 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")

Задание для самостоятельной работы

  1. Запустить код из раздела "Создаём синтетический skewed датасет" с параметрами:
  2. 5 000 000 строк
  3. hot_key_1 - 35% строк
  4. hot_key_2 - 25% строк
  5. Остальные 40% - 100 000 уникальных ключей

  6. Выполнить join с таблицей profiles (100 003 строки: по одной на каждый ключ)

  7. Открыть Spark UI → Stage с join → сделать скриншот Task Duration metrics

  8. Выполнить функцию assess_skew и зафиксировать:

  9. max/median ratio для duration
  10. max/median ratio для shuffle read
  11. наличие spill

  12. Написать аудиторское заключение: подтвердить 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