Storage Partitioning: выбор колонок для partitionBy и анти-паттерны

Физика partitionBy в файловых хранилищах, цена Directory Listing в S3 и HDFS, правила выбора ключа партиционирования, три главных анти-паттерна и как их диагностировать, Hidden Partitioning в Iceberg и Liquid Clustering в Delta Lake.

optimization

В предыдущем уроке мы разобрали Bucketing - механизм, который устраняет Shuffle при join двух больших таблиц. Сегодня поговорим о вещи более фундаментальной: как именно данные лежат на диске и почему неправильный выбор ключа partitionBy может сделать таблицу бесполезной - даже если формат Parquet, AQE включён и всё по учебнику.

Storage Partitioning - первый и самый дешёвый уровень оптимизации в любом Lakehouse. Перед тем как Spark прочтёт хоть один байт данных, он делает Directory Listing и решает, какие папки вообще открывать. Правильно спроектированное партиционирование сокращает объём чтения на порядки. Неправильное - убивает кластер ещё на стадии планирования запроса.


1. Физика процесса: как Spark и хранилище видят partitionBy

Анатомия директорий на диске

Когда вы вызываете df.write.partitionBy("dt", "country").parquet(path), Spark не просто записывает файлы - он создаёт иерархию директорий, где каждый уровень вложенности соответствует одной partition-колонке, а имя директории содержит значение этой колонки.

df.write \
    .mode("overwrite") \
    .partitionBy("event_date", "country") \
    .parquet("s3a://data/events/")

# Результирующая структура:
# s3a://data/events/
# ├── event_date=2026-01-01/
# │   ├── country=RU/
# │   │   └── part-00000-abc123.snappy.parquet   (128 МБ)
# │   │   └── part-00001-abc123.snappy.parquet   (95 МБ)
# │   ├── country=US/
# │   │   └── part-00000-def456.snappy.parquet   (512 МБ)
# │   └── country=DE/
# │       └── part-00000-ghi789.snappy.parquet   (43 МБ)
# ├── event_date=2026-01-02/
# │   ├── country=RU/
# │   │   └── ...
# ...

Значения partition-колонок не хранятся внутри Parquet-файлов - они закодированы в пути к файлу. Когда Spark читает такую таблицу, он реконструирует значения event_date и country из имён папок, а не из содержимого файлов. Это экономит место: колонки партиционирования не дублируются в каждой строке файла.

Directory Listing: самая медленная часть запроса

До того как Spark прочитает первый байт данных, Driver выполняет Directory Listing - обходит дерево директорий, собирает список файлов и формирует план чтения. Именно здесь кроется главная проблема неправильного партиционирования.

При правильном партиционировании Directory Listing занимает миллисекунды: Driver обходит несколько папок и сразу знает, где нужные данные. При неправильном - он обходит миллионы папок только чтобы найти несколько нужных файлов.

Почему S3 ListObjects - это дорого

В HDFS Directory Listing - локальная операция на NameNode: ответ за миллисекунды. В облачных объектных хранилищах (S3, GCS, Azure ADLS) ситуация принципиально иная:

  • Каждый ListObjects возвращает не более 1000 объектов за один HTTP-запрос
  • При 100 000 папок нужно минимум 100 HTTP-запросов только для листинга
  • Каждый запрос - это сетевой RTT 5–50 мс плюс время сервера
  • S3 тарифицирует каждый ListObjects запрос (0.005$ за 1000 запросов в AWS)
  • ListObjects не транзакционен: при конкурентной записи возможна несогласованная картина
# Диагностика: сколько файлов в таблице?
import subprocess
result = subprocess.run(
    ["aws", "s3", "ls", "--recursive", "s3://data/events/", "--summarize"],
    capture_output=True, text=True
)
# Total Objects: 2 847 392   ← если такое число - у вас проблема
# Total Size: 1.2 TiB

# Через Spark: посмотреть число файлов на партицию
spark.read.parquet("s3a://data/events/") \
     .select("event_date", "country") \
     .groupBy("event_date", "country") \
     .count() \
     .orderBy("count") \
     .show(20)
# Это покажет число строк, но не файлов; для файлов нужен ls

При 2.8 миллионах файлов только их перечисление займёт ~2847 HTTP-запросов к S3 - это несколько секунд задержки до начала любого запроса к этой таблице. А Driver, который держит в памяти список всех файлов, расходует сотни мегабайт heap.


2. Стратегия выбора ключа партиционирования

Три правила идеального partition-ключа

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

Правило 1: Умеренная кардинальность

Кардинальность (cardinality) - число уникальных значений ключа. Для partition-ключа нужен баланс: слишком мало уникальных значений даёт слабый pruning эффект, слишком много - катастрофа с файлами.

Практическое правило: от 50 до ~10 000 уникальных значений partition-ключа - рабочий диапазон. За 10 000 нужно очень веское обоснование и чёткий план обслуживания.

# Проверить кардинальность перед выбором ключа
df.select("event_date", "country", "status", "user_id").describe().show()

# Точный подсчёт уникальных значений
from pyspark.sql import functions as F
df.agg(
    F.countDistinct("event_date").alias("dates"),
    F.countDistinct("country").alias("countries"),
    F.countDistinct("status").alias("statuses"),
    F.countDistinct("user_id").alias("users"),
).show()
# +-----+---------+--------+-----------+
# |dates|countries|statuses|      users|
# +-----+---------+--------+-----------+
# | 1095|      195|       4|100,000,000|  ← user_id нельзя!

Правило 2: Равномерное распределение данных

Кардинальность - необходимое, но не достаточное условие. Важно ещё и то, как данные распределяются по значениям ключа. Если 95% строк имеют country = 'US', то партиция US будет в десятки раз больше остальных.

# Проверить распределение данных по потенциальному partition-ключу
df.groupBy("country") \
  .count() \
  .withColumn("pct", F.round(F.col("count") / df.count() * 100, 2)) \
  .orderBy(F.col("count").desc()) \
  .show(20)
# +---------+----------+-----+
# |  country|     count|  pct|
# +---------+----------+-----+
# |       US|  4800000 |48.0 |  ← почти половина таблицы в одной партиции
# |       IN|   800000 | 8.0 |
# |       BR|   600000 | 6.0 |
# |       RU|   400000 | 4.0 |
# ...  ещё 191 страна ...
# |      SJ|        50 | 0.0 |  ← 50 строк в партиции = 1 крошечный файл

Критерий нормального распределения: ни одна партиция не занимает больше 20–30% от общего объёма данных. Если одна страна содержит 48% данных - запросы по ней будут медленными независимо от партиционирования.

Правило 3: Ключ присутствует в бизнес-запросах

Партиционирование даёт выгоду только если WHERE-условие запроса использует partition-ключ. Если аналитики никогда не фильтруют по status, то партиционирование по status бесполезно - оно только увеличивает число файлов без какой-либо выгоды.

# ПЕРЕД выбором ключа: проанализировать паттерны существующих запросов
# Что чаще всего в WHERE-условиях?

# Типичный анализ через Spark History Server / query logs:
# SELECT * FROM events WHERE event_date = ?         → 95% запросов
# SELECT * FROM events WHERE country = 'RU'         → 40% запросов
# SELECT * FROM events WHERE user_id = ?            → 60% запросов, НО! user_id → high cardinality
# SELECT * FROM events WHERE status = 'active'      → 30% запросов

# Вывод: partitionBy("event_date") - однозначно; partitionBy("country") - опционально;
# partitionBy("user_id") - НЕЛЬЗЯ (high cardinality); user_id можно через Bucketing

Золотой стандарт: дата как первый partition-ключ

Для подавляющего большинства аналитических таблиц дата - лучший выбор. Почему:

  • Аналитические запросы почти всегда фильтруют по времени: «за вчера», «за последние 30 дней», «за Q1 2026»
  • Данные поступают append-only: каждый день - новая папка, старые не меняются
  • Кардинальность предсказуема: ~365 значений в год, управляемо
  • Retention Policy: удалять старые данные = удалять папки по дате, просто и безопасно
  • Инкрементальная обработка: читать только новые партиции (WHERE dt > last_processed_dt)
# ─── Правильный дизайн для event-таблицы ─────────────────────────────────────
events_df.write \
    .mode("overwrite") \
    .partitionBy("event_date") \           # ← первый ключ: дата
    .parquet("s3a://data/events/")

# Для таблиц с региональной аналитикой - добавить вторым ключом регион:
events_df.write \
    .partitionBy("event_date", "region") \ # ← два уровня: дата + регион
    .parquet("s3a://data/events_regional/")
# Но проверьте: кардинальность region? Распределение?
# 50 регионов × 365 дней = 18 250 партиций - на грани

# ─── Выбор гранулярности времени ─────────────────────────────────────────────
# Для больших таблиц (>100 ГБ/день) дневная гранулярность норма
events_df.withColumn("dt", F.to_date("event_ts")) \
         .write.partitionBy("dt").parquet(path)

# Для средних таблиц (1-10 ГБ/день) можно месячную гранулярность
events_df.withColumn("ym", F.date_format("event_ts", "yyyy-MM")) \
         .write.partitionBy("ym").parquet(path)

# Для маленьких таблиц (<1 ГБ/день) - годовую или вообще без партиционирования
events_df.withColumn("yr", F.year("event_ts")) \
         .write.partitionBy("yr").parquet(path)

3. Главные анти-паттерны и катастрофы в продакшне

Анти-паттерн 1: Высокая кардинальность - «Король боли»

Это самый разрушительный анти-паттерн. Партиционирование по user_id, order_id, uuid или timestamp создаёт миллионы папок, каждая из которых содержит один-два крошечных файла.

Что происходит физически:

s3a://data/events/
├── user_id=1/          part-00000.parquet    (2 KB)
├── user_id=2/          part-00000.parquet    (2 KB)
├── user_id=3/          part-00000.parquet    (2 KB)
...  100 миллионов директорий ...
├── user_id=99999998/   part-00000.parquet    (2 KB)
├── user_id=99999999/   part-00000.parquet    (2 KB)
└── user_id=100000000/  part-00000.parquet    (2 KB)

100 миллионов директорий × 2 КБ файлов = 200 ГБ данных в 100 миллионах файлов. Реальный размер данных 200 ГБ, но производительность системы - как при работе с 10 ТБ.

Что происходит с Spark Driver:

В HDFS аналогичная проблема: NameNode хранит метаданные всех файлов в памяти. 100 миллионов файлов × ~150 байт метаданных = 15 ГБ на NameNode только для одной таблицы. NameNode падает или деградирует, и проблема затрагивает весь кластер.

# ❌ КАТАСТРОФА: партиционирование по user_id
events_df.write \
    .partitionBy("user_id") \  # user_id: 100 млн уникальных значений!
    .parquet("s3a://data/events/")
# Результат: 100 млн папок, каждая с 1-2 файлами по 2 КБ

# ❌ ЕЩЁ ХУЖЕ: партиционирование по timestamp
events_df.write \
    .partitionBy("event_ts") \  # timestamp: каждая строка уникальна
    .parquet("s3a://data/events/")
# Результат: столько папок, сколько строк в таблице

# ✅ Решение: агрегировать до умеренной кардинальности
events_df \
    .withColumn("event_date", F.to_date("event_ts")) \  # timestamp → date
    .write \
    .partitionBy("event_date") \  # ~1000 значений за 3 года - нормально
    .parquet("s3a://data/events/")

Анти-паттерн 2: Перекос данных (Data Skew)

Равномерная кардинальность - необходимое, но не достаточное условие. Если данные распределены неравномерно, одна партиция становится «горячей» - и всё упирается в скорость её обработки.

Типичные примеры перекоса:

  • country: 80% данных глобального e-commerce - США. Партиция country=US в 50 раз больше country=FR.
  • event_date для временного трафика: партиции за Black Friday в 10 раз больше обычного дня.
  • tenant_id в multi-tenant SaaS: топ-5 клиентов занимают 90% данных.
# Диагностика перекоса: смотрим размер файлов по партициям
from pyspark.sql import functions as F

# Метод 1: считаем строки на партицию
df.groupBy("country") \
  .count() \
  .withColumn("pct", F.round(F.col("count") / F.lit(df.count()) * 100, 2)) \
  .orderBy(F.col("count").desc()) \
  .show(10)

# Метод 2: смотрим на Task Duration в Spark UI
# Если большинство задач заняли 2 секунды, а одна - 45 минут: это skew

# Метод 3: через AQE Skew Detection
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
# AQE автоматически разбивает горячие партиции при чтении

Что происходит при skewed партициях в запросе:

Три Executor-а простаивают 44 минуты, ожидая когда обработается партиция US. Эффективная утилизация кластера - 25%. Весь кластер из 100 машин работает со скоростью одной.

Решения для skewed партиций:

# Решение 1: Добавить второй ключ для дробления горячей партиции
events_df.write \
    .partitionBy("event_date", "country") \
    # Теперь US делится по дням: country=US/event_date=2026-01-01/ и т.д.
    # Один запрос по дате = только один день US вместо всего US
    .parquet(path)

# Решение 2: Синтетический ключ для балансировки (salting)
events_df \
    .withColumn("partition_shard", (F.hash("user_id") % 10).cast("int")) \
    # Разбиваем горячий ключ на 10 подпартиций случайно
    .write \
    .partitionBy("country", "partition_shard") \
    .parquet(path)
# Недостаток: запрос по country должен знать про partition_shard

# Решение 3: AQE Skew Join Handling (при join)
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionFactor", "5")
# AQE автоматически разбивает партиции, которые в 5x больше медианы

Анти-паттерн 3: Слишком глубокая вложенность

Многоуровневое партиционирование кажется привлекательным: «Если два уровня хорошо, то пять - отлично!» На практике каждый дополнительный уровень умножает число директорий.

# ❌ ПРОБЛЕМА: экспоненциальный рост директорий
events_df.write \
    .partitionBy("year", "month", "day", "hour", "country") \
    .parquet(path)

# Математика катастрофы:
# year: 3 значения
# month: 12 значений
# day: 31 значение
# hour: 24 значения
# country: 195 значений
# Итого: 3 × 12 × 31 × 24 × 195 = 65 304 960 директорий (для 3 лет данных)

# При реальных данных: не все комбинации существуют, но порядок цифр тот же.
# Листинг 65 миллионов директорий = десятки минут overhead на каждый запрос.

Золотое правило многоуровневого партиционирования: число директорий не должно превышать ~10 000 – 50 000. За этой границей Directory Listing становится узким местом.

# ✅ Решение: агрегировать временные уровни, убрать лишние измерения
events_df \
    .withColumn("dt", F.to_date("event_ts")) \  # year+month+day → одна дата
    .write \
    .partitionBy("dt") \                         # только один уровень
    .parquet(path)
# 365 директорий в год - идеально

# Или максимум два уровня с умеренной кардинальностью:
events_df \
    .withColumn("dt", F.to_date("event_ts")) \
    .write \
    .partitionBy("dt", "region") \               # дата + регион (50 значений)
    .parquet(path)
# 365 × 50 = 18 250 директорий - на верхней границе нормы

4. Практика: моделируем и чиним проблемы

Подготовка данных для демонстрации

from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.types import *

spark = SparkSession.builder \
    .appName("Partitioning-Practice") \
    .config("spark.sql.adaptive.enabled", "true") \
    .getOrCreate()

# Генерируем события: 10 млн строк, 3 года данных
events_df = (
    spark.range(1, 10_000_001)
    .withColumn("event_ts",
        F.from_unixtime(
            1_672_531_200 + (F.col("id") % (365 * 24 * 3600 * 3)).cast("long")
        ))
    .withColumn("event_date",  F.to_date("event_ts"))
    .withColumn("user_id",     (F.col("id") % 1_000_000).cast("long"))
    .withColumn("country",
        F.when(F.col("id") % 10 < 5, "US")   # 50% US → намеренный skew
        .when(F.col("id") % 10 < 7, "IN")   # 20% India
        .when(F.col("id") % 10 < 8, "BR")   # 10% Brazil
        .when(F.col("id") % 10 < 9, "RU")   # 10% Russia
        .otherwise("OTHER"))                  # 10% всё остальное
    .withColumn("amount",  (F.rand(seed=42) * 1000 + 1).cast("double"))
    .withColumn("status",
        F.when(F.col("id") % 20 == 0, "cancelled").otherwise("completed"))
    .drop("id")
)

Демонстрация 1: Правильное партиционирование по дате

import time

# ─── Запись с правильным партиционированием ───────────────────────────────────
t0 = time.time()
events_df.write \
    .mode("overwrite") \
    .partitionBy("event_date") \
    .parquet("/tmp/events_by_date/")
print(f"Запись заняла: {time.time() - t0:.1f}s")

# ─── Чтение с partition pruning: только один день из 3 лет ───────────────────
events_good = spark.read.parquet("/tmp/events_by_date/")

t1 = time.time()
result_good = events_good \
    .filter(F.col("event_date") == "2026-01-15") \
    .agg(F.sum("amount").alias("total"), F.count("*").alias("cnt"))
result_good.show()
print(f"Запрос за один день: {time.time() - t1:.3f}s")

# Проверить план: должен быть PartitionFilters
events_good.filter(F.col("event_date") == "2026-01-15").explain(mode="formatted")
# == Physical Plan ==
# FileScan parquet [event_ts, user_id, country, amount, status, event_date]
#   PartitionFilters: [isnotnull(event_date), (event_date = 18277)]   ← PRUNING!
#   PushedFilters: []
#   ReadSchema: ...
# Spark читает ОДИН файл из ~1095 партиций

Демонстрация 2: Катастрофа с high-cardinality ключом

# ВНИМАНИЕ: не запускайте это на production с миллионами user_id!
# Для demo ограничим до 10 000 уникальных user_id

demo_df = events_df.withColumn(
    "user_id_demo", (F.col("user_id") % 10000).cast("long")
)

# ─── Запись: 10 000 партиций по user_id ──────────────────────────────────────
t0 = time.time()
demo_df.write \
    .mode("overwrite") \
    .partitionBy("user_id_demo") \     # 10 000 уникальных значений!
    .parquet("/tmp/events_by_user/")
print(f"Запись: {time.time() - t0:.1f}s")

# Посмотрим сколько файлов создалось
import subprocess
result = subprocess.run(
    ["find", "/tmp/events_by_user/", "-name", "*.parquet", "-type", "f"],
    capture_output=True, text=True
)
file_count = len(result.stdout.strip().split("\n"))
print(f"Файлов создано: {file_count}")  # → ~10 000-20 000 мелких файлов

# ─── Чтение: листинг 10 000 директорий ───────────────────────────────────────
t1 = time.time()
events_user = spark.read.parquet("/tmp/events_by_user/")
events_user.filter(F.col("event_date") == "2026-01-15").count()
print(f"Запрос (с листингом): {time.time() - t1:.3f}s")
# Значительно медленнее чем первый пример, хотя данных то же количество

# ─── Проверить размер файлов: крошечные = проблема ───────────────────────────
import os
sizes = []
for root, dirs, files in os.walk("/tmp/events_by_user/"):
    for f in files:
        if f.endswith(".parquet"):
            sizes.append(os.path.getsize(os.path.join(root, f)))

if sizes:
    avg_kb = sum(sizes) / len(sizes) / 1024
    print(f"Среднй размер файла: {avg_kb:.1f} КБ")  # → ~50-100 КБ вместо 128-256 МБ!

Метрики здоровья партиционирования

Золотой стандарт размера файла в партиции: 128 МБ - 512 МБ на файл в формате Parquet/ORC. Меньше - overhead на открытие файлов и metadata; больше - теряется параллелизм.

def analyze_partition_health(path, spark):
    """Анализ здоровья партиционированной таблицы."""
    df = spark.read.parquet(path)

    # Число партиций (Spark-партиций при чтении)
    num_partitions = df.rdd.getNumPartitions()

    # Размер каждой партиции в МБ
    partition_sizes = df.rdd.mapPartitions(
        lambda it: [sum(1 for _ in it)]
    ).collect()

    # Подсчёт строк на партицию
    row_counts = (
        df.groupBy(*[c for c in df.columns if "date" in c or "dt" in c])
        .count()
        .orderBy("count", ascending=False)
    )

    return {
        "num_spark_partitions": num_partitions,
        "max_rows_per_partition": max(partition_sizes) if partition_sizes else 0,
        "min_rows_per_partition": min(partition_sizes) if partition_sizes else 0,
        "skew_ratio": max(partition_sizes) / max(min(partition_sizes), 1)
    }

health = analyze_partition_health("/tmp/events_by_date/", spark)
print(health)
# {'num_spark_partitions': 1095, 'max_rows_per_partition': 9185,
#  'min_rows_per_partition': 9153, 'skew_ratio': 1.003}  ← идеальный баланс!

Контроль числа файлов внутри партиции

Даже при правильном ключе партиционирования число файлов внутри одной партиции может быть избыточным - если, например, записи шли маленькими батчами или Spark создавал по одному файлу на Task.

# Проблема: 200 мелких файлов в одной партиции event_date=2026-01-15
# Решение 1: repartition перед записью
events_df \
    .filter(F.col("event_date") == "2026-01-15") \
    .repartition(4) \                    # контролируем число файлов: 4
    .write \
    .mode("overwrite") \
    .parquet("/tmp/events_by_date_clean/event_date=2026-01-15/")

# Решение 2: coalesce (без shuffle, только уменьшение числа файлов)
events_df \
    .partitionBy("event_date") \
    .coalesce(2) \                       # максимум 2 файла на партицию
    .write.parquet(path)
# Осторожно: coalesce может дать несбалансированные партиции!

# Решение 3: правильный расчёт при записи всей таблицы
# Целевой размер файла: 256 МБ
# Объём одной партиции event_date: 500 МБ → нужно 2 файла
# Объём задаётся через spark.sql.files.maxRecordsPerFile или spark.sql.shuffle.partitions

spark.conf.set("spark.sql.files.maxRecordsPerFile", "1000000")  # ~128 МБ на файл

Динамическая перезапись партиций

В продакшн-пайплайнах данные часто перезаписываются инкрементально: обновить только вчерашнюю партицию, не трогая исторические данные.

# ─── Проблема: статический overwrite стирает ВСЮ таблицу ─────────────────────
events_df.write \
    .mode("overwrite") \       # ← стирает s3a://data/events/ целиком!
    .partitionBy("event_date") \
    .parquet("s3a://data/events/")

# ─── Решение: Dynamic Partition Overwrite ────────────────────────────────────
# Перезаписывает ТОЛЬКО те партиции, которые есть в записываемом DataFrame
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")

yesterday_events.write \
    .mode("overwrite") \       # теперь только event_date=2026-01-14 перезапишется
    .partitionBy("event_date") \
    .parquet("s3a://data/events/")
# Остальные партиции (2026-01-01 ... 2026-01-13) не тронуты!

# ─── Альтернатива: явное удаление партиции перед записью ─────────────────────
import shutil
shutil.rmtree("/data/events/event_date=2026-01-14/", ignore_errors=True)
yesterday_events.write \
    .mode("append") \           # дописываем в теперь пустую партицию
    .partitionBy("event_date") \
    .parquet("/data/events/")

5. Partition Pruning в explain-плане

Умение читать explain() применительно к партиционированию - необходимый навык для диагностики и оптимизации.

events = spark.read.parquet("/tmp/events_by_date/")

# ─── Запрос с pruning ─────────────────────────────────────────────────────────
q1 = events.filter(F.col("event_date") == "2026-01-15").select("amount", "country")
q1.explain(mode="formatted")

# == Physical Plan ==
# Project [amount#5, country#3]
# +- Filter (isnotnull(event_date#6) AND (event_date#6 = 18277))
#    +- FileScan parquet [country#3,amount#5,event_date#6] Batched: true, DataFilters: [],
#       Format: Parquet, Location: InMemoryFileIndex(1 paths)[file:/tmp/events_by_date],
#       PartitionFilters: [isnotnull(event_date#6), (event_date#6 = 18277)],
#       PushedFilters: [], ReadSchema: struct<country:string,amount:double>
#
# PartitionFilters: [isnotnull(event_date#6), (event_date#6 = 18277)]
# ↑ Spark знает, что читать только одну партицию

# ─── Запрос БЕЗ pruning: нет фильтра по partition-колонке ────────────────────
q2 = events.filter(F.col("amount") > 500).select("amount", "country")
q2.explain(mode="formatted")

# == Physical Plan ==
# Project [amount#5, country#3]
# +- Filter (isnotnull(amount#5) AND (amount#5 > 500.0))
#    +- FileScan parquet [country#3,amount#5,event_date#6] Batched: true,
#       Format: Parquet, Location: InMemoryFileIndex(1095 paths)[...1094 more paths],
#       PartitionFilters: [],          ← ПУСТО: pruning нет
#       PushedFilters: [IsNotNull(amount), GreaterThan(amount,500.0)],
#       ReadSchema: struct<country:string,amount:double>
#
# Spark читает ВСЕ 1095 партиций и фильтрует в памяти
# PushedFilters сработают внутри Parquet (row-group level), но не на уровне партиций

Ключевые маркеры в explain для партиционирования:

  • PartitionFilters: [...] - какие условия применяются к partition-ключам для отсечения папок. Если пусто - нет pruning.
  • Location: InMemoryFileIndex(N paths) - N показывает сколько файлов Spark планирует читать. Чем меньше N по сравнению с общим числом файлов - тем лучше pruning.
  • PushedFilters - фильтры, которые проталкиваются в Parquet (row-group level), но это уже после pruning.

6. Modern Lakehouse: Hidden Partitioning и Liquid Clustering

Ограничения классического Hive-style Partitioning

Классическое партиционирование имеет фундаментальный недостаток: пользователь должен знать о физической организации данных при написании запросов. Если таблица партиционирована по event_date, то фильтр WHERE event_ts > '2026-01-01' (по timestamp, не по date) не воспользуется pruning - Spark не знает о связи между event_ts и event_date.

Ещё одна проблема: изменить ключ партиционирования после записи данных - значит переписать всю таблицу.

Apache Iceberg: Hidden Partitioning

Iceberg решает первую проблему через Hidden Partitioning: физическое распределение данных скрыто от пользователя. Пользователь фильтрует по исходной колонке (event_ts), а Iceberg автоматически транслирует это в partition-фильтр.

# Создание Iceberg таблицы с hidden partitioning
spark.sql("""
    CREATE TABLE iceberg.events (
        event_ts    TIMESTAMP,
        user_id     BIGINT,
        country     STRING,
        amount      DOUBLE
    )
    USING iceberg
    PARTITIONED BY (days(event_ts), country)
    -- days() - это transform: Iceberg сам вычислит partition-ключ из event_ts
    -- Физически: данные лежат в 2026-01-15/RU/ и т.д.
    -- Но пользователь фильтрует по event_ts напрямую!
""")

# Запись данных
events_df.write \
    .format("iceberg") \
    .mode("append") \
    .save("iceberg.events")

# Запрос: пользователь фильтрует по event_ts (не по partition-ключу!)
spark.sql("""
    SELECT country, sum(amount) AS total
    FROM iceberg.events
    WHERE event_ts >= '2026-01-01' AND event_ts < '2026-02-01'
    GROUP BY country
""")
# Iceberg автоматически транслирует event_ts-фильтр в days(event_ts)-partition-фильтр
# Читаются только партиции за январь 2026

Доступные трансформации в Iceberg:

  • years(col) - по году
  • months(col) - по году-месяцу
  • days(col) - по дате
  • hours(col) - по часу
  • bucket(N, col) - хэш-бакетизация (N бакетов)
  • truncate(N, col) - обрезание строки или числа до N символов/разрядов
  • identity(col) - то же что классическое partitionBy

Partition Evolution в Iceberg

Iceberg позволяет менять схему партиционирования без переписывания данных. Старые данные остаются в старом layout, новые пишутся в новый. Iceberg корректно обрабатывает оба набора при чтении.

# Изначально: партиционирование по месяцам
spark.sql("""
    ALTER TABLE iceberg.events
    REPLACE PARTITION FIELD months(event_ts) WITH days(event_ts)
""")
# Данные до изменения: лежат в month-партициях → Iceberg читает их через manifest
# Данные после изменения: пишутся в day-партиции
# Пользователь не замечает разницы - запросы работают прозрачно

Delta Lake: Liquid Clustering

Delta Lake 3.1+ представил Liquid Clustering - адаптивную замену классическому partitionBy. Вместо жёсткой иерархии директорий Liquid Clustering организует данные внутри файлов таким образом, чтобы похожие значения кластерных колонок лежали рядом. Автоматически применяет Data Skipping через Transaction Log.

# Создание таблицы с Liquid Clustering
spark.sql("""
    CREATE TABLE delta.events
    USING DELTA
    CLUSTER BY (event_date, country)
    AS SELECT * FROM events_df
""")
-- Нет жёсткого partitionBy!
-- Delta сам решает как физически организовать файлы

# Периодическая оптимизация кластеризации
spark.sql("OPTIMIZE delta.events")
# OPTIMIZE переупорядочивает файлы для лучшего Data Skipping

# Проверить эффективность: посмотреть сколько файлов пропустилось
spark.sql("""
    SELECT * FROM delta.events WHERE event_date = '2026-01-15' AND country = 'RU'
""")
# В Spark UI / Delta logs: numFilesSkipped покажет сколько файлов пропущено

7. Best Practices и чек-лист

Алгоритм выбора partition-ключа

Прежде чем вызвать partitionBy(), пройдите по этому алгоритму:

# Шаг 1: узнать паттерны запросов
# Какие колонки чаще всего в WHERE? Спросите аналитиков, посмотрите Query History.

# Шаг 2: проверить кардинальность кандидатов
df.select(
    F.countDistinct("event_date").alias("dates"),
    F.countDistinct("country").alias("countries"),
    F.countDistinct("region").alias("regions"),
).show()
# Цель: 50 – 10 000 уникальных значений

# Шаг 3: проверить распределение
df.groupBy("country").count() \
  .withColumn("pct", F.round(F.col("count") / df.count() * 100, 2)) \
  .filter(F.col("pct") > 20) \
  .show()
# Если есть значения с pct > 20-30% → риск skew

# Шаг 4: оценить итоговое число партиций
# num_partitions = cardinality_key1 × cardinality_key2 × ...
# Цель: не более 10 000 – 50 000 партиций

# Шаг 5: оценить средний размер партиции
total_gb = 500  # ГБ
num_partitions = 365  # дней
avg_partition_gb = total_gb / num_partitions  # ~1.37 ГБ на партицию
# Цель: 0.5 – 5 ГБ на партицию (128 МБ – 512 МБ на файл внутри партиции)

Сводная таблица: что выбрать

Сценарий Рекомендация
Аналитическая таблица событий partitionBy("event_date")
Глобальная таблица с регионами partitionBy("event_date", "region") если регионов < 200
Таблица с high-cardinality ключом join partitionBy("event_date") + bucketBy(N, "user_id")
Частые точечные запросы по user_id Bucketing, не partitioning
Multi-tenant SaaS с неравными клиентами partitionBy("event_date"), tenant через Bucketing
Lakehouse с частыми обновлениями Delta Liquid Clustering
Разные паттерны запросов (WHERE по разным колонкам) Iceberg Hidden Partitioning
Маленькая таблица (< 1 ГБ) Не партиционировать вообще

Признаки нездоровья партиционирования

Симптом Вероятная причина Лечение
Spark Driver OOM при запуске Слишком много файлов (high-cardinality ключ) Снизить кардинальность ключа
Один Task занимает в 10× дольше других Data Skew в партициях Добавить второй ключ или salting
Directory Listing занимает > 10 секунд Слишком много директорий Уменьшить глубину вложенности
PartitionFilters: [] в explain Запрос не использует partition-ключ Добавить фильтр по partition-колонке
Тысячи мелких файлов per партиция Много маленьких batch-записей Compaction (OPTIMIZE) или coalesce при записи
Одна партиция в 100× больше остальных Structural skew Переработать модель партиционирования

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

Дан следующий дизайн таблицы в продакшн-системе:

# Таблица: логи API-запросов, 5 млрд строк в год
api_logs.write \
    .partitionBy("year", "month", "day", "hour", "endpoint_id") \
    .parquet("s3a://data/api_logs/")

# endpoint_id: 50 000 уникальных эндпоинтов
# Запросы аналитиков:
# SELECT count(*) FROM api_logs WHERE day = '2026-01-15' AND endpoint_id = 12345
# SELECT avg(latency_ms) FROM api_logs WHERE month = '2026-01' AND status_code >= 500

Задание 1 - Посчитать катастрофу. Рассчитайте максимально возможное число директорий при 3 годах данных: year × month × day × hour × endpoint_id. Сравните с рекомендуемым максимумом 10 000–50 000.

Задание 2 - Предложить исправление. Переработайте дизайн партиционирования: выберите подходящие ключи с учётом паттернов запросов и кардинальности endpoint_id. Если endpoint_id нужен как фильтр, предложите альтернативу партиционированию (Bucketing? Predicate Pushdown через Parquet statistics?).

Задание 3 - Проверить через explain. Напишите код, который:

  • Записывает тестовый DataFrame с вашим новым дизайном партиционирования
  • Читает его и запускает explain(mode="formatted") для обоих типичных запросов из задания
  • Находит и интерпретирует строки PartitionFilters и Location: InMemoryFileIndex(N paths)

Задание 4 - Modern Lakehouse. Перепишите решение из Задания 2 используя Iceberg Hidden Partitioning с трансформацией hours(request_ts). Объясните: почему пользователю не нужно знать о физической организации данных при написании запроса WHERE request_ts >= '2026-01-15 10:00:00'?


В следующем уроке разберём проблему маленьких файлов подробнее: откуда она берётся в streaming и batch-пайплайнах, как она убивает производительность и какие инструменты - от OPTIMIZE в Delta до специфических настроек Spark - позволяют с ней справляться.