Storage Partitioning: выбор колонок для partitionBy и анти-паттерны
Физика partitionBy в файловых хранилищах, цена Directory Listing в S3 и HDFS, правила выбора ключа партиционирования, три главных анти-паттерна и как их диагностировать, Hidden Partitioning в Iceberg и Liquid Clustering в Delta Lake.
В предыдущем уроке мы разобрали 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 - позволяют с ней справляться.