Data Engineering и ETL Хранилища и форматы: Parquet, колоночные БД, lakehouse
0%

Хранилища и форматы: Parquet, колоночные БД, lakehouse

Хранилища и форматы: Parquet, колоночные БД, lakehouse

Все предыдущие статьи трека говорили о движении данных: как их извлечь (https://courses.digitable.life/post/data-engineering/01-etl-vs-elt/), в какую форму привести (https://courses.digitable.life/post/data-engineering/02-data-modeling/), чем посчитать (https://courses.digitable.life/post/data-engineering/03-batch-processing/, https://courses.digitable.life/post/data-engineering/04-streaming/) и как запустить по расписанию (https://courses.digitable.life/post/data-engineering/05-orchestration/). Эта статья — про то, где данные лежат между запусками, и это не техническая деталь: формат хранения определяет и скорость запросов, и счёт за облако, и то, можно ли вообще безопасно писать в таблицу двумя процессами одновременно.

Правило, вокруг которого всё вертится: самый быстрый способ прочитать данные — не читать их. Колоночные форматы, партиционирование, статистики, индексы и табличные форматы — это всё разные способы сказать движку «эти байты можно пропустить». Дальше мы разберём каждый механизм до уровня «что именно лежит на диске и какой байт кто читает первым».


1. Первый принцип: почему раскладка байт решает всё

На диске и в сети данные линейны — это последовательность байт. Таблица двумерна. Значит, кто-то должен выбрать порядок обхода: сначала по строкам или сначала по колонкам. Именно этот выбор — водораздел между OLTP и OLAP.

Строчная и колоночная раскладка одной таблицы

Аналитический запрос почти всегда выглядит так: «просканируй сотни миллионов строк, но возьми 3 колонки из 200 и посчитай агрегат». Транзакционный — так: «найди одну строку по ключу и верни/измени её целиком». Первому нужна колоночная раскладка, второму — строчная. Это не мода, а прямое следствие того, какие байты нужны.

Три эффекта колоночной раскладки, и их стоит различать:

  1. Projection pushdown. Читаются только запрошенные колонки. Таблица из 200 колонок, запрос по трём — экономия ~98% ввода-вывода. Это самый крупный выигрыш, и он бесплатный.
  2. Сжатие в разы лучше. В колонке соседние значения однотипны и часто повторяются: country — это 200 уникальных строк на миллиард записей, event_type — десяток. Словарное кодирование + RLE ужимают такую колонку в 50–100 раз. В строчной раскладке рядом лежат int, string и timestamp, и общий компрессор не находит закономерности.
  3. Векторное выполнение. Колонка amount — это плотный массив int64 в памяти. По нему процессор идёт с предсказуемым префетчем и SIMD-инструкциями, обрабатывая 4–8 значений за такт. Строчная раскладка даёт кэш-промах почти на каждой записи.

Порядок величин, которые полезно держать в голове: чтение с NVMe — единицы ГБ/с, чтение из S3 одним потоком — сотни МБ/с (и растёт линейно с числом параллельных запросов), чтение из RAM — десятки ГБ/с, декомпрессия zstd — порядка 1 ГБ/с на ядро, snappy — 2–4 ГБ/с на ядро. Отсюда важное следствие: при чтении из объектного хранилища узкое место — сеть, и почти любое сжатие окупается; при чтении из локального RAM-кэша сильное сжатие может стать узким местом, и снаппи выигрывает.


2. Кодирование и сжатие — это разные вещи

Частая путаница. Кодирование (encoding) знает тип и семантику колонки и работает до сжатия. Компрессия — это общий алгоритм над байтами (snappy, zstd, gzip, lz4). Их эффекты перемножаются, и первый обычно важнее второго.

Основные кодировки, которые применяют Parquet, ORC и колоночные СУБД:

Кодировка Как работает Где выигрывает
Dictionary Уникальные значения в словарь, в данных — индексы Низкая кардинальность: страны, статусы, UUID-справочники
RLE (run-length) RU RU RU RU(RU, 4) Отсортированные или редко меняющиеся колонки
Bit-packing Индексы 0..7 хранятся по 3 бита, а не по 32 Любые словарные индексы, флаги
Delta Хранится разница с предыдущим значением Возрастающие ключи, timestamp, счётчики
Delta length / byte-stream-split Отдельно длины строк и байты; байты float по позициям Строки, float-колонки метрик
Frame of reference (FOR) Значения минус минимум блока Числа в узком диапазоне

Ключевая мысль: эффективность кодирования зависит от порядка строк. Одни и те же данные, отсортированные по country, а не по случайному user_id, сжимаются кратно лучше — RLE наконец находит длинные серии. Поэтому сортировка перед записью (ORDER BY/sortWithinPartitions) — это не косметика, а способ снизить счёт за хранение и ускорить чтение.

# Наглядный эксперимент: одни и те же данные, разный порядок строк.
import numpy as np, pyarrow as pa, pyarrow.parquet as pq

n = 5_000_000
rng = np.random.default_rng(42)
country = rng.choice(["RU", "US", "DE", "FR", "BR", "IN"], size=n)
amount = rng.integers(1, 500, size=n)

tbl = pa.table({"country": country, "amount": amount})
pq.write_table(tbl, "/tmp/unsorted.parquet", compression="zstd")

# Та же таблица, отсортированная по country
tbl_sorted = tbl.sort_by([("country", "ascending")])
pq.write_table(tbl_sorted, "/tmp/sorted.parquet", compression="zstd")

# На типичных данных: unsorted ~9 МБ, sorted ~4 МБ — только за счёт RLE по country
# и за счёт того, что min/max статистики стали селективными.

Про компрессоры коротко и по делу:

  • snappy — исторический дефолт Parquet. Быстрый, ратио посредственное. Берите, когда данные читаются с локальных дисков и CPU дороже байт.
  • zstd (уровень 1–3) — сегодня разумный дефолт для озера. На тех же данных даёт на 20–40% меньше байт, чем snappy, при сопоставимой скорости декомпрессии. Для холодных данных — уровень 9–12.
  • gzip — не берите: медленнее zstd при худшем ратио.
  • lz4 — когда критична латентность декомпрессии (горячий кэш, ClickHouse по умолчанию).

3. Parquet изнутри

Parquet (официальная спецификация) — это самодостаточный файл: он несёт схему, статистики и данные. Иерархия ровно такая: файл → row group → column chunk → page.

Анатомия файла Parquet

  • Row group — горизонтальный срез: несколько миллионов строк, все колонки. Целевой размер — 128–512 МБ (несжатых). Row group — минимальная единица параллелизма: один воркер Spark берёт один row group.
  • Column chunk — все значения одной колонки внутри row group, лежащие непрерывно. Именно поэтому можно прочитать колонку одним диапазонным GET.
  • Page — единица сжатия и кодирования, по умолчанию около 1 МБ. Внутри — dictionary page и data pages.
  • Footer — Thrift-структура в конце файла: схема, для каждого column chunk смещение, размер, кодировки и статистики (min, max, null_count), плюс опциональные ColumnIndex/OffsetIndex (те же min/max на уровне страниц) и Bloom-фильтры.

Почему footer в конце: писателю не нужно знать статистики заранее — он пишет данные потоком, а метаданные складывает в конце. Читателю поэтому нужно два обращения: прочитать последние 8 байт (длина footer + магия), затем сам footer. Отсюда практическое следствие для S3: чтение одного Parquet-файла — это минимум 2 GET, и когда файлов миллион, вы платите не за байты, а за запросы.

Как выглядит pruning на практике

import pyarrow.parquet as pq

pf = pq.ParquetFile("/tmp/sorted.parquet")
print(pf.metadata)              # число строк, row groups, версия писателя
print(pf.schema_arrow)

rg = pf.metadata.row_group(0)
for i in range(rg.num_columns):
    col = rg.column(i)
    print(
        col.path_in_schema,
        col.compression,
        col.encodings,
        col.total_compressed_size,
        col.statistics.min if col.statistics else None,
        col.statistics.max if col.statistics else None,
    )
# Именно эти min/max движок сравнивает с предикатом WHERE и решает,
# читать ли row group вообще.

Уровни отсечения выстраиваются в воронку — каждый следующий дешевле предыдущего только потому, что предыдущий уже отбросил большую часть данных:

Сложность запроса-скана честно записывается так:

Прочитанные байты ≈ N_партиций_после_отсечения
                  × N_файлов_в_партиции
                  × доля_row_groups_прошедших_статистики
                  × (размер_нужных_колонок / размер_всех_колонок)
                  × коэффициент_сжатия

Все пять множителей — рычаги, и все под вашим контролем: партиционирование, компакция, сортировка, ширина проекции, компрессор. Оптимизация «сделаем кластер побольше» не входит ни в один из них.

Вложенные структуры и модель Dremel

Parquet умеет хранить вложенность (struct, list, map), и делает это не наивно. Каждое значение листовой колонки сопровождается двумя числами — definition level (на каком уровне вложенности значение стало null) и repetition level (на каком уровне началось новое повторение). Эта схема описана в статье Google Dremel: Interactive Analysis of Web-Scale Datasets (VLDB 2010) и позволяет разложить произвольно вложенный JSON в плоские колонки без потери информации.

Практический вывод: не храните аналитические данные как строку JSON в одной колонке. Вы теряете и projection pushdown, и статистики, и сжатие по типу. Если схема нестабильна — разложите стабильную часть в колонки, а «хвост» оставьте в JSON-поле (или используйте тип variant, который сейчас появляется в спецификациях Parquet и Iceberg).


4. Форматы файлов: Parquet, ORC, Avro, Arrow, CSV

Формат Раскладка Схема Где применять Где не применять
Parquet колоночная в файле аналитика, озеро, витрины, обмен между движками построчные апдейты, потоковая дозапись мелкими порциями
ORC колоночная в файле экосистема Hive/Trino, где он исторически новые проекты вне Hive — экосистема Parquet шире
Avro строчная в файле + Schema Registry landing-зона, CDC, сообщения Kafka, манифесты Iceberg аналитические сканы по подмножеству колонок
Arrow / Arrow IPC колоночная в схеме обмен в памяти между процессами, zero-copy, Flight долговременное хранение (нет сжатия по умолчанию, формат не оптимизирован под диск)
CSV / JSONL строчная нет обмен с внешним миром, отладка, ручная выгрузка всё остальное

Разница между Parquet и ORC сегодня невелика: у ORC исторически лучше встроенные Bloom-фильтры и row-index на 10 000 строк, у Parquet — существенно более широкая поддержка (pandas, DuckDB, Polars, BigQuery, Snowflake, Spark, Trino, ClickHouse). Если нет исторической причины — Parquet.

Разница между Parquet и Arrow принципиальна и её часто путают: Parquet оптимизирован под минимальный объём на диске (сжатие, кодирование, дорогая распаковка), Arrow — под мгновенный доступ в памяти (фиксированная раскладка, никакой распаковки, zero-copy между процессами). Они дополняют друг друга: читаем Parquet → материализуем в Arrow → считаем. Спецификация: arrow.apache.org/docs/format.

Avro остаётся правильным выбором там, где данные пишутся построчно и читаются целиком: приём CDC-потока, «сырой» слой, сообщения. Он же используется внутри Iceberg для манифестов — именно потому, что манифест читается целиком.


5. Колоночные СУБД: когда файлов недостаточно

Файлы в озере хороши для больших сканов, но плохи для запроса «покажи дашборд за 200 мс». Колоночные СУБД — это те же принципы (колонки, кодирование, векторизация), но с собственным управлением раскладкой, индексами и мержами.

ClickHouse: MergeTree как LSM для аналитики

Таблица MergeTree состоит из партов (parts) — независимых директорий с колоночными файлами. Каждый INSERT создаёт новый парт, фоновые мержи объединяют мелкие парты в крупные. Внутри парта строки отсортированы по ORDER BY, а первичный индекс разрежённый: одна запись на каждые index_granularity строк (по умолчанию 8192). То есть индекс для миллиарда строк занимает мегабайты и живёт в памяти, а поиск даёт не строку, а гранулу, которую надо дочитать.

-- Витрина событий в ClickHouse: ключ сортировки важнее, чем всё остальное
CREATE TABLE events
(
    event_time  DateTime,
    dt          Date MATERIALIZED toDate(event_time),
    country     LowCardinality(String),   -- словарное кодирование на уровне типа
    user_id     UInt64,
    event_type  LowCardinality(String),
    amount      Decimal(18, 2)
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(event_time)         -- крупная партиция: месяц, не день и не час
ORDER BY (country, event_type, event_time)  -- от низкой кардинальности к высокой
TTL event_time + INTERVAL 18 MONTH        -- автоматическое удаление старого
SETTINGS index_granularity = 8192;

-- Пропускающий индекс для колонки, которой нет в ключе сортировки
ALTER TABLE events
  ADD INDEX idx_user user_id TYPE bloom_filter(0.01) GRANULARITY 4;

Три правила ClickHouse, которые экономят недели:

  1. ORDER BY — от низкой кардинальности к высокой. Тогда префикс ключа даёт длинные серии для RLE и селективные засечки в индексе.
  2. Партиций должно быть мало (десятки–сотни). PARTITION BY по дню на трёхлетней истории — это 1000+ партиций и деградация мержей; месяц почти всегда правильнее.
  3. Вставляйте большими батчами (десятки–сотни тысяч строк). Тысяча мелких INSERT в секунду породит тысячу партов, и сервер утонет в мержах — та же болезнь мелких файлов, что и в озере.

DuckDB: «SQLite для аналитики»

DuckDB — встраиваемая векторная СУБД, которая умеет читать Parquet прямо с диска и из S3 без загрузки. Это лучший инструмент для отладки пайплайнов и для витрин в единицы–десятки гигабайт: ноутбук на 16 ГБ RAM спокойно агрегирует сотни миллионов строк.

-- Проверить раскладку озера, не поднимая ни одного кластера
SELECT
    file_name,
    row_group_id,
    num_rows,
    total_compressed_size / 1024 / 1024 AS mb
FROM parquet_metadata('s3://lake/events/dt=2026-07-15/*.parquet')
ORDER BY mb;

-- Сколько у нас «мелких файлов»?
SELECT
    count(*) FILTER (WHERE size_mb < 16) AS small_files,
    count(*)                             AS total_files,
    round(avg(size_mb), 1)               AS avg_mb
FROM (
    SELECT file_name, sum(total_compressed_size) / 1048576.0 AS size_mb
    FROM parquet_metadata('s3://lake/events/dt=2026-07-15/*.parquet')
    GROUP BY file_name
);

Идея, которую стоит унести: прежде чем оптимизировать Spark-джобу, посмотрите на метаданные Parquet через DuckDB. В 80% случаев проблема видна там — мелкие файлы, отсутствие сортировки, бесполезные min/max (когда min = 0, max = 2^63 в каждом файле, потому что данные записаны в случайном порядке).


6. Почему «просто файлы в S3» ломаются

Классическая раскладка Hive — это каталог с партициями: s3://lake/events/dt=2026-07-15/part-0001.parquet. Схема таблицы и список партиций живут в Hive Metastore, а данные — в объектном хранилище. Пока данные только дописываются, это работает. Ломается оно в четырёх местах:

  1. Нет атомарности. Джоба пишет 400 файлов в 30 партиций. На 200-м файле она падает. Читатель, пришедший в этот момент, видит половину результата и не может об этом узнать. «Записать во временный каталог и переименовать» тоже не спасает: в объектном хранилище нет атомарного rename каталога — это копирование объектов по одному.
  2. Дорогой листинг. Чтобы понять, какие файлы читать, движок делает LIST по префиксам. Таблица на 500 000 файлов — это тысячи LIST-запросов и десятки секунд только на планирование, до чтения первого байта данных.
  3. Партиции зашиты в путь. Захотели перейти с дня на час — надо переписать всю таблицу и все запросы. А ещё пользователь, забывший написать WHERE dt = ..., случайно сканирует три года истории.
  4. Нет апдейтов и удалений. GDPR-запрос «удалить пользователя» превращается в перезапись всех партиций, где он встречался.

Все четыре проблемы решает один и тот же приём: добавить над файлами слой метаданных, который знает точный список файлов таблицы и меняется атомарно. Это и есть табличный формат — и есть суть lakehouse.


7. Lakehouse: Iceberg, Delta Lake, Hudi

Lakehouse = дешёвое объектное хранилище + открытые файловые форматы + табличный формат, дающий ACID, эволюцию схемы и time travel. Термин закрепила статья Lakehouse: A New Generation of Open Platforms that Unify Data Warehousing and Advanced Analytics (CIDR 2021).

Модель метаданных Iceberg

Читайте эту схему снизу вверх, и вся конструкция станет очевидной. DATA_FILE несёт min/max по колонкам — значит, планировщик отсекает файлы не листингом каталогов, а чтением манифеста. Манифесты перечислены в MANIFEST_LIST, тот принадлежит снапшоту, снапшот — версии метаданных, а каталог хранит один указатель на текущую версию. Коммит = атомарная подмена этого указателя. Всё.

Отсюда автоматически следуют все свойства:

  • ACID. Читатель, зафиксировавший snapshot_id, видит согласованный набор файлов, что бы ни писали параллельно (snapshot isolation).
  • Time travel. Старые снапшоты никуда не делись — можно прочитать таблицу «как на вчера».
  • Эволюция схемы. Колонки идентифицируются field_id, а не позицией и не именем: переименование безопасно, удаление и добавление не ломают старые файлы.
  • Скрытое партиционирование (hidden partitioning). Пользователь пишет WHERE event_time > ..., а Iceberg сам сопоставляет предикат с трансформацией days(event_time). Партиции не зашиты в путь и их схему можно менять — старые данные останутся в старой схеме, новые пойдут в новую.

Коммит: оптимистичная конкуренция

Ключевое: конкуренция оптимистичная, блокировок нет, а вся атомарность сводится к одной compare-and-swap операции в каталоге. Именно поэтому каталог обязан её поддерживать: обычный файл в S3 исторически не годился, отсюда Glue, Hive Metastore с транзакцией, Nessie, Polaris и стандарт Iceberg REST Catalog. И именно поэтому конфликтующие писатели в одну партицию — плохая идея: append-операции почти никогда не конфликтуют, а вот два одновременных overwrite/MERGE в одну партицию будут ретраиться, пока один не сдастся.

Copy-on-write и merge-on-read

Апдейт строки в неизменяемом файле невозможен — файл придётся переписать. Отсюда две стратегии:

Copy-on-write (CoW) Merge-on-read (MoR)
Что делает запись переписывает файлы целиком пишет маленький delete-файл / deletion vector
Стоимость записи высокая низкая
Стоимость чтения обычная выше: нужно применять удаления
Когда брать редкие апдейты, много чтений (витрины) частые апдейты, CDC-приём, стриминг

Delta Lake и Iceberg поддерживают обе; Hudi исторически строился вокруг upsert-сценария и record-level индекса. В свежих версиях обеих экосистем MoR реализован через deletion vectors — компактный bitmap удалённых позиций вместо файла со списком строк; это заметно дешевле при чтении.

Delta Lake: то же самое, другой лог

В Delta метаданные — это каталог _delta_log/ с JSON-файлами по одному на коммит (00000000000000000042.json), где записаны действия add/remove/metaData/protocol. Каждые 10 коммитов складывается контрольная точка в Parquet, чтобы не переигрывать историю с нуля. Атомарность коммита обеспечивается атомарным созданием файла с номером N (put-if-absent). Спецификация: Delta Lake protocol.

Практическая разница между Iceberg и Delta в 2026 году невелика: обе дают ACID, time travel, эволюцию схемы, обе читаются основными движками. Выбор чаще диктуется экосистемой (Databricks — Delta; мультидвижковое озеро на Trino/Flink/Spark — чаще Iceberg) и наличием REST-каталога.

Жизненный цикл файла данных

Из этой диаграммы вытекает главная эксплуатационная обязанность владельца lakehouse-таблицы: помимо записи данных надо регулярно запускать три служебные операции — компакцию, expire_snapshots и remove_orphan_files. Без первой деградируют запросы, без второй и третьей растёт счёт за хранение.


8. Рабочий пример: таблица Iceberg от создания до обслуживания

-- Spark SQL / Trino: создание таблицы со скрытым партиционированием
CREATE TABLE lake.analytics.events (
    event_id     STRING,
    event_time   TIMESTAMP,
    user_id      BIGINT,
    country      STRING,
    event_type   STRING,
    amount       DECIMAL(18, 2)
)
USING iceberg
PARTITIONED BY (days(event_time), bucket(16, user_id))
TBLPROPERTIES (
    'write.format.default'          = 'parquet',
    'write.parquet.compression-codec' = 'zstd',
    'write.target-file-size-bytes'  = '536870912',   -- целимся в 512 МБ на файл
    'write.distribution-mode'       = 'hash',        -- убираем мелкие файлы на записи
    'format-version'                = '2'
);

-- Идемпотентная перезапись одного дня: MERGE вместо DELETE + INSERT
MERGE INTO lake.analytics.events t
USING staging.events_2026_07_15 s
ON  t.event_id = s.event_id
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *;

-- Time travel: что было в таблице до вчерашней загрузки
SELECT count(*) FROM lake.analytics.events
FOR TIMESTAMP AS OF TIMESTAMP '2026-07-14 23:59:59';

-- Диагностика: сколько файлов и какого размера в проблемной партиции
SELECT
    partition,
    count(*)                          AS files,
    round(avg(file_size_in_bytes)/1048576, 1) AS avg_mb,
    sum(record_count)                 AS rows
FROM lake.analytics.events.files
GROUP BY partition
ORDER BY files DESC
LIMIT 20;

Обслуживание — обычно отдельный DAG в Airflow (см. https://courses.digitable.life/post/data-engineering/05-orchestration/), который бежит ночью:

# Служебные процедуры Iceberg, вызываемые из Spark.
# Запускать по расписанию для КАЖДОЙ активно пишущейся таблицы.

TABLE = "lake.analytics.events"

# 1. Компакция: слить мелкие файлы и заодно отсортировать данные внутри партиции.
#    Сортировка по country делает min/max селективными и улучшает RLE.
spark.sql(f"""
    CALL lake.system.rewrite_data_files(
        table => '{TABLE}',
        strategy => 'sort',
        sort_order => 'country ASC NULLS LAST, event_time ASC',
        options => map(
            'target-file-size-bytes', '536870912',
            'min-input-files', '10',
            'partial-progress.enabled', 'true'   -- коммитить пачками, а не всё разом
        )
    )
""")

# 2. Слить манифесты: планирование запроса читает их целиком.
spark.sql(f"CALL lake.system.rewrite_manifests(table => '{TABLE}')")

# 3. Удалить снапшоты старше окна time travel (например, 7 суток).
spark.sql(f"""
    CALL lake.system.expire_snapshots(
        table => '{TABLE}',
        older_than => TIMESTAMP '2026-07-11 00:00:00',
        retain_last => 10
    )
""")

# 4. Удалить осиротевшие файлы. ВАЖНО: older_than должен быть заведомо больше
#    длительности самой долгой пишущей джобы, иначе вы удалите файлы,
#    которые прямо сейчас пишет другой процесс.
spark.sql(f"""
    CALL lake.system.remove_orphan_files(
        table => '{TABLE}',
        older_than => TIMESTAMP '2026-07-15 00:00:00'
    )
""")

9. Экономика: за что вы на самом деле платите

Разложим счёт за аналитическое хранилище на составляющие (порядки величин для объектного хранилища класса S3 Standard; уточняйте по своему прайсу):

  • Хранение — порядка 0.02 $ за ГБ в месяц. Терабайт ≈ 20 $–25/мес. Дёшево — и именно поэтому все забывают про старые снапшоты и осиротевшие файлы, которые тихо удваивают объём.
  • Запросы — GET порядка 0.4 $ за миллион, PUT/LIST порядка 5 $ за миллион. Заметьте разницу в порядок: писать и листить в 10 раз дороже, чем читать. Таблица из 2 млн мелких файлов стоит ~1.6 $ только на GET-ах при одном полном скане, плюс LIST-ы на планирование, плюс задержка.
  • Трафик наружу — самая дорогая строка (0.05 $–0.09 за ГБ). Держите вычисление в том же регионе, что и данные.
  • Вычисление — обычно доминирует, но напрямую зависит от того, сколько байт пришлось прочитать.

Отсюда следует контринтуитивный вывод: проблема мелких файлов — это в первую очередь проблема денег и латентности, а не «неаккуратности». Сравните два способа хранить 100 ГБ:

200 000 файлов по 0.5 МБ 200 файлов по 512 МБ
GET-запросов на полный скан ~400 000 ~600
Время планирования десятки секунд миллисекунды
Эффективность сжатия низкая (словарь на файл крошечный) высокая
Полезность min/max почти нулевая высокая
Стоимость скана в 50–100 раз выше базовая

Правило большого пальца: целевой размер файла 128 МБ – 1 ГБ, партиция — не меньше сотни мегабайт. Если партиция за час весит 3 МБ — партиционируйте по дню, а не по часу.


10. Как выбрать раскладку: практический алгоритм

И отдельно про выбор ключа партиционирования — он остаётся самым частым источником боли (подробно о двух смыслах слова «партиция» — в https://courses.digitable.life/post/data-engineering/03-batch-processing/):

  • партиционируйте по колонке, которая есть в фильтре почти каждого запроса (обычно дата события);
  • целевой размер партиции — сотни мегабайт – единицы гигабайт;
  • не больше двух уровней вложенности; dt/country/city/device — это гарантированные мелкие файлы;
  • никогда не партиционируйте по колонке высокой кардинальности (user_id, session_id) — вместо этого используйте bucket(N, user_id) или сортировку;
  • если запросы фильтруют по двум измерениям одновременно — вместо второго уровня партиций возьмите сортировку/Z-order (кривая Гильберта или Мортона даёт кластеризацию сразу по нескольким колонкам).

11. Типичные ошибки

  1. CSV/JSON как основной формат озера. Нет схемы, нет статистик, нет pushdown; парсинг съедает больше CPU, чем вся полезная работа. Landing-зона — Avro или Parquet; JSON — только на входе от внешнего источника.
  2. Мелкие файлы. Стриминг с чекпоинтом раз в минуту даёт 1440 файлов в сутки на партицию. Компакция — обязательный фоновый процесс, а не «когда-нибудь потом».
  3. Данные записаны в случайном порядке. min/max в каждом файле покрывают весь диапазон, отсечение не работает вовсе. Лечится ORDER BY перед записью или sort-компакцией.
  4. Партиционирование по высококардинальной колонке. 3 млн директорий, каждая по 20 КБ; планирование запроса длится дольше, чем сам запрос.
  5. SELECT * в аналитическом запросе. Убивает главное преимущество колоночного хранения. В витринах и в BI это ещё и приводит к неявной зависимости от всех колонок.
  6. Забытые expire_snapshots / VACUUM. Объём хранения растёт вдвое-втрое, никто не понимает, почему счёт растёт при постоянном размере таблицы.
  7. Слишком короткий older_than в remove_orphan_files или VACUUM. Прямой путь к удалению файлов у себя из-под ног работающей джобы и к битой таблице. Ставьте окно заведомо больше самой долгой записи (обычно 3–7 суток).
  8. Изменение типа колонки «переписыванием файлов». Используйте эволюцию схемы табличного формата; безопасные расширения (intlong, floatdouble, добавление поля) поддерживаются штатно.
  9. Один гигантский row group на файл. Теряется параллелизм и отсечение внутри файла; 128–512 МБ — разумный целевой размер.
  10. Тысячи мелких INSERT в ClickHouse. Каждый порождает парт; сервер уходит в бесконечные мержи. Батчируйте на стороне приложения или используйте буферизацию/асинхронные вставки.

12. Как это выглядит в проде

Типичная зрелая платформа выглядит так — и почти никогда не состоит из одного хранилища:

  • Landing / raw — Avro или Parquet ровно в том виде, в котором пришли данные, партиции по дате приёма, retention 30–90 дней. Иммутабельный слой: если в логике нашли баг, backfill делается отсюда.
  • Core / staging — Iceberg или Delta, Parquet + zstd, партиции по дате события, компакция и сортировка ночью. Именно здесь живут ACID, апдейты и GDPR-удаления.
  • Marts / витрины — либо Iceberg-таблицы, оптимизированные под конкретные запросы, либо загрузка в ClickHouse/Druid, когда нужна субсекундная латентность на дашборде.
  • Каталог — REST-каталог (Polaris/Nessie/Glue) как единая точка правды о таблицах; движки (Spark, Trino, Flink, ClickHouse, DuckDB) читают одни и те же файлы, не копируя их.
  • Обслуживание — отдельный DAG: компакция → rewrite_manifests → expire_snapshots → remove_orphan_files, с алертами на «партиций с >1000 файлов» и «объём хранения вырос >20% за неделю».

Три метрики, которые стоит вывести на дашборд платформы: среднее число файлов на партицию, медианный размер файла и объём хранения на активные данные против общего объёма (отношение показывает, сколько вы платите за мусор). Все три считаются прямо из метаданных таблиц.

Отдельно про совместимость: главный практический смысл открытых форматов — данные переживают движки. Кластер Spark можно заменить на Trino, Trino на DuckDB, а таблицы останутся на месте. Это единственная известная защита от вендорского замка в аналитике, и она работает ровно до тех пор, пока вы не начали использовать проприетарные расширения формата.


13. Мини-итог

  • Раскладка байт — не деталь реализации, а главный рычаг производительности и стоимости. Колоночное хранение выигрывает трижды: projection pushdown, сжатие, векторизация.
  • Кодирование (dictionary, RLE, delta, bit-packing) важнее компрессора, и его эффективность зависит от порядка строк. Сортировка перед записью — самая дешёвая оптимизация в озере.
  • Parquet: файл → row group (128–512 МБ) → column chunk → page (~1 МБ), footer со статистиками в конце. Отсечение идёт воронкой: партиции → файлы → row groups → страницы → колонки.
  • Arrow — про память, Parquet — про диск, Avro — про построчную запись. Это не конкуренты.
  • Колоночные СУБД (ClickHouse, DuckDB) добавляют к тем же принципам управляемую раскладку, разрежённый индекс и фоновые мержи; ключ сортировки там важнее всех остальных настроек.
  • «Просто файлы в S3» ломаются на атомарности, листинге, зашитых в путь партициях и апдейтах. Табличный формат (Iceberg/Delta/Hudi) чинит всё это одним приёмом: атомарной подменой указателя на метаданные.
  • Lakehouse требует обслуживания: компакция, expire_snapshots, remove_orphan_files. Без них деградируют и запросы, и бюджет.
  • Мелкие файлы — самая частая и самая дорогая болезнь озера. Цель: 128 МБ – 1 ГБ на файл.

Источники

Что дальше

Мы научились хранить данные так, чтобы их было дёшево и быстро читать. Остался вопрос, который никакой формат не решает: можно ли этим данным верить. Схема в Parquet гарантирует типы, но не гарантирует, что amount не отрицательный, что источник не начал присылать в поле country коды в другом регистре, и что витрину не сломает изменение в чужом сервисе. Для этого нужны контракты данных, проверки качества, линидж и правила владения.

Дальше: Качество данных, контракты, линидж и governance.

Нашли неточность? Выделите фрагмент текста — рядом появится жучок.

Нужен разбор именно вашей ситуации?

Статья описывает общий случай. Если у вас частный — можно разобрать его отдельно, платно. А если не хватает целого материала, предложите тему: её оплачивают вскладчину, и она выходит открытой для всех.

Доска запросов