Parquet на S3: vectorized reader, column pruning и row group filtering

Физическая анатомия Parquet: как column pruning, predicate pushdown и vectorized reader позволяют Spark читать только нужные байты из S3 и экономить сетевой трафик в 10–100 раз.

storage

Формат хранения - это половина производительности

В традиционном 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 не может знать, как читать файл.

Когда 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: пропускаем целые блоки данных

Каждый 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. Для каждой строки:

  1. Выделить память под объект Row.
  2. Скопировать значения из декодированных буферов в поля объекта.
  3. Передать объект в следующий оператор.
  4. 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

  1. Никогда SELECT * на таблицах с > 20 колонками. Column Pruning экономит пропорционально количеству ненужных колонок.

  2. Фильтры по оригинальным полям, не по вычисляемым. WHERE year(event_date) = 2024WHERE event_date >= '2024-01-01' AND event_date < '2025-01-01'.

  3. Партиционировать по низкокардинальным колонкам (дата, регион, категория). Максимум 10 000 партиций. Высокая кардинальность → Row Group Filtering внутри партиции.

  4. Сортировать данные внутри партиций по часто используемым колонкам фильтрации. sortWithinPartitions("user_id", "event_type") перед записью.

  5. Row Group Size 256–512 МБ для аналитических нагрузок на S3. Меньше файлов → меньше Footer GET запросов.

  6. Bloom Filter для точечных запросов по высококардинальным колонкам (user_id, order_id).

  7. Векторизованный читатель - убедитесь, что в Physical Plan Batched: true. Избегайте Python UDF (используйте Pandas UDF или встроенные функции Spark).

  8. Iceberg/Delta Lake для дополнительного file-level pruning поверх Row Group Filtering Parquet.

  9. Компакция мелких файлов - целевой размер файла 256–512 МБ. Тысячи мелких файлов убивают производительность на этапе планирования.

  10. Мониторинг через explain() (ReadSchema + PushedFilters) и Spark UI (bytes scanned, files pruned).