Parquet на S3: vectorized reader, column pruning и row group filtering
Физическая анатомия Parquet: как column pruning, predicate pushdown и vectorized reader позволяют Spark читать только нужные байты из S3 и экономить сетевой трафик в 10–100 раз.
Формат хранения - это половина производительности¶
В традиционном on-premise Hadoop диски были быстрыми (HDFS на локальных HDD/SSD), задержки сети - малыми. Можно было себе позволить читать «лишние» данные: накладные расходы были небольшими.
В облаке всё иначе. Каждый байт, прочитанный из S3 - это HTTP-запрос через публичную сеть. Латентность GET-запроса к S3 - 5–50 мс. Пропускная способность - ограничена и тарифицируется. Чтение 1 ТБ из S3, когда реально нужно 10 ГБ, означает:
- в 100 раз больше сетевого трафика - сотни долларов за egress;
- в 100 раз больше времени - запрос на 10 минут вместо 6 секунд;
- в 100 раз больше CPU на декомпрессию ненужных данных.
Именно поэтому физический формат хранения становится ключевым архитектурным решением при работе со Spark + S3. Parquet - это не просто «сжатый CSV». Это формат, спроектированный с единственной целью: читать из хранилища ровно столько, сколько нужно для ответа на запрос, и ни байтом больше.
Этот урок - детальный разбор физической механики Parquet на S3: как он устроен изнутри, как Spark использует эту структуру для минимизации IO, и как неправильное использование превращает эти преимущества в ничто.
Физическая анатомия Parquet-файла¶
Row-oriented vs Column-oriented хранение¶
Прежде чем разбирать Parquet, нужно понять фундаментальный принцип колоночного хранения.
Row-oriented хранение (CSV, JSON, AVRO): данные записаны строка за строкой. Чтобы прочитать поле amount из миллиона строк, нужно прочесть все строки целиком - включая user_id, timestamp, country, product_name и ещё 196 полей.
Column-oriented хранение (Parquet, ORC): данные одной колонки хранятся вместе. Чтобы прочитать amount из миллиона строк - читаем только блок данных колонки amount. Остальные 199 колонок не трогаем.
Аналитические запросы - это всегда «выбрать несколько колонок из миллиардов строк». Для них колоночное хранение даёт принципиальное преимущество.
Физическая структура Parquet-файла¶
Parquet-файл - это не монолитный поток данных. Он состоит из строго структурированных частей:
┌─────────────────────────────────────────────────────────────────┐
│ Magic Number: PAR1 (4 байта) │
├─────────────────────────────────────────────────────────────────┤
│ ROW GROUP 1 (128–512 MB по умолчанию) │
│ ┌──────────────────────────────────────────────────────────┐ │
│ │ Column Chunk: user_id │ │
│ │ ┌────────────┐ ┌────────────┐ ┌────────────┐ │ │
│ │ │ Page 1 │ │ Page 2 │ │ Page 3 │ ... │ │
│ │ │ (8-16KB) │ │ (8-16KB) │ │ (8-16KB) │ │ │
│ │ └────────────┘ └────────────┘ └────────────┘ │ │
│ │ Column Chunk: amount │ │
│ │ ┌────────────┐ ┌────────────┐ ... │ │
│ │ └────────────┘ └────────────┘ │ │
│ │ Column Chunk: event_date │ │
│ │ ┌────────────┐ ... │ │
│ └──────────────────────────────────────────────────────────┘ │
├─────────────────────────────────────────────────────────────────┤
│ ROW GROUP 2 │
│ ... │
├─────────────────────────────────────────────────────────────────┤
│ ROW GROUP N │
│ ... │
├─────────────────────────────────────────────────────────────────┤
│ FOOTER (переменный размер, обычно КБ–МБ) │
│ ┌──────────────────────────────────────────────────────────┐ │
│ │ Схема данных (Schema) │ │
│ │ Для каждой Row Group: │ │
│ │ Для каждой колонки: │ │
│ │ - byte_offset: где в файле начинается Column Chunk │ │
│ │ - total_compressed_size │ │
│ │ - statistics: min/max/null_count │ │
│ └──────────────────────────────────────────────────────────┘ │
│ Footer Length (4 байта) │
│ Magic Number: PAR1 (4 байта) │
└─────────────────────────────────────────────────────────────────┘
Разберём каждую часть:
Row Group - горизонтальный срез данных, содержащий N строк. Размер по умолчанию в Spark - 128 МБ (настраивается через parquet.block.size). Это ключевая единица параллелизма: каждый Spark-таск обрабатывает одну Row Group. Также это ключевая единица для row group filtering: если Row Group не проходит фильтр → вся она пропускается без чтения данных.
Column Chunk - все данные одной колонки в рамках одной Row Group. Физически хранится непрерывным блоком байт. Именно это позволяет Spark сделать один HTTP Range Request к S3 для чтения одной колонки.
Page - минимальная единица сжатия и кодирования внутри Column Chunk. Размер страницы - 8–16 КБ по умолчанию. Spark не может читать «половину страницы» - это минимальная единица IO на уровне Parquet.
Footer - критически важная часть. Это карта всего файла: где находится каждый Column Chunk, его размер, статистика (min/max/null count). Без Footer Spark не может знать, как читать файл.
Магия Footer: почему Spark всегда читает конец файла первым¶
Когда Spark открывает Parquet-файл из S3, он не начинает читать с начала. Первый HTTP-запрос всегда направлен в конец файла - туда, где находится Footer.
Почему это важно: из файла в 10 ГБ с 200 колонками и 20 Row Groups Spark может прочитать всего несколько сотен МБ, ответив на конкретный аналитический запрос. Это возможно только потому, что Footer содержит точные byte offsets каждого Column Chunk.
Footer хранится в конце файла по историческим причинам: при записи Parquet-файла размер каждой части неизвестен заранее (данные компрессируются, размер зависит от содержимого). Поэтому сначала пишутся данные, потом - Footer с точными offsets.
Column Pruning: читаем только нужные колонки¶
Что происходит без Column Pruning¶
Таблица логов событий: 200 колонок (user_id, session_id, event_type, timestamp, client_ip, user_agent, country, city, device_type, OS, browser, app_version, ... и ещё 190 полей). Один файл - 10 ГБ.
Запрос аналитика:
SELECT user_id, SUM(amount) as total
FROM events
WHERE event_date = '2024-06-15'
GROUP BY user_id
Без Column Pruning - Spark читает все 200 колонок → 10 ГБ из S3.
С Column Pruning - Spark читает только user_id, amount, event_date → ~150 МБ из S3.
67-кратное сокращение трафика из одного только Column Pruning.
Как работает Column Pruning в Spark¶
Spark Catalyst Optimizer анализирует запрос и определяет минимальный набор колонок, необходимых для его выполнения. Этот набор передаётся в scan-оператор, который при чтении Parquet из Footer извлекает byte offsets только нужных Column Chunks и делает HTTP Range Requests только для них.
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
spark = SparkSession.builder.appName("parquet-demo").getOrCreate()
# Читаем таблицу с 200 колонками (10 GB):
events = spark.read.parquet("s3a://data-lake/events/")
# ❌ SELECT * - Column Pruning не работает, читаем все 200 колонок:
events.select("*") \
.filter(F.col("event_date") == "2024-06-15") \
.groupBy("user_id") \
.agg(F.sum("amount")) \
.show()
# ✅ Точечный SELECT - Column Pruning читает 3 из 200 колонок:
events.select("user_id", "amount", "event_date") \
.filter(F.col("event_date") == "2024-06-15") \
.groupBy("user_id") \
.agg(F.sum("amount")) \
.show()
# ✅ Spark автоматически применяет Column Pruning даже без explicit select:
# Catalyst видит, что groupBy/agg использует только user_id и amount,
# а filter использует только event_date → автоматически запрашивает 3 колонки
events \
.filter(F.col("event_date") == "2024-06-15") \
.groupBy("user_id") \
.agg(F.sum("amount")) \
.show()
Проверка Column Pruning через Explain Plan¶
Лучший способ убедиться, что Column Pruning работает - посмотреть физический план:
# Посмотреть Physical Plan:
events \
.filter(F.col("event_date") == "2024-06-15") \
.groupBy("user_id") \
.agg(F.sum("amount")) \
.explain(extended=True)
# В выводе ищем ReadSchema в операторе FileScan:
# == Physical Plan ==
# ...
# +- *(1) FileScan parquet [user_id#0L, amount#1, event_date#2]
# Batched: true
# Location: InMemoryFileIndex(...)
# PushedFilters: [IsNotNull(event_date), EqualTo(event_date,2024-06-15)]
# ReadSchema: struct<user_id:bigint, amount:double, event_date:date>
#
# ReadSchema содержит только 3 колонки из 200 - Column Pruning работает!
ReadSchema в выводе - это ключевой индикатор. Если в нём только нужные колонки, Column Pruning применён. Если там struct<user_id:bigint, name:string, country:string, ... (200 полей)> - Column Pruning отсутствует.
Nested Column Pruning для STRUCT-полей¶
Parquet поддерживает вложенные структуры (STRUCT, ARRAY, MAP). Column Pruning работает и на уровне вложенных полей:
# Схема: events.metadata STRUCT<device: STRUCT<type: STRING, os: STRING>, app_version: STRING>
# Нужно только metadata.device.type:
events.select(
"user_id",
F.col("metadata.device.type").alias("device_type")
).show()
# Physical Plan покажет:
# ReadSchema: struct<user_id:bigint, metadata:struct<device:struct<type:string>>>
# Spark читает только поле type внутри device внутри metadata - не весь STRUCT!
Антипаттерн: передача DataFrame через UDF теряет Column Pruning¶
from pyspark.sql.functions import udf
from pyspark.sql.types import DoubleType
# ❌ UDF принимает весь Row объект - Column Pruning отключается!
@udf(returnType=DoubleType())
def complex_calc(row):
return row.amount * 1.1
# Spark не знает, какие поля использует UDF → читает все колонки
events.withColumn("adjusted_amount", complex_calc(F.struct("*"))).show()
# ✅ UDF получает только нужные колонки:
@udf(returnType=DoubleType())
def calc_adjusted_amount(amount: float) -> float:
return amount * 1.1
# Spark знает: UDF нужна только колонка amount
events.withColumn(
"adjusted_amount",
calc_adjusted_amount(F.col("amount"))
).select("user_id", "adjusted_amount").show()
Row Group Filtering: пропускаем целые блоки данных¶
Статистика в Footer: min/max/null_count¶
Каждый Column Chunk в Footer содержит статистику для своей Row Group:
min: минимальное значение колонки в данной Row Group.max: максимальное значение.null_count: количество NULL.distinct_count(опционально): количество уникальных значений (включается отдельно).
Эта статистика записывается при создании Parquet-файла и не меняется. Именно она позволяет Spark принимать решение «читать или пропустить» Row Group, не загружая данные.
В примере выше из 4 Row Groups Spark читает только одну (25% данных). На практике при правильной физической организации данных можно пропускать 90–99% файла.
Predicate Pushdown: как фильтры передаются в scan¶
Predicate Pushdown - это оптимизация Catalyst Optimizer, при которой условие WHERE «проталкивается» как можно глубже в план выполнения - вплоть до оператора чтения файла.
# Запрос с фильтром:
result = events \
.filter(F.col("event_date") == "2024-06-15") \
.filter(F.col("amount") > 100) \
.select("user_id", "amount")
result.explain()
# Physical Plan (ключевые строки):
# +- *(1) Project [user_id#0L, amount#1]
# +- *(1) Filter (isnotnull(amount#1) AND (amount#1 > 100.0))
# +- *(1) FileScan parquet [user_id#0L, amount#1, event_date#2]
# Batched: true
# PushedFilters: [IsNotNull(event_date), EqualTo(event_date,2024-06-15),
# IsNotNull(amount), GreaterThan(amount,100.0)]
# ReadSchema: struct<user_id:bigint, amount:double, event_date:date>
PushedFilters - это именно те условия, которые переданы в Parquet Reader. Он использует их для Row Group Filtering через статистику Footer.
Обратите внимание: не все фильтры можно «протолкнуть». Фильтры, применённые после трансформации (withColumn, join), применяются уже к данным в памяти и не помогают пропускать Row Groups.
Почему сортировка данных критически важна для Row Group Filtering¶
Row Group Filtering эффективен только когда диапазоны min/max не перекрываются между разными Row Groups. Это возможно только если данные отсортированы по колонке фильтрации.
# ❌ Несортированные данные: фильтр event_date = '2024-06-15' неэффективен
# Каждая Row Group содержит случайный микс дат:
# RG1: min=2024-01-01, max=2024-12-31 (почти весь год)
# RG2: min=2024-01-02, max=2024-12-30
# RG3: min=2024-01-01, max=2024-12-31
# → Все Row Groups «могут содержать нужные строки» → читаем всё
# ✅ Отсортированные данные: каждая Row Group = узкий диапазон дат
events_sorted = events.sortWithinPartitions("event_date")
# RG1: min=2024-01-01, max=2024-01-31 (январь)
# RG2: min=2024-02-01, max=2024-02-28 (февраль)
# RG3: ...
# RG6: min=2024-06-01, max=2024-06-30 (июнь) ← только эта RG читается!
# Запись с правильной сортировкой:
events \
.sortWithinPartitions("event_date", "user_id") \
.write \
.option("parquet.block.size", 134217728) \ # 128 MB Row Groups
.parquet("s3a://data-lake/events-sorted/")
Аналогия: представьте библиотеку с книгами. Несортированные книги разбросаны по полкам случайно - чтобы найти все книги про Python, нужно просмотреть каждую полку. Отсортированные по теме - все Python-книги стоят на одной полке, остальные не трогаем.
Статистика страниц (Page Statistics) - более тонкая гранулярность¶
Parquet поддерживает статистику не только на уровне Row Group, но и на уровне страниц (Page). Это позволяет пропускать данные с ещё более тонкой гранулярностью.
# Включить запись page-level statistics:
spark.conf.set("spark.sql.parquet.enablePageIndexing", "true")
# При чтении Spark использует Page Index для additional filtering:
# Вместо «читать весь Column Chunk» → «читать только страницы с нужными данными»
# Это особенно эффективно для частичных фильтров внутри Row Group
Page-level statistics появились в Parquet 2.0 и поддерживаются в Spark 3.2+.
Bloom Filter: точечный поиск по некластеризованным данным¶
Для точечных запросов по высококардинальным колонкам (например, WHERE user_id = 12345) min/max статистика неэффективна - диапазоны слишком широкие. Parquet поддерживает Bloom Filter на уровне Column Chunk - вероятностную структуру данных, позволяющую быстро определить, есть ли значение в Row Group.
# Включить Bloom Filter при записи:
df.write \
.option("parquet.bloom.filter.enabled", "true") \
.option("parquet.bloom.filter.expected.ndv", "1000000") \ # ожидаемое кол-во уникальных значений
.option("parquet.bloom.filter.fpp", "0.05") \ # false positive rate 5%
.parquet("s3a://data-lake/users/")
# Spark использует Bloom Filter автоматически при наличии EqualTo фильтров:
spark.conf.set("spark.sql.parquet.enableBloomFilterPushdown", "true")
# Теперь: WHERE user_id = 12345
# → Spark проверяет Bloom Filter каждой RG (несколько KB из Footer)
# → Если Bloom Filter говорит «нет» → RG пропускается без чтения данных
# → Bloom Filter никогда не даёт false negative, только false positive
False positive rate: если Bloom Filter говорит «значение МОЖЕТ быть», это не значит, что оно есть - это требует чтения данных. Если говорит «значение точно НЕТ» - Row Group пропускается гарантированно.
Vectorized Parquet Reader: скорость на уровне CPU¶
Проблема строковой обработки¶
Классический Parquet Reader читает данные постранично и преобразует каждую строку в объект Row JVM. Для каждой строки:
- Выделить память под объект Row.
- Скопировать значения из декодированных буферов в поля объекта.
- Передать объект в следующий оператор.
- Garbage Collector в конечном счёте освобождает объекты.
При 100 миллионах строк - 100 миллионов JVM-объектов, 100 миллионов аллокаций памяти. GC-паузы, cache misses, медленное выполнение.
Vectorized Reader: батчевая обработка¶
Vectorized Parquet Reader читает данные батчами - группами по 4096 строк (по умолчанию). Вместо объектов Row данные хранятся в Column Vectors - плоских массивах примитивных типов в off-heap памяти:
Batch size: 4096 строк
user_id: int64[4096] - плоский массив int64
amount: float64[4096] - плоский массив float64
is_active: bool[4096] - плоский массив bool (или bit array)
Это даёт несколько преимуществ:
1. CPU Cache Locality: данные одной колонки лежат непрерывно в памяти. CPU cache работает эффективно - префетч следующих элементов происходит автоматически.
2. SIMD-подобные операции: современные CPU (AVX-512) обрабатывают 512 бит за один такт. Операция SUM(amount) над массивом float64 - это буквально 8 float64 за один CPU-такт через векторные инструкции.
3. Нет GC-давления: Column Vectors в off-heap памяти (за пределами JVM heap) - GC их не трогает.
4. Vectorized decoding: декодирование Parquet-кодирований (RLE, DELTA, DICTIONARY) также выполняется батчами, что значительно эффективнее построчной декомпрессии.
Конфигурация Vectorized Reader¶
# Включён по умолчанию в Spark 2.0+:
spark.conf.set("spark.sql.parquet.enableVectorizedReader", "true")
# Размер batch (строк за одну итерацию):
spark.conf.set("spark.sql.parquet.columnarReaderBatchSize", "4096")
# Увеличение до 8192 может ускорить простые запросы
# Уменьшение помогает при очень широких таблицах (много колонок)
# Максимальный размер Column Vector (в памяти):
# При batch=4096 и 200 колонок int64:
# 4096 × 200 × 8 bytes = 6.5 MB per batch - нормально
# Проверить, используется ли Vectorized Reader:
# В Physical Plan ищем "Batched: true":
# +- *(1) FileScan parquet [...]
# Batched: true ← Vectorized Reader активен
# Batched: false ← Row-by-Row Reader (деградация)
Когда Vectorized Reader отключается¶
Vectorized Reader - не серебряная пуля. Он отключается в нескольких ситуациях:
1. Сложные вложенные типы
# ❌ Vectorized Reader отключается для ARRAY/MAP/глубоких STRUCT:
events.select(
"user_id",
F.explode("tags") # tags: ARRAY<STRING>
).show()
# Physical Plan: Batched: false → Row-by-Row Reader
# ✅ Для простых STRUCT верхнего уровня Vector Reader работает:
events.select(
"user_id",
F.col("metadata.device_type") # STRUCT с простым полем
).show()
# Physical Plan: Batched: true (если нет глубокой вложенности)
2. Python UDF
from pyspark.sql.functions import udf
# ❌ Python UDF требует сериализации в Python-объекты → no vectorized:
@udf("double")
def python_calc(x):
return x * 1.1
events.withColumn("adj", python_calc("amount")).show()
# Batched: false - каждая строка сериализуется в Python
# ✅ Pandas UDF (Vectorized UDF) работает с Column Vectors:
from pyspark.sql.functions import pandas_udf
import pandas as pd
@pandas_udf("double")
def pandas_calc(series: pd.Series) -> pd.Series:
return series * 1.1
events.withColumn("adj", pandas_calc("amount")).show()
# Batched: true - работает с pandas Series (векторизовано)
3. Некоторые типы данных
# Decimal с precision > 38 → Row-by-Row
# BINARY columns → Row-by-Row (в некоторых версиях Spark)
# Проверяйте через explain() + ищите Batched: true/false
Размер Row Group: компромисс между параллелизмом и эффективностью¶
Почему размер Row Group важен¶
Размер Row Group влияет на два конкурирующих фактора:
Большие Row Groups → лучший Row Group Filtering (статистика охватывает больше строк → чаще «срабатывает» пропуск) → меньше файлов, меньше Footer-запросов к S3.
Маленькие Row Groups → лучший параллелизм (больше тасков в Spark, каждый таск обрабатывает меньше данных) → быстрее время до первого результата.
Практическое правило для S3-нагрузок:
| Тип нагрузки | Рекомендуемый размер Row Group | Причина |
|---|---|---|
| OLAP (batch аналитика) | 256–512 МБ | Меньше HTTP-запросов, лучший filtering |
| Streaming micro-batch | 64–128 МБ | Баланс параллелизма и filtering |
| Point queries | 64 МБ | Мало данных, Bloom Filter эффективнее |
| ML Feature Store | 256 МБ | Большие батчи для векторизованного чтения |
Настройка Row Group при записи¶
# Размер Row Group в байтах:
spark.conf.set("spark.hadoop.parquet.block.size", str(256 * 1024 * 1024)) # 256 MB
# Или через write options:
df.write \
.option("parquet.block.size", 268435456) \ # 256 MB
.option("parquet.page.size", 1048576) \ # 1 MB страницы
.option("parquet.dictionary.enabled", "true") \ # Dictionary encoding
.parquet("s3a://data-lake/events-optimized/")
# Проверить фактический Row Group layout Parquet-файла:
from pyspark.sql import functions as F
spark.read \
.parquet("s3a://data-lake/events-optimized/part-00000.parquet") \
.select("*") \
.explain()
# Или через parquet-tools CLI:
# parquet-tools meta part-00000.parquet | grep "row group"
Проблема мелких файлов на S3¶
Мелкие файлы - один из главных врагов производительности Parquet на S3. Если Spark генерирует тысячи файлов по несколько МБ, каждый файл:
- Требует отдельного HTTP GET для Footer (5–50 мс).
- Содержит мало Row Groups → Row Group Filtering почти не работает.
- Создаёт тысячи тасков при чтении → overhead на планирование.
# ❌ Плохо: 10 000 файлов по 1 MB
streaming_df.writeStream \
.trigger(processingTime="10 seconds") \ # частые триггеры = много мелких файлов
.parquet("s3a://bucket/events/")
# Посмотреть, сколько файлов:
spark.read.parquet("s3a://bucket/events/").inputFiles.__len__()
# → 10 000 файлов: 10 000 Footer GET запросов при каждом чтении
# ✅ Компакция через Spark:
def compact_parquet(
spark,
input_path: str,
output_path: str,
target_file_size_mb: int = 256,
target_files: int = None
) -> None:
"""
Объединить мелкие Parquet-файлы в крупные.
target_files: None = авто (spark.sql.shuffle.partitions)
"""
df = spark.read.parquet(input_path)
if target_files is None:
# Оценить оптимальное количество файлов:
# total_size / target_file_size_mb
# Spark не знает размер сам, используем shuffle.partitions
target_files = max(1, spark.sparkContext.defaultParallelism)
df.coalesce(target_files) \
.sortWithinPartitions("event_date", "user_id") \ # сортировка для Row Group Filtering
.write \
.option("parquet.block.size", target_file_size_mb * 1024 * 1024) \
.mode("overwrite") \
.parquet(output_path)
# После компакции: 50 файлов по 256 MB
# 50 Footer GET запросов вместо 10 000 = в 200 раз меньше S3 API вызовов
Parquet-оптимизации в Iceberg и Delta Lake¶
Как Lakehouse-форматы усиливают Parquet¶
Delta Lake и Apache Iceberg добавляют дополнительный слой метаданных поверх Parquet, который расширяет возможности data skipping:
Iceberg: Hidden Partitioning и Partition Statistics¶
# Iceberg с hidden partitioning по дате:
spark.sql("""
CREATE TABLE catalog.events (
user_id BIGINT,
event_type STRING,
amount DOUBLE,
created_at TIMESTAMP
)
USING iceberg
PARTITIONED BY (days(created_at)) -- hidden partition
""")
# Запрос с фильтром по времени:
spark.sql("""
SELECT user_id, SUM(amount)
FROM catalog.events
WHERE created_at >= '2024-06-15'
AND created_at < '2024-06-16'
GROUP BY user_id
""")
# Iceberg делает:
# 1. Читает Manifest List (1 HTTP GET)
# 2. Прочитает только Manifests для partition days=19889 (2024-06-15)
# 3. В пределах манифеста: проверяет min/max статистику каждого файла
# 4. Для выбранных файлов: Row Group Filtering из Parquet Footer
# → 3 уровня pruning вместо 1 (только Parquet Footer)
Delta Lake: Z-ordering для многомерного Row Group Filtering¶
# Z-ordering данных по нескольким колонкам:
spark.sql("""
OPTIMIZE delta.`s3a://data-lake/events`
ZORDER BY (user_id, event_date)
""")
# После Z-ordering:
# Данные кластеризованы по (user_id, event_date) одновременно
# Запрос WHERE user_id = 12345 AND event_date = '2024-06-15':
# → Row Groups с нужным user_id И нужной датой сгруппированы вместе
# → Большинство Row Groups пропускается по min/max обоих полей
# → Работает для ЛЮБОГО подмножества Z-ordered колонок
# Аналог в Iceberg через Sort Order:
spark.sql("""
ALTER TABLE catalog.events
WRITE ORDERED BY user_id, event_date
""")
# Или через rewrite:
spark.sql("""
CALL spark_catalog.system.rewrite_data_files(
table => 'catalog.events',
strategy => 'sort',
sort_order => 'user_id ASC, event_date ASC'
)
""")
Anti-patterns: как убить производительность Parquet¶
Антипаттерн 1: SELECT * в аналитических запросах¶
# ❌ SELECT * отменяет Column Pruning:
events.select("*") \
.filter("event_date = '2024-06-15'") \
.count()
# Читает все 200 колонок из S3 → 10 GB трафика
# ✅ Точечный SELECT:
events.select("user_id", "amount", "event_date") \
.filter("event_date = '2024-06-15'") \
.count()
# Читает 3 колонки → ~150 MB трафика (67× меньше)
Антипаттерн 2: Фильтры после join без предварительной фильтрации¶
# ❌ Чтение обеих таблиц полностью, потом join и фильтрация:
orders = spark.read.parquet("s3a://data-lake/orders/") # 500 GB
customers = spark.read.parquet("s3a://data-lake/customers/") # 50 GB
result = orders.join(customers, "customer_id") \
.filter("orders.order_date = '2024-06-15'") # фильтр применяется к результату join
# ✅ Предварительная фильтрация перед join:
orders_filtered = orders.filter("order_date = '2024-06-15'") # Row Group Filtering работает
# → читает только нужные Row Groups: 5 GB вместо 500 GB
result = orders_filtered.join(customers, "customer_id")
Антипаттерн 3: Фильтр по вычисляемому полю¶
# ❌ Фильтр по производному полю - Predicate Pushdown НЕ работает:
events \
.withColumn("year", F.year("event_date")) \
.filter(F.col("year") == 2024)
# PushedFilters: [] - пустой! Читается ВСЁ, потом фильтрация в памяти
# ✅ Фильтр по оригинальному полю:
events \
.filter((F.col("event_date") >= "2024-01-01") & (F.col("event_date") < "2025-01-01"))
# PushedFilters: [GreaterThanOrEqual(event_date,2024-01-01), LessThan(event_date,2025-01-01)]
# ✅ Или через SQL (Catalyst обрабатывает умнее):
events.createOrReplaceTempView("events")
spark.sql("SELECT * FROM events WHERE YEAR(event_date) = 2024")
# Catalyst может преобразовать YEAR(...) = 2024 в диапазонный фильтр
Антипаттерн 4: Высококардинальные партиции на S3¶
# ❌ Партиционирование по user_id (миллионы уникальных значений):
events.write \
.partitionBy("user_id") \ # 10М уникальных user_id = 10М директорий!
.parquet("s3a://data-lake/events/")
# Результат:
# - 10М партиций = 10М ListObjects запросов для планирования
# - Каждая партиция: 1 крошечный файл в несколько KB
# - Spark Driver OutOfMemory на этапе листинга директорий
# ✅ Партиционирование по низкокардинальному полю + Row Group Filtering для user_id:
events.write \
.partitionBy("event_date") \ # ~365 партиций за год
.sortWithinPartitions("user_id") \ # Сортировка для Row Group Filtering
.parquet("s3a://data-lake/events/")
# Запрос WHERE event_date='2024-06-15' AND user_id = 12345:
# 1. Partition pruning: читаем только event_date=2024-06-15 (1 из 365 директорий)
# 2. Row Group Filtering: пропускаем Row Groups где max(user_id) < 12345 или min > 12345
# 3. Bloom Filter (если включён): точечный поиск user_id внутри RG
Антипаттерн 5: Отключение словарного кодирования¶
# ❌ Отключение Dictionary Encoding для строковых колонок:
df.write \
.option("parquet.enable.dictionary", "false") \ # Отключили dictionary!
.parquet("s3a://data-lake/events/")
# Без dictionary каждое значение кодируется как raw string
# "purchase" (8 bytes) × 1M строк = 8 MB
# С dictionary: создаётся словарь {"purchase": 0, "view": 1, ...},
# каждое значение = 1 byte index → 1 MB (8× сжатие)
# ✅ Dictionary Encoding включён по умолчанию:
df.write \
.option("parquet.enable.dictionary", "true") \ # по умолчанию true
.option("parquet.dictionary.page.size", "1048576") \ # 1 MB словарь
.parquet("s3a://data-lake/events/")
Production-кейс: оптимизация Parquet-запроса в 50 раз¶
Рассмотрим реальный сценарий: аналитический запрос, который выполняется 45 минут. После поэтапной оптимизации - 54 секунды.
Исходный состав¶
Таблица: s3a://data-lake/events/
Размер: 2 TB, 10 000 файлов, 200 колонок
Формат: Parquet, файлы ~200 MB, Row Groups ~128 MB
Данные: НЕ отсортированы, партиций нет
# Запрос ДО оптимизации (45 минут):
spark.read.parquet("s3a://data-lake/events/") \
.filter(F.col("event_date") == "2024-06-15") \
.filter(F.col("country") == "RU") \
.groupBy("user_id") \
.agg(F.sum("amount").alias("total")) \
.show()
# Spark UI показывает:
# Scanned: 2 TB
# Rows output: 15,230
# Scan time: 42 minutes
# S3 GET requests: ~10,000 (footer) + ~200,000 (data)
Шаг 1: Column Pruning - убедиться что он работает¶
# Проверяем Physical Plan:
spark.read.parquet("s3a://data-lake/events/") \
.filter(F.col("event_date") == "2024-06-15") \
.filter(F.col("country") == "RU") \
.groupBy("user_id") \
.agg(F.sum("amount").alias("total")) \
.explain()
# ReadSchema: struct<user_id:bigint, amount:double, event_date:date, country:string>
# ✅ Column Pruning работает: 4 из 200 колонок
# Эффект: 2 TB → ~40 GB (50× меньше данных из S3 по объёму колонок)
# Но Row Groups пока не пропускаются (данные не отсортированы)
# Время: 45 min → ~12 min
Шаг 2: Партиционирование таблицы по event_date¶
# Переписать данные с партиционированием:
spark.read.parquet("s3a://data-lake/events/") \
.write \
.partitionBy("event_date") \ # 365 партиций за год
.sortWithinPartitions("country", "user_id") \
.option("parquet.block.size", 268435456) \ # 256 MB Row Groups
.mode("overwrite") \
.parquet("s3a://data-lake/events-v2/")
# Теперь запрос с event_date фильтром:
spark.read.parquet("s3a://data-lake/events-v2/") \
.filter(F.col("event_date") == "2024-06-15") \
.filter(F.col("country") == "RU") \
.groupBy("user_id") \
.agg(F.sum("amount").alias("total")) \
.show()
# Partition pruning: читаем только event_date=2024-06-15 (1/365 данных)
# Column Pruning: 4 из 200 колонок
# Row Group Filtering по country (данные отсортированы по country внутри партиции)
# Эффект: ~40 GB → ~500 MB
# Время: 12 min → ~90 секунд
Шаг 3: Bloom Filter для точечного поиска по country¶
# Переписать с Bloom Filter на колонке country:
spark.read.parquet("s3a://data-lake/events-v2/") \
.write \
.partitionBy("event_date") \
.sortWithinPartitions("country", "user_id") \
.option("parquet.block.size", 268435456) \
.option("parquet.bloom.filter.enabled", "true") \
.option("parquet.bloom.filter.column.names", "country") \
.option("parquet.bloom.filter.expected.ndv", "200") \ # ~200 стран
.mode("overwrite") \
.parquet("s3a://data-lake/events-v3/")
spark.conf.set("spark.sql.parquet.enableBloomFilterPushdown", "true")
# Запрос WHERE country = 'RU':
# → Bloom Filter мгновенно исключает Row Groups без записей RU
# → Ещё более точное пропускание по min/max (данные отсортированы)
# Время: 90 sec → ~54 секунды
Итоговое сравнение¶
def print_query_stats(df_plan_str: str, description: str):
"""Вывести ключевые метрики из Physical Plan."""
print(f"\n{'='*50}")
print(f"Версия: {description}")
print(f"{'='*50}")
print(df_plan_str)
# Метрики ДО → ПОСЛЕ оптимизации:
stats = {
"S3 трафик": "2 TB → ~500 MB (4000× меньше)",
"Scanned Row Groups": "16 000 → ~60 (267× меньше)",
"S3 GET requests": "210 000 → ~800 (262× меньше)",
"Scan time": "42 мин → ~54 сек (47× быстрее)",
"S3 стоимость ($/GB)": "$46.00 → $0.012 (3800× дешевле)",
}
for metric, comparison in stats.items():
print(f" {metric}: {comparison}")
Аудит через Spark UI: что искать¶
При анализе Parquet-запросов в Spark UI нужно смотреть на вкладку SQL → конкретный запрос → DAG или детали оператора:
FileScan parquet:
├── number of files read: 27 ← сколько файлов прочитано
├── scan time total: 00:01:23 ← время чтения из S3
├── metadata time: 00:00:02 ← время на Footer запросы
├── number of output rows: 15,230 ← строки после Row Group Filtering
├── bytes scanned: 512,345,678 ← байты прочитанные из S3
└── files pruned: 9,973 ← файлы пропущены partition pruning
Ключевые индикаторы проблем:
bytes scanned>> ожидаемого → нет Column Pruning или Row Group Filtering.number of output rows<<bytes scanned/ (avg_row_size) → данные не кластеризованы.scan time>>metadata time× 1000 → много данных читается, мало пропускается.files pruned == 0→ нет партиционирования или фильтр не совпадает с partition key.
Итого: чеклист оптимизации Parquet на S3¶
-
Никогда
SELECT *на таблицах с > 20 колонками. Column Pruning экономит пропорционально количеству ненужных колонок. -
Фильтры по оригинальным полям, не по вычисляемым.
WHERE year(event_date) = 2024→WHERE event_date >= '2024-01-01' AND event_date < '2025-01-01'. -
Партиционировать по низкокардинальным колонкам (дата, регион, категория). Максимум 10 000 партиций. Высокая кардинальность → Row Group Filtering внутри партиции.
-
Сортировать данные внутри партиций по часто используемым колонкам фильтрации.
sortWithinPartitions("user_id", "event_type")перед записью. -
Row Group Size 256–512 МБ для аналитических нагрузок на S3. Меньше файлов → меньше Footer GET запросов.
-
Bloom Filter для точечных запросов по высококардинальным колонкам (
user_id,order_id). -
Векторизованный читатель - убедитесь, что в Physical Plan
Batched: true. Избегайте Python UDF (используйте Pandas UDF или встроенные функции Spark). -
Iceberg/Delta Lake для дополнительного file-level pruning поверх Row Group Filtering Parquet.
-
Компакция мелких файлов - целевой размер файла 256–512 МБ. Тысячи мелких файлов убивают производительность на этапе планирования.
-
Мониторинг через
explain()(ReadSchema + PushedFilters) и Spark UI (bytes scanned, files pruned).