Spark + Iceberg на собеседовании Data Engineer
created_at с временем. Какой фильтр надёжнее задаёт диапазон «весь месяц» без потерь по краям?Содержание:
Почему Iceberg спрашивают
Apache Iceberg — это табличный формат для озера данных, который приносит в S3/HDFS то, чего там исторически не было: ACID-транзакции, эволюцию схемы без перезаписи данных и путешествия во времени. Если раньше «таблица» в озере была просто папкой с parquet-файлами, где два параллельных писателя легко ломали друг другу данные, то Iceberg добавляет поверх файлов слой метаданных, который делает из этой папки нормальную транзакционную таблицу.
На собесе Data Engineer связку Spark + Iceberg спрашивают, потому что она стала стандартом lakehouse-архитектур: Spark считает, Iceberg хранит и версионирует. Интервьюер обычно проверяет четыре вещи: как подключить каталог, как читать/писать транзакционно, как эволюционировать схему без боли и как бороться с деградацией из-за мелких файлов. Разберём по порядку.
Настройка
Iceberg подключается к Spark как отдельный каталог. Указываем реализацию каталога, его тип (здесь — Hive Metastore), адрес метастора и путь к warehouse в объектном хранилище:
spark.conf.set("spark.sql.catalog.iceberg", "org.apache.iceberg.spark.SparkCatalog")
spark.conf.set("spark.sql.catalog.iceberg.type", "hive")
spark.conf.set("spark.sql.catalog.iceberg.uri", "thrift://hive-metastore:9083")
spark.conf.set("spark.sql.catalog.iceberg.warehouse", "s3://my-warehouse/")Плюс к этому нужно, чтобы jar-файл Iceberg лежал в classpath Spark (обычно подключается через --packages при старте). Тип каталога может быть разным — hive, hadoop, rest, glue, nessie; выбор зависит от того, где вы храните метаданные о таблицах. На собесе достаточно понимать, что каталог — это «телефонная книга» таблиц, а не сами данные.
Каталог
После настройки с таблицами Iceberg работают обычным SQL. Создание таблицы с партиционированием по дню от timestamp:
USE iceberg.warehouse;
CREATE TABLE events (
id BIGINT, ts TIMESTAMP, user_id BIGINT
) USING iceberg
PARTITIONED BY (days(ts));Ключевая фича здесь — скрытое партиционирование (hidden partitioning). Функция days(ts) означает, что Iceberg сам вычисляет партицию из колонки ts, и в запросах не нужно писать отдельную колонку-партицию вроде dt. Пользователь фильтрует по ts, а Iceberg под капотом отсекает лишние партиции. Это избавляет от классической ошибки Hive, где забыли добавить фильтр по партиционной колонке и прочитали всю таблицу.
Чтение и запись
Читать и писать можно через DataFrame API или SQL:
df = spark.read.format("iceberg").load("iceberg.warehouse.events")
# Или через SQL
df = spark.sql("SELECT * FROM iceberg.warehouse.events WHERE date(ts) = '2026-05-07'")
# Запись (добавление)
new_df.writeTo("iceberg.warehouse.events").append()Главное отличие от «просто parquet в папке» — ACID-гарантии. Каждая запись создаёт новый снапшот таблицы, а читатели всегда видят согласованную версию: либо коммит целиком, либо ничего. Iceberg использует оптимистичную конкурентность — несколько писателей могут работать параллельно, и если их изменения не пересекаются, оба коммита проходят; при конфликте один из коммитов повторяется на свежих метаданных или падает, а не молча затирает чужие данные. Именно поэтому Iceberg безопасен там, где голый parquet ломался.
Эволюция схемы
Схему таблицы можно менять на лету, не переписывая данные:
ALTER TABLE iceberg.warehouse.events ADD COLUMN device_type STRING;
ALTER TABLE iceberg.warehouse.events RENAME COLUMN ts TO event_ts;Это работает без перезаписи файлов, потому что Iceberg отслеживает колонки по числовым ID, а не по имени и позиции. Добавление, удаление, переименование и перестановка колонок меняют только метаданные — старые файлы читаются как есть, а новые пишутся уже по новой схеме. В Hive такое переименование обычно означало «переписать всю таблицу»; в Iceberg это мгновенная операция над метаданными. На собесе это любимый вопрос: почему Iceberg умеет RENAME COLUMN бесплатно, а Hive — нет.
Компакция и снапшоты
Стриминг и частые мелкие записи порождают «проблему маленьких файлов»: тысячи крохотных parquet сильно замедляют чтение, потому что на каждый файл — отдельное открытие и метаданные. Iceberg умеет схлопывать их процедурой компакции:
CALL iceberg.system.rewrite_data_files('warehouse.events');Процедура перепаковывает мелкие файлы в крупные (bin-packing) — это регулярная операция обслуживания, которую обычно ставят по расписанию. Отдельная история — снапшоты: каждая запись оставляет версию таблицы для time travel, и со временем их накапливается много. Их чистят, чтобы освободить место:
CALL iceberg.system.expire_snapshots(
'warehouse.events',
TIMESTAMP '2026-04-01'
);expire_snapshots удаляет старые снапшоты и файлы, на которые больше никто не ссылается, и высвобождает место в хранилище. Плата за это — теряется возможность откатиться к состоянию до указанной даты. Поэтому на собесе важно проговорить trade-off: держать историю дольше — дороже по хранению, но безопаснее для восстановления и аудита.
Как это спрашивают на собесе
Реальные формулировки и что за ними проверяют:
- «Зачем нужен Iceberg, если есть parquet в S3?» Ответ: ACID, эволюция схемы без перезаписи, time travel, скрытое партиционирование — всё, чего нет у голых файлов.
- «Что происходит при двух параллельных записях?» Оптимистичная конкурентность: непересекающиеся коммиты проходят оба, конфликтующий повторяется или падает, данные не затираются.
- «Как переименовать колонку в таблице на терабайт?»
ALTER TABLE ... RENAME COLUMN— мгновенно, потому что колонки трекаются по ID, данные не переписываются. - «У нас деградировали чтения из стриминга — почему?» Проблема мелких файлов; лечится
rewrite_data_filesпо расписанию. - «Как откатить таблицу к вчерашнему состоянию?» Через снапшоты и time travel (
VERSION AS OF/TIMESTAMP AS OF), если снапшот ещё не удалёнexpire_snapshots.
Частые ошибки
- Не настроена компакция. Стриминг пишет мелкие файлы, чтения деградируют, а никто не запускает
rewrite_data_files. Обслуживание таблицы — часть архитектуры, а не опция. - Слишком агрессивный expire_snapshots. Удалили снапшоты за вчера — и потеряли возможность откатиться и провести аудит. Ретеншн снапшотов надо согласовывать с требованиями восстановления.
- Партиционирование по сырой колонке высокой кардинальности. Партиция на каждый
user_idпорождает миллионы мелких партиций. Партиционируют по трансформациям —days(ts),bucket(N, id). - Путать каталог с данными. Каталог хранит указатели на таблицы, а не сами данные. Смена типа каталога (
hive/rest/glue) не трогает файлы в warehouse. - Считать, что Iceberg сам всё оптимизирует. Компакция, чистка снапшотов и orphan-файлов — это явные процедуры, которые надо ставить в расписание.
Связанные темы
- Iceberg deep для DE
- Spark RDD vs DataFrame для DE
- Lakehouse Iceberg Delta для DE
- Spark Catalog для DE
- Подготовка к собесу Data Engineer
FAQ
Чем Iceberg отличается от обычного parquet в S3?
Iceberg добавляет поверх файлов слой метаданных, который даёт ACID-транзакции, атомарные коммиты, эволюцию схемы без перезаписи данных, скрытое партиционирование и путешествия во времени. Голый parquet — это просто папка с файлами: параллельные записи небезопасны, а изменение схемы часто требует переписать таблицу.
Как Iceberg обеспечивает ACID при параллельной записи?
Через оптимистичную конкурентность и снапшоты. Каждая запись создаёт новый снапшот; читатели видят согласованную версию целиком. Если два писателя не пересекаются по данным — проходят оба коммита. При конфликте один коммит повторяется на свежих метаданных или завершается ошибкой, но чужие данные не затираются.
Почему RENAME COLUMN в Iceberg бесплатный, а в Hive нет?
Iceberg идентифицирует колонки по числовым ID, а не по имени и позиции. Переименование, добавление, удаление и перестановка колонок меняют только метаданные, а файлы данных остаются как есть. Hive же сопоставляет колонки по позиции, поэтому изменение схемы обычно означает перезапись данных.
Что делать с проблемой мелких файлов?
Регулярно запускать компакцию CALL system.rewrite_data_files(...) — она перепаковывает множество мелких файлов в крупные и возвращает чтениям скорость. Обычно ставят по расписанию, особенно если в таблицу пишет стриминг.
Это официальная информация?
Нет. Статья основана на документации Apache Iceberg и Spark и опыте кандидатов. Конкретный стек, тип каталога и настройки зависят от команды и инфраструктуры.
Тренируйте Data Engineering — откройте тренажёр с 1500+ вопросами для собесов.