Форматы данных: Parquet, ORC, Avro - выбор, настройка, pushdown

Архитектура Parquet, ORC, Avro: row groups, predicate pushdown, column pruning, compression codecs, schema evolution и small files problem

core optimization

Почему формат хранения - фундамент производительности Spark

В реальном production pipeline 60–80% времени выполнения Spark Job тратится не на вычисления, а на чтение данных с диска или из облачного хранилища. Выбор формата хранения влияет на:

  • Объём IO: сколько байт читается с диска на один запрос
  • Стоимость сжатия: насколько меньше места занимают данные и насколько быстро они сжимаются/разжимаются
  • Возможности pushdown: может ли Spark пропускать целые файлы, блоки или строки без чтения
  • Скорость shuffle: меньше данных = меньше времени на перераспределение между executors
  • Стоимость хранения: на S3/GCS/ADLS вы платите и за хранение, и за операции чтения

Плохой выбор формата (например, CSV в production) может замедлить аналитический запрос в 10–100× по сравнению с правильно настроенным Parquet. Это не преувеличение - это повседневная реальность большинства дата-платформ.

Row-based vs Columnar: анатомия форматов

Всё многообразие форматов делится на два фундаментальных типа по способу физического расположения данных на диске.

Строчные форматы (Row-based)

В строчном формате данные записываются строка за строкой. Все поля одной записи лежат рядом:

[row1: id=1, name="Alice", age=30, amount=100.50]
[row2: id=2, name="Bob",   age=25, amount=75.00]
[row3: id=3, name="Carol", age=35, amount=200.25]

Преимущество строчного формата - атомарность строки. Чтобы получить всю строку (например, одну запись пользователя), нужно одно последовательное чтение. Именно поэтому строчные форматы идеальны для OLTP (транзакций): вставка одной записи, выборка одной записи по ключу.

Недостаток: если вам нужна только колонка amount для агрегации по 100 миллионам строк - вы читаете весь файл, включая name, age и все остальные поля, которые вам не нужны.

Примеры: CSV, JSON, Avro.

Колончатые форматы (Columnar)

В колончатом формате данные одной колонки хранятся вместе:

[ids:     1, 2, 3, ...]
[names:   "Alice", "Bob", "Carol", ...]
[ages:    30, 25, 35, ...]
[amounts: 100.50, 75.00, 200.25, ...]

Для аналитического запроса SELECT SUM(amount) WHERE age > 28 нужно прочитать только колонки amount и age - две из N. Остальные N-2 колонок не читаются вообще. Это column pruning в действии.

Дополнительное преимущество: значения одной колонки имеют одинаковый тип и часто схожие значения. Это делает их крайне эффективно сжимаемыми - алгоритмы сжатия видят паттерны (RLE, Dictionary encoding) намного лучше, чем в смешанных строчных данных.

Примеры: Parquet, ORC.

Сравнение форматов: когда что выбрать

Критерий CSV/JSON Avro Parquet ORC
Тип Row Row Columnar Columnar
Чтение всей строки Хорошо Хорошо Медленно Медленно
Аналитические запросы Плохо Плохо Отлично Отлично
Column pruning Нет Нет Да Да
Predicate pushdown Нет Нет Да Да
Сжатие Плохо Хорошо Отлично Отлично
Schema evolution Плохо Отлично Хорошо Хорошо
Streaming / CDC Хорошо Отлично Плохо Плохо
Spark экосистема Хорошо Хорошо Отлично Хорошо
Hive/Trino оптимизация Нет Нет Хорошо Отлично

Внутри Parquet файла: deep dive

Parquet - бинарный колончатый формат, разработанный Twitter и Cloudera (2013), сегодня de facto стандарт для аналитических workloads в Hadoop/Spark экосистеме. Понимание внутренней структуры файла напрямую объясняет, почему Parquet такой быстрый.

Иерархия структур Parquet

Row Groups

Row Group - горизонтальный раздел данных, охватывающий все строки определённого диапазона. Каждый Row Group содержит данные нескольких тысяч-миллионов строк. По умолчанию в Spark размер Row Group - 128 MB (настраивается через parquet.block.size).

Почему важен размер Row Group? Потому что это единица пропуска при predicate pushdown. Если statistics (min/max) Row Group не удовлетворяют фильтру - весь Row Group пропускается без чтения. Чем больше Row Group, тем лучше компрессия и тем меньше метаданных. Но чем меньше Row Group - тем точнее прунинг (меньше ложных чтений).

Оптимальный размер: 128 MB–256 MB. Меньше - слишком много метаданных, меньше компрессия. Больше - прунинг становится грубым.

Column Chunks и Pages

Внутри каждого Row Group данные каждой колонки хранятся в отдельном Column Chunk - непрерывном блоке данных для одной колонки всего Row Group.

Column Chunk состоит из Pages - единиц декомпрессии и декодирования:

  • Dictionary Page: словарь уникальных значений (для Dictionary Encoding)
  • Data Page: закодированные данные строк
  • Data Page v2: расширенный формат с встроенной статистикой на уровне page

Размер page по умолчанию - 1 MB. Это минимальный объём данных, который читается за раз. Меньший page = более точный pushdown, но больше overhead на метаданные.

Parquet footer - это блок метаданных в конце файла. Он содержит:

  • Полную схему файла (все колонки, типы, nested структуры)
  • Для каждого Row Group: byte offset (где он начинается на диске), размер
  • Для каждого Column Chunk в каждом Row Group: min/max значения, количество null, количество строк

При открытии Parquet файла Spark сначала читает footer (Seek to end of file → read footer size → seek back → read footer bytes). Это обязательная операция, даже если вы читаете только 10 строк. Именно поэтому spark.table("...").limit(10) на Parquet может занять несколько секунд - нужно прочитать footers всех файлов.

На HDFS это быстро (NameNode знает расположение блоков). На S3 каждый HEAD+GET range - отдельный HTTP запрос. Тысячи файлов = тысячи HTTP запросов только для footers. Это small files problem в контексте облачных хранилищ.

Encoding: как Parquet сжимает данные без compression codec

Прежде чем применить compression codec (Snappy, ZSTD), Parquet применяет encoding - логическое преобразование данных, которое убирает избыточность.

Dictionary Encoding - для колонок с небольшим количеством уникальных значений (страны, категории, статусы). Вместо хранения строки "COMPLETED" 10 миллионов раз хранится словарь {0: "COMPLETED", 1: "PENDING", 2: "FAILED"} и массив ID: [0, 0, 1, 0, 2, ...]. Вместо строк - компактные числа.

# Без dictionary encoding: 10M × "COMPLETED" (9 bytes) = 90 MB
# С dictionary encoding: 3-элементный словарь + 10M × 1 byte = ~10 MB + overhead
# Экономия: ~9×

Run Length Encoding (RLE) - для колонок с повторяющимися последовательными значениями. Вместо [1, 1, 1, 1, 1, 2, 2, 2] хранится [(1, 5), (2, 3)] - «значение 1 встречается 5 раз подряд, затем 2 встречается 3 раза». Особенно эффективно для boolean колонок и null bitmasок.

Delta Encoding - для монотонно растущих колонок (timestamp, sequence ID). Вместо абсолютных значений хранятся разницы: [1000, 1001, 1003, 1007][1000, 1, 2, 4]. Маленькие дельты хорошо кодируются небольшим количеством бит.

Bit Packing - если значения попадают в диапазон, скажем, 0–15 (4 бита), то вместо хранения 8-байтных int64 можно упаковать их в 4 бита на значение, экономя 16×.

Encoding применяется до compression codec. Это значит, что Snappy работает уже с данными, из которых убрана большая часть избыточности - и сжимает их ещё дополнительно.

ORC: Hive-native columnar format

ORC (Optimized Row Columnar) разработан Hortonworks для Hive (2013). Структурно похож на Parquet, но с рядом специфических особенностей.

Stripe в ORC - аналог Row Group в Parquet, но с другим дефолтным размером (256 MB). Каждый Stripe содержит:

  • Index data: min/max для каждых 10 000 строк (Row Index Entry) - более гранулярная статистика
  • Row data: закодированные колончатые данные
  • Stripe footer: метаданные и кодировки

ACID транзакции - ORC нативно поддерживает ACID в Hive через delta files и merge компакцию. Для Spark это менее актуально, так как Delta Lake реализует ACID поверх Parquet.

Bloom Filters в ORC - встроены в формат на уровне спецификации. В Parquet bloom filters появились позже (в спецификации Parquet 2.0, реализация в Spark с 3.x).

Векторизованный ридер в Hive был реализован для ORC раньше, чем для Parquet. Сегодня в Spark оба формата поддерживают векторизованное чтение. Но в экосистемах Hive/Presto/Trino ORC традиционно показывает чуть лучшие результаты для complex queries.

Когда выбирать ORC вместо Parquet:

  • Production Hive workloads с часто обновляемыми таблицами (ACID)
  • Запросы на Presto/Trino, где традиционно сильнее оптимизация под ORC
  • Если существующая инфраструктура уже основана на ORC

Для новых Spark-ориентированных платформ (Delta Lake, Iceberg) - Parquet предпочтительнее.

Avro: format для транспорта и CDC

Avro - строчный бинарный формат с JSON-схемой, разработанный Apache для Hadoop (2009). В отличие от Parquet и ORC, Avro оптимизирован не для аналитики, а для сериализации и передачи данных.

Почему Avro - стандарт для Kafka и CDC

Avro файл содержит схему (в JSON формате) прямо в заголовке - или схема хранится в Schema Registry, а в данных только schema ID (2-байтовый magic byte + 4-байтовый ID). Это позволяет:

  • Десериализовать одно сообщение независимо от контекста (нет нужды знать схему заранее)
  • Менять схему эволюционно, сохраняя обратную и прямую совместимость
  • Интегрировать с Confluent Schema Registry для централизованного управления схемами

Schema evolution в Avro: backward и forward compatibility

Avro строго разграничивает типы изменений схемы:

Backward compatible (новая схема может читать старые данные):

  • Добавление поля с "default" - старые записи без этого поля будут использовать default значение
  • Удаление поля без "default" - новая схема просто игнорирует поле при чтении старых данных

Forward compatible (старая схема может читать новые данные):

  • Добавление поля без "default" - старая схема проигнорирует новое поле

Полная совместимость (Full): изменение допустимо в обе стороны - только добавление полей с "default".

Именно эта строгая система совместимости делает Avro идеальным для CDC (Change Data Capture) пайплайнов: базы данных меняют схему, Debezium генерирует новую версию Avro-схемы, Schema Registry проверяет совместимость, потребители (Spark) могут читать и старые и новые записи.

Compression codecs: баланс CPU и storage

Compression работает поверх encoding - после того, как данные уже логически сжаты через Dictionary/RLE.

Snappy - исторически дефолт в Spark. Быстрая компрессия и декомпрессия, умеренное сжатие. Хороший выбор для промежуточных результатов и данных с высокой частотой чтения.

GZIP - высокое сжатие, но медленная декомпрессия. Стоит использовать для архивных данных, которые читаются редко, или для cold-tier storage, где экономия места важнее скорости.

ZSTD (Zstandard) - современный стандарт от Facebook (2015). Лучший баланс: сжатие сравнимо с GZIP, скорость декомпрессии сравнима со Snappy. Поддерживает настройку уровня компрессии (1–22, дефолт 3). Начиная со Spark 3.0 рекомендуется использовать ZSTD как дефолт.

LZ4 - самая быстрая декомпрессия. Используется в Spark для сжатия shuffle данных в памяти (spark.io.compression.codec = lz4). Для хранения на диске обычно предпочтительнее ZSTD или Snappy.

Настройка codec в Spark:

# Для Parquet
spark.conf.set("spark.sql.parquet.compression.codec", "zstd")

# Или при записи конкретного DataFrame
df.write.option("compression", "zstd") \
        .parquet("/path/to/output")

# ZSTD с настройкой уровня (требует Spark 3.2+)
df.write.option("compression", "zstd") \
        .option("parquet.compression.codec.zstd.level", "6") \
        .parquet("/path/to/output")

# Для ORC
df.write.option("compression", "zlib").orc("/path/to/output")

# Для Avro
df.write.option("compression", "snappy").format("avro").save("/path/to/output")

Рекомендация по умолчанию: ZSTD с уровнем 3 для Silver/Gold слоёв (баланс), Snappy для временных файлов и shuffle, GZIP/ZSTD-6 для архивных данных.

Pushdown оптимизации: Spark читает меньше

Pushdown - это ключевой механизм, который делает Parquet/ORC принципиально быстрее CSV. Это семейство оптимизаций, при которых Spark передаёт ограничения (фильтры, проекции) на уровень файлового ридера, чтобы он сам пропускал ненужные данные - не читая их вообще.

Column Pruning (Column Pushdown)

Самая простая и самая мощная оптимизация. Если ваш запрос использует только 3 колонки из 50 - Parquet читает только эти 3 Column Chunks. Остальные 47 не читаются вообще.

# DataFrame с 50 колонками
df = spark.read.parquet("/data/wide_table/")

# Без Column Pruning (CSV): читается весь файл
# С Column Pruning (Parquet): читается только amount + status

df.select("amount", "status").filter(col("status") == "COMPLETED").show()

Column Pruning работает автоматически - Catalyst анализирует запрос, определяет нужные колонки и передаёт их список в Parquet ридер. Никакой дополнительной конфигурации не требуется.

Эффект: для таблицы с 50 колонками и запросом на 3 - экономия ~94% IO.

Predicate Pushdown (Row Group Skipping)

Если в запросе есть фильтр (WHERE, .filter()), Spark передаёт его в Parquet ридер. Ридер проверяет min/max статистику каждого Row Group из footer. Если max(amount) в Row Group = 500, а фильтр требует amount > 1000 - весь Row Group пропускается без чтения.

# Spark генерирует predicate и передаёт в Parquet reader
df.filter(col("event_date") >= "2024-01-01") \
  .filter(col("amount") > 10000) \
  .select("user_id", "amount") \
  .show()

# Parquet reader проверяет для каждого Row Group:
# event_date: max < "2024-01-01"? → Skip
# amount: max < 10000? → Skip
# Только Row Groups с подходящими min/max читаются

Эффективность predicate pushdown зависит от корреляции между данными и физическим расположением. Если данные случайным образом перемешаны - все Row Groups будут содержать значения из всего диапазона, pushdown не поможет. Если данные отсортированы по фильтруемой колонке - большинство Row Groups попадут в один диапазон, и pushdown уберёт 99% IO.

# Для максимального эффекта от predicate pushdown - сортировка перед записью
df.orderBy("event_date", "user_id") \
  .write.parquet("/data/events_sorted/")
# Теперь фильтры по event_date пропускают почти все Row Groups

Partition Pruning: file-level skipping

Partition Pruning - это оптимизация на уровне файловой системы, не Row Groups. Если данные партиционированы (физически разложены по папкам по значению колонки), Spark вообще не обращается к папкам, не удовлетворяющим фильтру.

# Запись с партиционированием
df.write.partitionBy("year", "month") \
        .parquet("/data/events/")

# Структура папок:
# /data/events/year=2024/month=01/part-00000.parquet
# /data/events/year=2024/month=02/part-00000.parquet
# /data/events/year=2023/month=12/part-00000.parquet

# Чтение - Spark видит фильтр year=2024 AND month=01
# и читает ТОЛЬКО /data/events/year=2024/month=01/
df.read.parquet("/data/events/") \
   .filter((col("year") == 2024) & (col("month") == 1)) \
   .count()

Partition Pruning работает только для колонок, по которым написан partitionBy. Это более грубое, но и более мощное инструмент: пропускаются целые файлы (и даже папки), а не только Row Groups.

Как Catalyst передаёт predicates в storage layer

Проверить, какие фильтры Spark смог протолкнуть в Parquet ридер:

df.filter(col("age") > 30).select("name", "amount").explain()

# Physical Plan:
# *(1) Project [name#5, amount#6]
# +- *(1) Filter (isnotnull(age#7) AND (age#7 > 30))
#    +- *(1) FileScan parquet [name#5, amount#6, age#7]
#            Batched: true,
#            DataFilters: [isnotnull(age#7), (age#7 > 30)],
#            Format: Parquet,
#            PushedFilters: [IsNotNull(age), GreaterThan(age,30)],
#            ReadSchema: struct<name:string, amount:double>

PushedFilters - это фильтры, которые переданы в Parquet ридер для Row Group прунинга.
ReadSchema - это projected схема (только нужные колонки).
Обратите внимание: age тоже входит в ReadSchema (нужна для проверки фильтра), хотя не входит в select.

Ограничения predicate pushdown

Не все фильтры можно протолкнуть. Catalyst отказывается от pushdown если:

  • Фильтр применяется к результату UDF - оптимизатор не знает семантику UDF
  • Фильтр использует OR с разными колонками - сложно оценить по min/max
  • Фильтр на вычисляемой колонке: .filter(col("a") + col("b") > 100) - нет статистики для суммы
  • Фильтр по вложенному полю struct: частично поддерживается в Spark 3.x
# ХОРОШО: фильтр по обычной колонке → pushdown работает
df.filter(col("amount") > 1000)

# ПЛОХО: UDF в фильтре → pushdown не работает
df.withColumn("score", my_udf(col("amount"))).filter(col("score") > 0.5)

# ХОРОШО: сначала filter, потом UDF
df.filter(col("amount") > 1000).withColumn("score", my_udf(col("amount")))

Vectorized Parquet Reader

Начиная со Spark 2.0, чтение Parquet (и ORC) выполняется через векторизованный ридер - он читает данные колончатыми батчами (columnar batch), а не строками. Это позволяет использовать CPU SIMD инструкции при декодировании и избежать создания тысяч Java объектов.

# Включён по умолчанию для Parquet
spark.conf.get("spark.sql.parquet.enableVectorizedReader")  # "true"

# Для ORC
spark.conf.get("spark.sql.orc.enableVectorizedReader")  # "true"

Vectorized reader работает совместно с WholeStageCodeGen: данные читаются колончатыми векторами и сразу обрабатываются сгенерированным JVM кодом без промежуточных Row объектов. Для CSV и JSON такой оптимизации нет - там каждая строка парсится в Java String, затем конвертируется в нужный тип.

Именно это объясняет разницу в 10–50× между Parquet и CSV на аналитических запросах: не только меньший объём IO, но и принципиально эффективнее обработка после чтения.

Schema evolution в Parquet

Parquet поддерживает additive schema evolution - безопасное добавление новых колонок. При этом:

  • Старые файлы без новой колонки → Spark читает их как null для отсутствующей колонки
  • Новые файлы с новой колонкой → читаются полностью
# Исходная схема: id, name, amount
df_v1 = spark.read.parquet("/data/events/2024-01/")

# Новая схема: id, name, amount, category (новая колонка)
df_v2 = spark.read.parquet("/data/events/2024-02/")

# Автоматическое слияние схем - ДОРОГО (читает все footers)
df_all = spark.read.option("mergeSchema", "true").parquet("/data/events/")

# Лучше: явно указать схему через schema= параметр
from pyspark.sql.types import StructType, StructField, LongType, StringType, DoubleType

target_schema = StructType([
    StructField("id", LongType()),
    StructField("name", StringType()),
    StructField("amount", DoubleType()),
    StructField("category", StringType()),  # nullable=True по умолчанию
])

df_all = spark.read.schema(target_schema).parquet("/data/events/")

Schema merge overhead

mergeSchema=true читает footer каждого файла чтобы составить объединённую схему. На таблице из 10 000 файлов - это 10 000 footer чтений при каждом открытии таблицы. На S3 это тысячи HTTP запросов перед тем, как прочитать хоть байт данных.

Решение: используйте mergeSchema=true только один раз для построения финальной схемы, затем сохраните её и передавайте явно через schema=.

Breaking changes (несовместимые изменения) в Parquet:

  • Изменение типа колонки (int → string) - невозможно без полного rewrite
  • Переименование колонки - обрабатывается только через metadata alias
  • Удаление колонки - данные остаются в файлах, но при чтении с явной схемой просто игнорируются

Small Files Problem: главная болезнь production Spark

Тысячи маленьких файлов - одна из самых частых проблем в production data pipelines. Каждый executor при чтении Parquet:

  1. Открывает файл (HTTP HEAD на S3)
  2. Читает footer (HTTP GET range)
  3. Принимает решение, читать ли данные

На 10 000 файлов по 1 MB: 20 000 HTTP запросов только для метаданных, плюс перегруженный NameNode (HDFS) или S3 Rate Limiting (429 ошибки). Реальный scan на S3 при small files может занимать 80% времени на метаданные.

# Диагностика: смотрим количество файлов
import subprocess
result = subprocess.run(
    ["hdfs", "dfs", "-count", "/data/events/"],
    capture_output=True, text=True
)
# Output: DIR_COUNT FILE_COUNT CONTENT_SIZE PATHNAME
# 1234 987654 50000000000 /data/events/

# Много маленьких файлов? Нужен compaction.

Оптимальный размер файлов

Цель: файлы от 128 MB до 1 GB. Меньше - слишком много метаданных. Больше - плохо параллелизуется (один task на один файл, часть executors будет простаивать).

# ПЕРЕД записью: контролируем количество файлов

# Вариант 1: coalesce - уменьшаем без shuffle (только объединение)
df.coalesce(100).write.parquet("/data/output/")
# Минус: может создать неравномерные файлы если данные неравномерны

# Вариант 2: repartition - с shuffle, равномерное распределение
df.repartition(100).write.parquet("/data/output/")
# Минус: дополнительный shuffle

# Вариант 3: через spark.sql.files.maxRecordsPerFile
spark.conf.set("spark.sql.files.maxRecordsPerFile", "5000000")  # 5M записей на файл
df.write.parquet("/data/output/")
# Spark сам разобьёт партиции если в одном файле слишком много записей

Для расчёта оптимального количества партиций:

# Оцениваем размер DataFrame
total_size_bytes = df.rdd.map(lambda row: len(str(row))).sum()  # грубо

# Цель: 256 MB на файл
target_file_size = 256 * 1024 * 1024
optimal_partitions = max(1, int(total_size_bytes / target_file_size))

df.repartition(optimal_partitions).write.parquet("/data/output/")

Compaction: исправляем small files

Compaction - процесс объединения мелких файлов в крупные. В Delta Lake это встроено (OPTIMIZE). Для raw Parquet нужно делать вручную:

def compact_parquet_table(
    spark,
    path: str,
    target_partitions: int = 100,
    partition_cols: list = None,
):
    df = spark.read.parquet(path)

    if partition_cols:
        df.repartition(target_partitions, *partition_cols) \
          .write.mode("overwrite") \
          .partitionBy(*partition_cols) \
          .parquet(path + "_compacted")
    else:
        df.repartition(target_partitions) \
          .write.mode("overwrite") \
          .parquet(path + "_compacted")

compact_parquet_table(spark, "/data/events/", target_partitions=50)

Parquet tuning: практические настройки

spark = SparkSession.builder \
    .config("spark.sql.parquet.compression.codec", "zstd") \
    .config("spark.hadoop.parquet.block.size", str(256 * 1024 * 1024)) \
    .config("spark.hadoop.parquet.page.size", str(1 * 1024 * 1024)) \
    .config("spark.sql.parquet.enableVectorizedReader", "true") \
    .config("spark.sql.files.maxPartitionBytes", str(256 * 1024 * 1024)) \
    .config("spark.sql.files.openCostInBytes", str(4 * 1024 * 1024)) \
    .getOrCreate()
  • parquet.block.size - размер Row Group. По умолчанию 128 MB. Увеличьте до 256 MB для больших таблиц (лучше сжатие, меньше метаданных). Уменьшите до 64 MB если нужна точность predicate pushdown на маленьких диапазонах.

  • parquet.page.size - размер Page. По умолчанию 1 MB. Уменьшение улучшает точность dictionary encoding для маленьких групп значений, но увеличивает overhead.

  • spark.sql.files.maxPartitionBytes - Spark разбивает большие Parquet файлы на партиции при чтении. По умолчанию 128 MB. Увеличьте до 256 MB для аналитических запросов (меньше tasks, больше данных на task).

  • spark.sql.files.openCostInBytes - оценка стоимости открытия одного файла. Используется при упаковке нескольких маленьких файлов в одну партицию. По умолчанию 4 MB. Если файлы на S3 и открытие дорогое - увеличьте.

Форматы в медальонной архитектуре

Bronze: сырые данные в максимальном сохранении исходной структуры.

  • Kafka → Avro (с Schema Registry): схема зафиксирована, schema evolution из коробки
  • REST API → JSON: точная копия source без потерь
  • CDC → Avro: возможность replay при смене downstream схемы

Silver: очищенные, нормализованные данные для аналитики.

  • Parquet + ZSTD: максимальная скорость запросов
  • partitionBy("year", "month", "day"): partition pruning по дате
  • orderBy("user_id") перед записью: predicate pushdown по user_id
  • Явная схема без mergeSchema

Gold: агрегированные витрины.

  • Parquet или Delta Lake: если нужны ACID обновления (SCD, corrections) - Delta
  • Маленькие таблицы: можно без партиционирования (broadcast join с dim)

Практика: benchmark CSV vs Parquet vs ORC

from pyspark.sql import SparkSession
from pyspark.sql import functions as F
import time

spark = SparkSession.builder \
    .config("spark.sql.parquet.enableVectorizedReader", "true") \
    .getOrCreate()

# Генерируем синтетический датасет: 50M строк
n = 50_000_000
df = spark.range(n).select(
    F.col("id"),
    (F.rand() * 10000).alias("amount"),
    F.date_add(F.lit("2020-01-01"), (F.rand() * 1460).cast("int")).alias("event_date"),
    F.element_at(F.array(F.lit("A"), F.lit("B"), F.lit("C"), F.lit("D")),
                 (F.rand() * 4 + 1).cast("int")).alias("category"),
)

# Записываем в разных форматах
df.write.mode("overwrite").csv("/tmp/bench/csv/")
df.write.mode("overwrite").parquet("/tmp/bench/parquet/")
df.write.mode("overwrite").option("compression", "zstd").parquet("/tmp/bench/parquet_zstd/")
df.write.mode("overwrite").orc("/tmp/bench/orc/")

# Тест 1: полный scan + агрегация
def bench(name, reader):
    start = time.time()
    reader.agg(F.sum("amount"), F.count("id")).collect()
    print(f"{name}: {time.time() - start:.1f}s")

bench("CSV", spark.read.csv("/tmp/bench/csv/", header=True, inferSchema=True))
bench("Parquet (Snappy)", spark.read.parquet("/tmp/bench/parquet/"))
bench("Parquet (ZSTD)", spark.read.parquet("/tmp/bench/parquet_zstd/"))
bench("ORC", spark.read.orc("/tmp/bench/orc/"))

# Ожидаемые результаты (YARN кластер, 4 executors, 8 cores each):
# CSV:               185s  (нет pushdown, нет columnar, медленный парсинг строк)
# Parquet (Snappy):   12s  (columnar, vectorized reader)
# Parquet (ZSTD):     10s  (лучше сжатие = меньше IO)
# ORC:                11s  (сравнимо с Parquet)

Практика: predicate pushdown analysis

# Записываем с сортировкой для эффективного predicate pushdown
df.orderBy("event_date") \
  .write.mode("overwrite") \
  .parquet("/tmp/bench/parquet_sorted/")

# Несортированные данные
df.write.mode("overwrite").parquet("/tmp/bench/parquet_unsorted/")

# Тест: фильтр по event_date (только последние 30 дней)
filter_expr = F.col("event_date") >= "2023-12-01"

def bench_with_filter(name, path):
    start = time.time()
    count = spark.read.parquet(path).filter(filter_expr).count()
    elapsed = time.time() - start
    print(f"{name}: {elapsed:.1f}s, count={count}")

bench_with_filter("Unsorted", "/tmp/bench/parquet_unsorted/")
bench_with_filter("Sorted", "/tmp/bench/parquet_sorted/")

# Что смотреть в Spark UI → SQL tab → FileScan parquet:
# "number of files pruned" - сколько Row Groups пропущено
# Для sorted: большинство Row Groups будут вне диапазона → пропущены

# Проверяем через explain:
spark.read.parquet("/tmp/bench/parquet_sorted/") \
     .filter(filter_expr) \
     .select("id", "amount") \
     .explain()
# PushedFilters: [GreaterThanOrEqual(event_date,2023-12-01)]
# ReadSchema: struct<id:bigint,amount:double>
# event_date НЕ в ReadSchema - читается только для фильтра на Row Group level,
# не загружается в executor memory

Практика: schema evolution pipeline

# Эмулируем добавление новой колонки в схему

# Шаг 1: данные v1 (без поля "region")
df_v1 = spark.createDataFrame([
    (1, "Alice", 100.0),
    (2, "Bob", 200.0),
], ["id", "name", "amount"])

df_v1.write.mode("overwrite").parquet("/tmp/evolution/v1/")

# Шаг 2: данные v2 (с новым полем "region")
df_v2 = spark.createDataFrame([
    (3, "Carol", 300.0, "North"),
    (4, "Dave", 400.0, "South"),
], ["id", "name", "amount", "region"])

df_v2.write.mode("overwrite").parquet("/tmp/evolution/v2/")

# Шаг 3: чтение обеих версий
# Вариант A: mergeSchema=true (дорого, но автоматично)
df_merged = spark.read.option("mergeSchema", "true") \
                 .parquet("/tmp/evolution/v1/", "/tmp/evolution/v2/")
df_merged.show()
# v1 строки: region = null
# v2 строки: region = "North"/"South"

# Вариант B: явная схема (предпочтительно в production)
from pyspark.sql.types import StructType, StructField, LongType, StringType, DoubleType

target_schema = StructType([
    StructField("id", LongType()),
    StructField("name", StringType()),
    StructField("amount", DoubleType()),
    StructField("region", StringType()),  # nullable=True → null для v1 данных
])

df_explicit = spark.read.schema(target_schema) \
                   .parquet("/tmp/evolution/v1/", "/tmp/evolution/v2/")
df_explicit.show()

Практика: оптимизация small files

from pyspark.sql import functions as F

# Эмулируем проблему: streaming pipeline создал 10 000 маленьких файлов
# (реальный кейс: micro-batch Spark Streaming, checkpoint=1 минута)
path = "/tmp/small_files/"

# Диагностика размера файлов
import subprocess
result = subprocess.run(["hdfs", "dfs", "-ls", "-R", path],
                        capture_output=True, text=True)
file_count = result.stdout.count(".parquet")
print(f"Файлов: {file_count}")

# Compaction: читаем все, объединяем, перезаписываем
df_messy = spark.read.parquet(path)
total_rows = df_messy.count()
print(f"Строк: {total_rows}")

# Рассчитываем оптимальное количество файлов
# Цель: 256 MB на файл
# Если средняя строка 1 KB, то в файле 256 * 1024 строк = 262 144
target_rows_per_file = 256 * 1024
optimal_files = max(1, total_rows // target_rows_per_file)
print(f"Целевое количество файлов: {optimal_files}")

# Compaction
df_messy.repartition(optimal_files) \
        .write.mode("overwrite") \
        .option("compression", "zstd") \
        .parquet(path + "_compacted/")

print("Compaction завершён")

Anti-patterns: что не делать с форматами

1. CSV в production как основной storage format

# ПЛОХО: CSV в production - нет pushdown, нет columnar, медленный парсинг
spark.read.csv("/data/events/*.csv", header=True, inferSchema=True)

# ХОРОШО: конвертируем CSV в Parquet при первом ingestion
csv_df = spark.read.csv("/landing/events.csv", header=True, schema=event_schema)
csv_df.write.partitionBy("year", "month").parquet("/bronze/events/")

2. inferSchema=True без кэширования результата

# ПЛОХО: inferSchema читает файл дважды - первый раз для определения схемы
spark.read.csv("/data/huge_file.csv", inferSchema=True)

# ХОРОШО: определяем схему один раз, передаём явно
schema = StructType([...])
spark.read.csv("/data/huge_file.csv", schema=schema)

3. mergeSchema=True на каждый запрос

# ПЛОХО: 10 000 footer чтений при каждом открытии таблицы
df = spark.read.option("mergeSchema", "true").parquet("/data/huge_table/")

# ХОРОШО: mergeSchema один раз → сохранить схему → использовать явно
merged_schema = spark.read.option("mergeSchema", "true").parquet("/data/huge_table/").schema
import json
with open("schema.json", "w") as f:
    f.write(merged_schema.json())

# В будущем:
from pyspark.sql.types import StructType
with open("schema.json") as f:
    schema = StructType.fromJson(json.load(f))
df = spark.read.schema(schema).parquet("/data/huge_table/")

4. Миллионы мелких файлов без compaction

# ПЛОХО: streaming без coalesce → тысячи файлов
stream_df.writeStream.parquet("/data/output/")  # 1 файл на micro-batch

# ХОРОШО: периодический compaction или настройка trigger
# В Spark Structured Streaming: используйте trigger(processingTime='5 minutes')
# и coalesce внутри query для укрупнения файлов

5. Parquet для Kafka CDC данных (вместо Avro)

# ПЛОХО: Parquet плохо подходит для частых мелких write операций
kafka_stream.writeStream.format("parquet").start()  # создаёт мини-файлы

# ХОРОШО: Avro для streaming, Parquet для batch storage
kafka_stream.writeStream.format("avro").start("/bronze/avro/")
# Затем batch job конвертирует в Parquet
spark.read.format("avro").load("/bronze/avro/") \
     .write.parquet("/silver/parquet/")

6. Partitioning по high-cardinality колонке

# ПЛОХО: partitionBy(user_id) при 10M пользователях → 10M папок
df.write.partitionBy("user_id").parquet("/data/events/")

# ХОРОШО: partitionBy по дате или категории (низкая кардинальность)
df.write.partitionBy("year", "month").parquet("/data/events/")
# Или: bucket by user_id (не папки, а hash-based)
df.write.bucketBy(1024, "user_id").sortBy("user_id").saveAsTable("events_bucketed")

Checklist: выбор формата и настройка

Задача Рекомендация
Production analytics storage Parquet + ZSTD
Kafka / CDC ingestion Avro + Confluent Schema Registry
Hive / Presto workloads ORC + ZLIB
Временные файлы, shuffle Snappy или LZ4
Архивные данные (редкое чтение) GZIP или ZSTD level 9
Compression дефолт для Spark 3.x ZSTD (лучший баланс)
Фильтры по дате → predicate pushdown orderBy("date") перед записью
Много уникальных значений в фильтруемой колонке Bloom filters в ORC
Новая колонка без полного rewrite Явная nullable схема при чтении
Много маленьких файлов Compaction: repartition(N).write.parquet()
Цель по размеру файла 128 MB - 1 GB на файл
Row Group size 128 MB (дефолт) или 256 MB для крупных таблиц
Партиционирование По дате (year/month) - низкая кардинальность
Партиционирование по высококардинальным ключам bucketBy() вместо partitionBy()
mergeSchema=true в production Только один раз, затем сохранять схему явно
Читать CSV на prod Немедленно конвертировать в Parquet