Форматы данных: Parquet, ORC, Avro - выбор, настройка, pushdown
Архитектура Parquet, ORC, Avro: row groups, predicate pushdown, column pruning, compression codecs, schema evolution и small files problem
Почему формат хранения - фундамент производительности 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 на метаданные.
Footer: самая важная операция при чтении¶
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:
- Открывает файл (HTTP HEAD на S3)
- Читает footer (HTTP GET range)
- Принимает решение, читать ли данные
На 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 |