Форматы файлов в HDFS: ORC vs Parquet vs Avro, Snappy vs ZSTD vs Gzip - когда что выбирать

Архитектурный разбор форматов хранения и кодеков сжатия для on-premise HDFS Data Lake: физика строчного и колоночного хранения, анатомия Parquet и ORC, Avro для шины данных, Snappy vs ZSTD vs Gzip и сплинтуемость, Vectorized Reader, Predicate Pushdown, матрица выбора для Medallion Architecture, полная реализация на PySpark.

storage

1. Архитектурная классификация: строчные против колоночных форматов

Выбор формата хранения данных - одно из фундаментальных архитектурных решений при построении Data Lake. Это не техническая деталь, которую можно отложить на потом: неправильный выбор на раннем этапе приводит к тому, что Spark-кластер тратит 80% времени либо на пустую распаковку гигабайт по сети (I/O Bound), либо упирается в CPU при банальном чтении (CPU Bound).

Почему классические форматы неприемлемы на терабайтных данных

CSV, JSON и текстовые логи были созданы для обмена данными между системами, а не для аналитики. Их главные недостатки при работе со Spark на HDFS:

Отсутствие встроенной схемы. CSV не хранит типы данных. Каждый раз при чтении Spark либо инферирует схему (медленно - второй проход по файлу), либо читает всё как StringType и вы ловите ошибки при агрегации. JSON хранит названия полей в каждой строке, что создаёт колоссальный overhead: для файла с 100 колонками каждый JSON-объект содержит 100 строковых ключей, которые Spark должен парсить при каждом чтении.

Полное сканирование всегда. Для выполнения SELECT amount FROM sales WHERE date = '2024-01-15' CSV-ридер прочитает 100% данных файла, включая все колонки кроме amount и date. Эти данные выброшены в /dev/null, но уже занимают пропускную способность сети и дисков.

Нет сжатия на уровне структуры. Числа в CSV хранятся как строки: 1000000 занимает 7 байт вместо 4 байт INT32. Строки без схемы не позволяют применить dictionary encoding.

Схема показывает фундаментальное различие: при строчном хранении Spark читает каждую строку целиком, затем применяет фильтр и отбрасывает ненужные колонки. При колоночном хранении Spark читает только нужные колонки и может отбросить целые блоки данных, не читая их вообще, на основе хранящейся в файле статистики (min/max значений).

Строчный подход: физика хранения

В строчном формате записи хранятся последовательно: все поля строки 1, затем все поля строки 2, затем строки 3... Это идеально для паттерна «прочитать/записать один объект целиком»:

[user_id=1, name="Alice", amount=500, date="2024-01-15", status="active"]
[user_id=2, name="Bob",   amount=200, date="2024-01-15", status="closed"]
[user_id=3, name="Carol", amount=800, date="2024-01-16", status="active"]

Строчные форматы отлично подходят для:

  • OLTP: транзакция всегда работает с одной записью целиком (INSERT / UPDATE / SELECT * WHERE id=X)
  • Event streaming: каждое событие - самодостаточный объект со всеми полями
  • Сериализация между сервисами: передача объектов по сети (Kafka, Debezium CDC)

Колоночный подход: физика хранения

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

[user_id]:  1,     2,     3,     4,     5
[name]:     Alice, Bob,   Carol, Dave,  Eve
[amount]:   500,   200,   800,   150,   300
[date]:     01-15, 01-15, 01-16, 01-15, 01-16
[status]:   active,closed,active,active,closed

Колоночные форматы отлично подходят для:

  • OLAP: аналитический запрос читает 2-3 колонки из 50
  • Агрегации: SUM(amount) читает только колонку amount, без name, date, status
  • Компрессия: значения одной колонки похожи друг на друга (числа, категории), сжимаются в разы лучше

2. Apache Avro: идеальный выбор для шины данных и слоя Bronze

Apache Avro - строчный бинарный формат с возможностью Schema Evolution. Создан в экосистеме Hadoop как альтернатива XML/JSON для межсистемного взаимодействия.

Анатомия Avro-файла

Ключевая особенность: схема данных хранится прямо в заголовке каждого Avro-файла в формате JSON. Это означает, что Avro-файл самодостаточен: любой читатель, даже не имея внешнего описания схемы, может понять структуру данных. Это критически важно для долгосрочного архивного хранения.

Sync Marker - 16 байт случайных данных, уникальных для каждого файла, которые разделяют блоки. Это делает Avro сплинтуемым: Spark может найти начало любого блока, отыскав Sync Marker в произвольном месте файла, и начать чтение с середины.

Schema Evolution: главное преимущество Avro

Schema Evolution - это возможность изменять схему данных со временем при сохранении совместимости со старыми данными. Именно это делает Avro стандартом для потоковой обработки и CDC.

# Исходная схема (версия 1.0)
schema_v1 = {
    "type": "record",
    "name": "UserEvent",
    "fields": [
        {"name": "user_id", "type": "long"},
        {"name": "event_type", "type": "string"},
        {"name": "timestamp", "type": "long"}
    ]
}

# Обновлённая схема (версия 2.0) - добавили новые поля
# Старые данные в Avro-файлах за прошлые месяцы остаются читаемыми!
schema_v2 = {
    "type": "record",
    "name": "UserEvent",
    "fields": [
        {"name": "user_id", "type": "long"},
        {"name": "event_type", "type": "string"},
        {"name": "timestamp", "type": "long"},
        # Новое поле: если его нет в старых данных - подставляется default
        {"name": "session_id", "type": ["null", "string"], "default": None},
        # Ещё одно новое поле
        {"name": "platform", "type": "string", "default": "web"}
    ]
}

# При чтении старых данных с новой схемой:
# - session_id будет null (default значение)
# - platform будет "web" (default значение)
# Никаких ошибок, никакой потери данных!

Три вида совместимости Avro:

  • Backward compatibility: новая схема может читать данные, написанные старой схемой (добавление новых полей с default - OK)
  • Forward compatibility: старая схема может читать данные, написанные новой схемой (удаление полей с default - OK)
  • Full compatibility: оба направления совместимы одновременно

Avro в экосистеме Kafka и Debezium

Avro - безальтернативный стандарт для Kafka в корпоративных Data Lake:

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

spark = SparkSession.builder \
    .appName("avro-bronze-ingestion") \
    .config("spark.jars.packages", "org.apache.spark:spark-avro_2.12:3.5.1") \
    .getOrCreate()

# Чтение потока событий из Kafka (Avro с Schema Registry)
kafka_df = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "kafka:9092") \
    .option("subscribe", "user_events") \
    .load()

# Десериализация Avro (байты → DataFrame)
# Schema Registry хранит схемы, Spark получает их по schema_id
from pyspark.sql.avro.functions import from_avro

avro_schema = """
{
  "type": "record",
  "name": "UserEvent",
  "fields": [
    {"name": "user_id", "type": "long"},
    {"name": "event_type", "type": "string"},
    {"name": "amount", "type": ["null", "double"], "default": null},
    {"name": "timestamp", "type": "long"}
  ]
}
"""

events_df = kafka_df.select(
    from_avro(F.col("value"), avro_schema).alias("event")
).select("event.*")

# Записываем в Bronze слой как Avro + Snappy
# Сохраняем ОРИГИНАЛЬНЫЙ формат - это важно для Bronze!
# Bronze должен хранить данные максимально близко к источнику
query = events_df.writeStream \
    .format("avro") \
    .option("path", "hdfs://cluster/data/bronze/user_events/") \
    .option("checkpointLocation", "/tmp/checkpoints/bronze/user_events") \
    .partitionBy("event_type") \
    .option("compression", "snappy") \
    .trigger(processingTime="30 seconds") \
    .start()

Когда использовать Avro:

  • Потоковая загрузка данных из Kafka (скорость записи критична)
  • CDC из RDBMS через Debezium (схема может меняться при ALTER TABLE)
  • Обмен данными между микросервисами
  • Bronze слой Data Lake (сохранение оригинальных событий)

Когда НЕ использовать Avro:

  • Аналитические запросы с GROUP BY, агрегациями (строчный формат - медленно)
  • BI-инструменты и Spark SQL на Silver/Gold слоях
  • Данные с множеством числовых колонок (колоночное сжатие лучше)

3. Apache Parquet: стандарт де-факто в экосистеме Spark

Apache Parquet - колоночный бинарный формат, созданный совместно Twitter и Cloudera в 2013 году. Сегодня это основной формат хранения данных в Spark, Delta Lake, Apache Iceberg и Apache Hudi.

Внутренняя физика: анатомия Parquet-файла

Понимание внутренней структуры Parquet критически важно для правильного тюнинга Spark. Каждый аспект структуры влияет на производительность конкретных видов запросов.

Row Group - горизонтальный срез данных. По умолчанию 128 MB, что намеренно выровнено с размером блока HDFS. Один Row Group = один блок HDFS = одна Spark-задача при чтении. Это ключевое: если Row Group = 128 MB и блок HDFS = 128 MB, то один Spark Task читает ровно один Row Group с одного DataNode.

Column Chunk - вертикальный срез Row Group для одной колонки. Физически все значения данной колонки в данном Row Group хранятся рядом. Это обеспечивает column pruning: при запросе SELECT amount Spark читает только Column Chunk для amount, пропуская остальные.

Page - минимальная единица чтения внутри Column Chunk. Обычно 8 KB – 1 MB. Page может быть Data Page (данные), Dictionary Page (словарь для dictionary encoding) или Index Page.

File Footer - метаданные всего файла. Spark читает Footer в первую очередь при открытии файла. Footer содержит статистику по каждому Column Chunk: min_value, max_value, null_count. Именно эта статистика позволяет Spark пропускать Row Groups при predicate pushdown.

Важно: Footer читается последним (он в конце файла). Spark делает 2 операции чтения при открытии Parquet: сначала последние 8 байт (размер Footer), потом сам Footer. Только после этого начинает читать данные. При большом числе мелких Parquet-файлов эти 2 × N metadata операций доминируют над временем фактического чтения данных.

Механика Vectorized Reader: батчевое чтение данных

Spark использует специальный Vectorized Parquet Reader при включённой конфигурации spark.sql.parquet.enableVectorizedReader=true (по умолчанию true). Это означает:

Вместо создания Java-объекта для каждой строки (что создаёт огромное GC-давление), Spark читает данные батчами по 4096 строк напрямую в нативную колоночную память (off-heap). Каждый батч - это ColumnarBatch с отдельными массивами для каждой колонки.

Это позволяет:

  • Применять SIMD-инструкции процессора (обработка нескольких значений за одну CPU-инструкцию)
  • Избежать GC-пауз (данные не в Java heap)
  • Применять предикаты к целому батчу за одну операцию (SIMD)
spark = SparkSession.builder \
    # Vectorized Reader включён по умолчанию.
    # Отключите только если видите проблемы с NULL-обработкой или UDF
    .config("spark.sql.parquet.enableVectorizedReader", "true") \

    # Размер батча: 4096 строк - баланс между памятью и эффективностью SIMD.
    # Увеличьте до 8192 если колонок мало и они числовые (лучше SIMD).
    # Уменьшите до 2048 если много широких строк (ограничение памяти Executor).
    .config("spark.sql.parquet.columnarReaderBatchSize", "4096") \

    # Если файлы Parquet содержат глубокие вложенные структуры (ARRAY of STRUCT),
    # Vectorized Reader не работает и Spark падает на Java-ридер автоматически.
    # В этом случае можно явно отключить:
    # .config("spark.sql.parquet.enableVectorizedReader", "false") \
    .getOrCreate()

Ограничения Parquet: где он сдаёт позиции

Parquet отлично подходит для табличных данных с примитивными типами (числа, строки, даты). Но есть сценарии, где он хуже альтернатив:

Глубоко вложенные структуры. Если схема содержит ARRAY<STRUCT<ARRAY<MAP<...>>>>, Parquet сериализует их через сложный алгоритм Dremel encoding (repetition levels + definition levels). Это увеличивает сложность чтения и снижает эффективность Column Pruning.

Очень маленькие файлы. При файлах < 5 MB overhead на чтение Footer (2 HDFS-запроса на каждый файл) начинает доминировать над полезным временем чтения данных.

Write-intensive workloads. При стриминговой записи каждые 30 секунд Parquet требует полного заполнения Row Group (128 MB) для хорошей компрессии. При микробатчах образуются мелкие файлы с неполными Row Groups.


4. Apache ORC: оптимизированный формат для экосистемы Hive

ORC (Optimized Row Columnar) - колоночный формат, созданный в проекте Hive в 2013 году как замена устаревшему RCFile. Разработан командой Hortonworks и оптимизирован под специфические паттерны нагрузки Hive-кластеров.

Анатомия ORC-файла

Ключевые отличия ORC от Parquet:

Row Index на уровне Stripe. ORC хранит Row Index внутри каждого Stripe: min/max значения для каждых 10 000 строк (Row Group внутри Stripe). При запросе с фильтром Spark проверяет Row Index и может пропустить конкретные Row Groups внутри Stripe, не читая их. У Parquet аналогичная функциональность через Page Index (добавлена в Parquet 1.13).

Bloom Filters. ORC по умолчанию строит Bloom Filters для строковых колонок. Bloom Filter позволяет за O(1) проверить «точно ли не существует значение X в этом Row Group». Это особенно эффективно для join-операций: при hash join Spark может пре-фильтровать таблицу на стороне хранилища. В Parquet Bloom Filters добавлены позднее (1.12+) и не включены по умолчанию.

Более агрессивное сжатие. ORC по умолчанию использует Zlib (уровень 6) вместо Snappy. Это даёт на 20-40% лучшее сжатие числовых данных ценой более высокого CPU при чтении/записи.

Lightweight Index (ACID support). ORC нативно поддерживает ACID-семантику в Hive: транзакционные записи, deltas, base files. Это позволяет Hive выполнять INSERT, UPDATE, DELETE с гарантиями транзакций. Для Spark это менее актуально (Delta Lake / Iceberg делают то же самое лучше).

Когда ORC выигрывает у Parquet

На практике выбор между ORC и Parquet определяется не столько техническими характеристиками (они близки), сколько экосистемой:

  • Hive-heavy инфраструктура: если рядом со Spark работает Apache Hive или Apache Impala, ORC обеспечивает более полную совместимость с их ACID-операциями и встроенными оптимизациями.
  • Trino/Presto с Hive Connector: Trino исторически оптимизирован под ORC (Hive Connector использует ORC Writer) и может читать ORC быстрее чем Parquet.
  • Числовые данные с высокой кардинальностью: ORC's integer encoding (RLE v2) часто даёт лучшую компрессию числовых последовательностей.

В чисто Spark-ориентированной инфраструктуре разница между Parquet и ORC минимальна, и Parquet предпочтителен из-за более широкой поддержки в экосистеме (Iceberg, Delta Lake, Hudi преимущественно используют Parquet).


5. Теория сжатия: кодеки и сплинтуемость

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

Компромисс Storage vs Compute

Любой алгоритм сжатия обменивает CPU-время на дисковое пространство. Это всегда компромисс:

Сплинтуемость: фундаментальное требование для HDFS

Сплинтуемость (Splittability) - это возможность параллельно читать разные части одного файла несколькими Spark Task'ами. Это критическое требование для работы в HDFS.

Представьте файл 10 GB в HDFS. HDFS разобьёт его на 80 блоков по 128 MB. Spark должен уметь начать читать данные с начала любого из этих 80 блоков независимо - иначе все 10 GB будет читать один Task, и весь параллелизм теряется.

Для этого кодек должен позволять найти границу сжатого блока в произвольном месте файла.

Несплинтуемые кодеки (raw .gz, raw .bz2):

Файл 10 GB = один непрерывный gzip-поток
Нет маркеров блоков, нельзя начать чтение с середины
→ Один Spark Task читает ВСЮ таблицу
→ Параллелизм = 1, независимо от числа Executor'ов

Сплинтуемые кодеки (в контейнерных форматах):

Snappy, ZSTD, LZ4 и даже Gzip становятся сплинтуемыми внутри Parquet/ORC/Avro, потому что сжатие применяется не к файлу целиком, а к каждому Row Group / Page / Block отдельно. Каждый Row Group - независимо сжатый блок с явным началом. Spark может читать Row Group 5 файла без чтения Row Groups 1-4.

# Демонстрация проблемы несплинтуемых кодеков

# ❌ АНТИПАТТЕРН: CSV + Gzip = не сплинтуемый!
# Один файл 10 GB = одна Spark-задача
df.write \
    .option("compression", "gzip") \
    .csv("hdfs://cluster/data/bad_idea.csv.gz")

# Проверим что получилось:
df_bad = spark.read.csv("hdfs://cluster/data/bad_idea.csv.gz")
print(df_bad.rdd.getNumPartitions())
# Вывод: 1 (одна партиция для всего файла!)

# ✅ ПРАВИЛЬНО: Parquet + Gzip = сплинтуемый!
# (Gzip применяется к каждому Page отдельно)
df.write \
    .option("compression", "gzip") \
    .parquet("hdfs://cluster/data/good.parquet")

df_good = spark.read.parquet("hdfs://cluster/data/good.parquet")
print(df_good.rdd.getNumPartitions())
# Вывод: 80 (80 Row Groups = 80 партиций)

Правило: CSV + Gzip или JSON + Gzip - всегда антипаттерн в Data Lake. Используйте их только в специфических сценариях (передача файла как единого объекта, не для аналитики в Spark).


6. Кодек Snappy: баланс по умолчанию

Snappy создан в Google и спроектирован с одной главной целью: максимальная скорость сжатия и распаковки, не максимальный коэффициент сжатия. Google использует Snappy в Google BigTable, Google Bigtable, LevelDB и многих других внутренних системах.

Характеристики Snappy

Скорость сжатия: 250–500 MB/s на одном ядре CPU (в 10-20 раз быстрее Gzip). Скорость распаковки: 1–2 GB/s на одном ядре CPU (в 5-10 раз быстрее Gzip). Степень сжатия: умеренная. Для текстовых данных: ~40-50% от исходного размера. Для числовых Parquet: ~60-70% от несжатого.

Snappy не является сплинтуемым как отдельный файл (.snappy). Но внутри Parquet/ORC он сплинтуем, потому что каждый Page сжимается отдельно.

Почему Snappy - дефолт в Spark для Parquet

# Snappy - дефолтный кодек для Parquet в Spark:
spark.conf.get("spark.sql.parquet.compression.codec")
# Вывод: snappy

# Это намеренное решение: Snappy даёт предсказуемую производительность
# на горячих данных Silver/Gold, которые читаются сотни раз в день.
# CPU-overhead на распаковку минимален.

# Явное указание:
df.write \
    .mode("overwrite") \
    .option("compression", "snappy") \
    .parquet("hdfs://cluster/data/silver/events/")

# Через глобальную конфигурацию:
spark.conf.set("spark.sql.parquet.compression.codec", "snappy")

Когда Snappy оптимален:

  • Silver и Gold слои Data Lake (часто читаются аналитиками)
  • Интерактивные Spark SQL запросы (низкая latency важнее compression ratio)
  • Кластеры с NVMe дисками (I/O уже быстрый, CPU overhead не компенсируется)
  • Данные с низкой компрессибельностью (уже сжатые изображения, видео)

7. Кодек ZSTD: современный король сжатия

Zstandard (ZSTD) создан Facebook в 2015 году и открыт как open source. ZSTD спроектирован как преемник Zlib/Gzip: предлагает сравнимую с Gzip степень сжатия при скорости, близкой к Snappy.

Уровни сжатия ZSTD

ZSTD поддерживает 22 уровня сжатия (от -5 до 22), что позволяет точно настраивать компромисс:

Уровень Скорость сжатия Скорость распаковки Коэффициент сжатия
1 (минимум) ~500 MB/s ~1.5 GB/s ~40% от оригинала
3 (дефолт) ~300 MB/s ~1.5 GB/s ~35% от оригинала
6 ~150 MB/s ~1.5 GB/s ~30% от оригинала
9 ~80 MB/s ~1.5 GB/s ~27% от оригинала
19 ~5 MB/s ~1.5 GB/s ~22% от оригинала
22 (максимум) ~1 MB/s ~1.5 GB/s ~20% от оригинала

Ключевое свойство: скорость распаковки почти не меняется с уровнем сжатия! Данные, сжатые на уровне 22 (очень медленно), читаются почти так же быстро как сжатые на уровне 1. Это означает: для архивных данных, которые пишутся редко и читаются ещё реже, можно применить максимальное сжатие без штрафа за чтение.

ZSTD в PySpark

spark = SparkSession.builder \
    # ZSTD с дефолтным уровнем (3) для Parquet
    .config("spark.sql.parquet.compression.codec", "zstd") \

    # ZSTD для ORC
    .config("spark.sql.orc.compression.codec", "zstd") \

    .getOrCreate()

# Запись с явным уровнем сжатия ZSTD для Parquet:
df.write \
    .mode("overwrite") \
    .option("compression", "zstd") \
    # Уровень задаётся через отдельную опцию (Spark 3.2+)
    .option("parquet.compression.codec.zstd.level", "9") \
    .parquet("hdfs://cluster/data/cold/events_2022/")

# Для горячих Silver данных (уровень 3 - оптимально):
df_silver.write \
    .mode("overwrite") \
    .option("compression", "zstd") \
    .option("parquet.compression.codec.zstd.level", "3") \
    .partitionBy("event_date") \
    .parquet("hdfs://cluster/data/silver/events/")

Когда ZSTD оптимален:

  • Основной кодек для новых Data Lakehouse на Iceberg/Delta Lake (экономия ~20-30% vs Snappy при минимальном overhead)
  • Архивные данные (уровень 9-15: отличное сжатие, быстрое чтение)
  • Кластеры с HDD (I/O медленный, CPU-оптимизация сжатия оправдана)
  • Gold слой с агрегатами (числовые данные сжимаются отлично)

ZSTD vs Snappy: практическое сравнение на реальных данных:

Датасет Snappy размер ZSTD L3 размер ZSTD L9 размер Чтение Snappy Чтение ZSTD
Транзакции (числа) 2.1 GB 1.4 GB (-33%) 1.2 GB (-43%) 100% 98%
Логи событий (строки) 3.8 GB 2.5 GB (-34%) 2.1 GB (-45%) 100% 96%
Пользовательский профиль 1.5 GB 1.1 GB (-27%) 0.9 GB (-40%) 100% 97%

ZSTD L3 даёт 27-34% экономии дискового пространства при почти нулевом влиянии на скорость чтения. На кластере с 100 TB данных это 27-34 TB экономии - стоимость нескольких серверов.


8. Кодеки Gzip и Bzip2: архивное хранение

Gzip: максимальное сжатие, проблемы со сплинтуемостью

Gzip реализует алгоритм DEFLATE (LZ77 + Huffman coding). Создан в 1992 году, долгое время был стандартом для сжатия файлов в Unix-системах.

Характеристики:

  • Степень сжатия: ~25-35% от оригинального размера (лучше ZSTD L3)
  • Скорость сжатия: ~30-50 MB/s на ядро (медленно)
  • Скорость распаковки: ~200-400 MB/s на ядро (умеренно)
  • Сплинтуемость: нет как standalone файл (.gz)

Ключевая проблема: сырой .gz файл несплинтуем. HDFS разобьёт data.csv.gz на 80 блоков, но Spark не сможет читать их параллельно - он создаст одну Task для чтения всего файла последовательно.

# Антипаттерн: CSV + Gzip = несплинтуемый!
df_bad = spark.read.csv("hdfs://cluster/data/logs.csv.gz")
# HDFS report:
# Файл: 10 GB на диске, но один блок для Spark!
# Spark создаст: 1 Task (не 80!)
# Производительность: 1/80 от максимальной

# НО: Parquet + Gzip = сплинтуемый!
df_ok = spark.read.parquet("hdfs://cluster/data/logs_parquet_gzip/")
# Каждый Row Group сжат отдельным gzip-потоком
# Spark создаст: 80 Tasks
# Производительность: максимальная

Когда Gzip уместен:

  • Файлы CSV/JSON от внешних систем, которые нужно принять «как есть»
  • Когда к файлу никогда не будет параллельного Spark-доступа (только hdfs dfs -get и обработка одним процессом)
  • Архивы логов для передачи по сети (gzip - стандартный HTTP-content-encoding)

Bzip2: нативная сплинтуемость, ужасная производительность

Bzip2 использует алгоритм Burrows-Wheeler Transform + Huffman coding. Интересная особенность: bzip2 нативно сплинтуемый на уровне файла - он делит файл на независимые блоки по 900 KB, каждый из которых сжимается отдельно.

Проблема: bzip2 катастрофически медленный. Скорость сжатия ~5-15 MB/s на ядро (в 20-50 раз медленнее Snappy). Скорость распаковки ~20-50 MB/s.

На практике bzip2 устарел. ZSTD и даже Gzip на уровне Parquet Row Groups дают лучшую производительность при сопоставимой или лучшей степени сжатия.


9. Матрица выбора: сводный гид для Data Engineer

Полная сравнительная таблица

Формат + Кодек Слой Medallion Паттерн нагрузки Сплинтуемость Степень сжатия Скорость Spark Применение
Avro + Snappy Bronze / Ingestion Write-heavy, streaming ✅ (через Sync Markers) ~50-60% Запись быстрая Kafka, CDC, event streaming
Parquet + Snappy Silver / Gold Read-heavy, интерактив ✅ (Row Groups) ~40-60% Чтение быстрое Spark SQL, BI, ML features
Parquet + ZSTD L3 Silver / Gold / Archive Mixed ✅ (Row Groups) ~35-45% Чтение ~97% Snappy Lakehouse, экономия дисков
Parquet + ZSTD L9 Cold Archive Редкое чтение ✅ (Row Groups) ~25-35% Чтение ~95% Snappy Долгосрочное хранение
ORC + ZSTD Silver / Gold (Hive) Read-heavy, Hive+Spark ✅ (Stripes) ~30-40% Хорошо для Hive Trino, Impala, Hive ACID
CSV + Snappy Только staging Обмен с внешними системами ~40-55% Медленно (нет column pruning) Только для обмена, не аналитики
CSV + Gzip ❌ Антипаттерн - ~25-35% Катастрофически медленно Никогда для Data Lake
JSON + Gzip ❌ Антипаттерн - ~25-35% 1 Task на файл Никогда для аналитики

Дерево принятия решений


10. Практика: настройка форматов и кодеков в PySpark

Полный набор конфигураций SparkSession

from pyspark.sql import SparkSession

# ── Сессия для Silver/Gold слоя (горячие данные) ──────────────────────
spark_hot = SparkSession.builder \
    .appName("silver-gold-etl") \

    # Parquet + Snappy: максимальная скорость чтения для интерактивной аналитики
    .config("spark.sql.parquet.compression.codec", "snappy") \

    # Vectorized Reader для колоночного чтения через SIMD
    .config("spark.sql.parquet.enableVectorizedReader", "true") \
    .config("spark.sql.parquet.columnarReaderBatchSize", "4096") \

    # Predicate Pushdown: фильтры применяются на уровне Parquet Footer
    .config("spark.sql.parquet.filterPushdown", "true") \

    # Row Group статистики: Spark использует min/max для пропуска Row Groups
    .config("spark.sql.parquet.recordLevelFilter.enabled", "true") \

    # Размер Row Group = размер блока HDFS = одна Task = оптимальная locality
    # 128 MB - дефолт Parquet, должен совпадать с dfs.blocksize
    .config("spark.hadoop.parquet.block.size", "134217728") \

    .getOrCreate()


# ── Сессия для Cold Archive (архивные данные) ─────────────────────────
spark_cold = SparkSession.builder \
    .appName("archive-etl") \

    # Parquet + ZSTD для максимального сжатия при хорошей скорости чтения
    .config("spark.sql.parquet.compression.codec", "zstd") \

    # ORC для архивов в Hive-экосистеме
    .config("spark.sql.orc.compression.codec", "zstd") \

    .getOrCreate()

# Уровень ZSTD задаётся при записи:
df_cold.write \
    .option("compression", "zstd") \
    .option("parquet.compression.codec.zstd.level", "9") \
    .parquet("hdfs://cluster/data/archive/2022/")


# ── Сессия для Avro (Bronze/Ingestion) ────────────────────────────────
spark_bronze = SparkSession.builder \
    .appName("kafka-bronze-ingestion") \
    .config("spark.jars.packages", "org.apache.spark:spark-avro_2.12:3.5.1") \
    .getOrCreate()

Тюнинг размеров Row Group и Stripe под HDFS

# Ключевой принцип: Row Group size = HDFS block size
# Это обеспечивает: 1 Row Group = 1 HDFS Block = 1 Spark Task = NODE_LOCAL чтение

# Если HDFS blocksize = 128 MB (стандарт):
spark.conf.set("spark.hadoop.parquet.block.size", "134217728")  # 128 MB

# Если HDFS blocksize = 256 MB (крупные файлы, меньше Tasks):
spark.conf.set("spark.hadoop.parquet.block.size", "268435456")  # 256 MB

# Для ORC:
spark.conf.set("spark.hadoop.orc.stripe.size", "67108864")  # 64 MB (ORC default)
# Или увеличить для крупных таблиц:
spark.conf.set("spark.hadoop.orc.stripe.size", "134217728")  # 128 MB

# Размер Page внутри Row Group (влияет на granularity статистики):
# Меньше Page = точнее predicate pushdown, но больше метаданных
spark.conf.set("spark.hadoop.parquet.page.size", "1048576")  # 1 MB (default)

Практические рецепты для каждого слоя Medallion Architecture

from pyspark.sql import SparkSession, functions as F
from datetime import date

spark = SparkSession.builder \
    .appName("medallion-formats-demo") \
    .enableHiveSupport() \
    .getOrCreate()

# ── Bronze Layer: Avro + Snappy (источник - Kafka с Debezium CDC) ─────
def write_bronze_avro(df, target_path: str, date_str: str) -> None:
    """
    Bronze слой: сохраняем данные максимально близко к источнику.
    Avro + Snappy оптимален для потоковой записи из Kafka/Debezium:
    - Схема хранится в файле (самодостаточный файл)
    - Schema Evolution (добавление колонок в исходной системе = OK)
    - Snappy: минимальный CPU overhead при высокоскоростной записи
    """
    df.write \
        .format("avro") \
        .mode("append") \
        .option("compression", "snappy") \
        .partitionBy("ingestion_date") \
        .save(f"{target_path}/date={date_str}/")


# ── Silver Layer: Parquet + Snappy (горячие очищенные данные) ─────────
def write_silver_parquet(df, table_name: str, date_str: str) -> None:
    """
    Silver слой: очищенные и нормализованные данные.
    Parquet + Snappy: оптимально для Spark SQL аналитики.
    - Columnar: Column Pruning (Spark читает только нужные колонки)
    - Snappy: быстрая распаковка при интерактивных запросах
    - Vectorized Reader: SIMD-оптимизация
    - partitionBy(date): partition pruning при ежедневных запросах
    """
    spark.conf.set("spark.sql.parquet.compression.codec", "snappy")

    df.repartition(50) \
        .write \
        .mode("overwrite") \
        .format("parquet") \
        .option("compression", "snappy") \
        .partitionBy("event_date") \
        .saveAsTable(f"silver.{table_name}")


# ── Gold Layer: Parquet + Snappy или ZSTD ──────────────────────────────
def write_gold_parquet(df, table_name: str, use_zstd: bool = False) -> None:
    """
    Gold слой: агрегаты и витрины для BI.
    - Если данные часто читаются → Snappy
    - Если данные большие и читаются реже → ZSTD (экономия 30%)
    """
    codec = "zstd" if use_zstd else "snappy"

    # Gold обычно небольшой после агрегации → меньше партиций
    df.coalesce(10) \
        .write \
        .mode("overwrite") \
        .format("parquet") \
        .option("compression", codec) \
        .saveAsTable(f"gold.{table_name}")


# ── Cold Archive: Parquet + ZSTD Level 9 ──────────────────────────────
def archive_to_cold_storage(source_table: str, target_path: str,
                             year: int, month: int) -> int:
    """
    Архивирование исторических данных в Cold Storage.
    Parquet + ZSTD Level 9:
    - Максимальное сжатие (читаются раз в квартал)
    - Скорость чтения почти не страдает (ZSTD decode быстрый)
    - Экономия 40-50% vs Snappy на архивных числовых данных
    """
    df = spark.table(source_table) \
        .filter(F.year("event_date") == year) \
        .filter(F.month("event_date") == month)

    row_count = df.count()

    # Максимальный ZSTD: пишется медленно, читается быстро
    df.repartition(20) \
        .write \
        .mode("overwrite") \
        .format("parquet") \
        .option("compression", "zstd") \
        .option("parquet.compression.codec.zstd.level", "9") \
        .save(f"{target_path}/year={year}/month={month:02d}/")

    print(f"Архивировано {row_count:,} строк за {year}-{month:02d}")
    print(f"  Codec: ZSTD Level 9")
    return row_count

Проверка реального сжатия через HDFS CLI

# Сравниваем физический размер файлов в разных форматах/кодеках
# После записи одного датасета в четырёх вариантах:

hdfs dfs -du -h /data/benchmark/
# Вывод (пример для 100M строк транзакционных данных):
# 3.2 G   9.6 G   /data/benchmark/json_no_compression    # JSON без сжатия
# 1.8 G   5.4 G   /data/benchmark/csv_gzip               # CSV + Gzip
# 1.5 G   4.5 G   /data/benchmark/parquet_snappy         # Parquet + Snappy
# 1.1 G   3.3 G   /data/benchmark/parquet_zstd_3         # Parquet + ZSTD L3
# 0.9 G   2.7 G   /data/benchmark/parquet_zstd_9         # Parquet + ZSTD L9
# 0.8 G   2.4 G   /data/benchmark/orc_zstd               # ORC + ZSTD

# Второе число = размер × replication factor (3x)
# Для экономии дисков смотрим на первое число

11. Лабораторная работа: Великий батл форматов

Следующий код запускает сравнительный бенчмарк всех форматов и кодеков на реальных данных.

# lab_file_formats_benchmark.py
# Сравниваем Avro, Parquet, ORC с разными кодеками

import time
import subprocess
from dataclasses import dataclass
from pyspark.sql import SparkSession, functions as F


@dataclass
class BenchmarkResult:
    name: str
    format: str
    codec: str
    write_time_sec: float
    disk_size_mb: float
    read_full_time_sec: float
    read_filtered_time_sec: float
    compression_ratio: float


def create_spark() -> SparkSession:
    return (
        SparkSession.builder
        .master("yarn")
        .appName("file-formats-benchmark")
        .config("spark.jars.packages",
                "org.apache.spark:spark-avro_2.12:3.5.1")
        .config("spark.executor.instances", "10")
        .config("spark.executor.cores", "4")
        .config("spark.executor.memory", "8g")
        .getOrCreate()
    )


def generate_dataset(spark: SparkSession, n_rows: int = 10_000_000):
    """Генерирует тестовый датасет транзакций."""
    return spark.range(n_rows).select(
        F.col("id").alias("txn_id"),
        (F.rand() * 1_000_000).cast("long").alias("user_id"),
        (F.rand() * 50_000).cast("long").alias("merchant_id"),
        F.array(
            F.lit("purchase"), F.lit("refund"),
            F.lit("transfer"), F.lit("fee")
        ).getItem((F.rand() * 4).cast("int")).alias("txn_type"),
        (F.rand() * 10_000).alias("amount"),
        F.array(F.lit("RUB"), F.lit("USD"), F.lit("EUR")).getItem(
            (F.rand() * 3).cast("int")
        ).alias("currency"),
        F.date_add(F.lit("2024-01-01"), (F.rand() * 365).cast("int")).alias("txn_date"),
        F.array(
            F.lit("web"), F.lit("mobile"), F.lit("pos"), F.lit("api")
        ).getItem((F.rand() * 4).cast("int")).alias("channel"),
    )


def get_hdfs_size_mb(path: str) -> float:
    """Получает реальный размер директории на HDFS."""
    result = subprocess.run(
        ["hdfs", "dfs", "-du", "-s", path],
        capture_output=True, text=True
    )
    parts = result.stdout.split()
    if len(parts) >= 1:
        return int(parts[0]) / 1024 / 1024
    return 0.0


def run_benchmark(
    spark: SparkSession,
    df,
    name: str,
    fmt: str,
    codec: str,
    base_path: str,
) -> BenchmarkResult:
    """Запускает один бенчмарк: пишет, измеряет размер, читает."""
    path = f"{base_path}/{name}/"
    print(f"\nБенчмарк: {name} (format={fmt}, codec={codec})")

    # Запись
    t0 = time.time()
    if fmt == "avro":
        df.write.format("avro").mode("overwrite") \
            .option("compression", codec).save(path)
    else:
        df.write.format(fmt).mode("overwrite") \
            .option("compression", codec).save(path)
    write_time = time.time() - t0
    print(f"  Запись: {write_time:.1f}с")

    # Размер на диске
    disk_mb = get_hdfs_size_mb(path)
    original_mb = get_hdfs_size_mb(f"{base_path}/parquet_snappy/")
    compression_ratio = disk_mb / original_mb if original_mb > 0 else 1.0
    print(f"  Размер: {disk_mb:.0f} MB ({compression_ratio:.2f}x vs Parquet+Snappy)")

    # Полное сканирование
    t0 = time.time()
    if fmt == "avro":
        spark.read.format("avro").load(path).agg(F.sum("amount")).collect()
    else:
        spark.read.format(fmt).load(path).agg(F.sum("amount")).collect()
    read_full = time.time() - t0
    print(f"  Полное чтение (SUM): {read_full:.1f}с")

    # Фильтрованный запрос (проверяем predicate pushdown + column pruning)
    t0 = time.time()
    if fmt == "avro":
        spark.read.format("avro").load(path) \
            .filter(F.col("txn_type") == "purchase") \
            .groupBy("currency") \
            .agg(F.sum("amount"), F.count("*")) \
            .collect()
    else:
        spark.read.format(fmt).load(path) \
            .filter(F.col("txn_type") == "purchase") \
            .groupBy("currency") \
            .agg(F.sum("amount"), F.count("*")) \
            .collect()
    read_filtered = time.time() - t0
    print(f"  Фильтрованное чтение (WHERE + GROUP BY): {read_filtered:.1f}с")

    return BenchmarkResult(
        name=name, format=fmt, codec=codec,
        write_time_sec=write_time,
        disk_size_mb=disk_mb,
        read_full_time_sec=read_full,
        read_filtered_time_sec=read_filtered,
        compression_ratio=compression_ratio,
    )


if __name__ == "__main__":
    spark = create_spark()
    BASE_PATH = "hdfs://cluster/benchmark/file-formats"
    df = generate_dataset(spark, n_rows=10_000_000)

    # Кешируем в памяти чтобы генерация не влияла на write benchmark
    df.cache()
    df.count()

    configs = [
        ("parquet_snappy",  "parquet", "snappy"),
        ("parquet_zstd_3",  "parquet", "zstd"),
        ("parquet_gzip",    "parquet", "gzip"),
        ("orc_snappy",      "orc",     "snappy"),
        ("orc_zstd",        "orc",     "zstd"),
        ("avro_snappy",     "avro",    "snappy"),
    ]

    results = []
    for name, fmt, codec in configs:
        r = run_benchmark(spark, df, name, fmt, codec, BASE_PATH)
        results.append(r)

    # Итоговая таблица
    print("\n" + "=" * 80)
    print(f"{'Формат':20s} {'Запись':>8} {'Диск MB':>9} {'Ratio':>6} {'Полн. чт':>10} {'Фильтр':>10}")
    print("-" * 80)
    for r in sorted(results, key=lambda x: x.read_filtered_time_sec):
        print(f"{r.name:20s} {r.write_time_sec:>7.1f}s {r.disk_size_mb:>8.0f} "
              f"{r.compression_ratio:>5.2f}x {r.read_full_time_sec:>9.1f}s "
              f"{r.read_filtered_time_sec:>9.1f}s")

    spark.stop()

Итоги: правило выбора формата в 2026 году

Формат хранения данных - это архитектурное решение с долгосрочными последствиями. Менять формат миллиардов строк в production - это дни работы и значительные риски. Поэтому важно принять правильное решение в начале.

Консенсус индустрии 2026 года:

  • Avro - для приёма потоковых данных из Kafka, Debezium CDC, WebHooks. Bronze слой. Schema Evolution из коробки. Кодек: Snappy.
  • Parquet - основной формат Data Lake и Lakehouse для Silver и Gold слоёв. Интерактивная аналитика, Spark SQL, Delta Lake, Iceberg. Кодек: Snappy (горячие), ZSTD L3 (баланс), ZSTD L9 (архив).
  • ORC - отличный выбор для инфраструктур с Hive, Trino, Impala. Лучше встроенная поддержка ACID в Hive. Кодек: ZSTD.
  • CSV/JSON - только для обмена данными с внешними системами. Никогда - как основной формат хранения в Data Lake.
  • Snappy - горячие данные, интерактивная аналитика. Минимальный CPU overhead.
  • ZSTD - универсальный выбор нового поколения. 30% экономии дисков vs Snappy при минимальном штрафе за чтение.
  • Gzip - только в Parquet/ORC (не как raw файл). Для архивов где важна максимальная компрессия.

Золотое правило для Spark + Iceberg + HDFS Data Lakehouse:

Bronze:  Avro + Snappy    (быстрая запись из потока, schema evolution)
Silver:  Parquet + Snappy (быстрое чтение, интерактивная аналитика)
Gold:    Parquet + ZSTD   (экономия дисков, хорошее чтение)
Archive: Parquet + ZSTD 9 (максимальная компрессия, редкие чтения)