Форматы файлов в 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.
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 (максимальная компрессия, редкие чтения)