dbt + Iceberg: file_format=iceberg, tblproperties и schema evolution
Apache Iceberg даёт dbt-моделям ACID-транзакции, скрытое партиционирование и безопасную эволюцию схем без перезаписи файлов. Разбираем интеграцию, tblproperties и time travel.
Почему dbt + Iceberg - сильный тандем¶
Проблемы классического Parquet-озера¶
До появления современных table formats хранение данных в Parquet на S3 или HDFS работало по принципу "файловой системы": таблица - это просто директория с .parquet файлами. Этот подход прост, но имеет фундаментальные ограничения.
Нет атомарности записи. Когда вы выполняете INSERT OVERWRITE на партицию, Spark сначала удаляет старые файлы, потом пишет новые. В промежутке между удалением и записью таблица находится в inconsistent-состоянии. Если другой процесс в этот момент читает таблицу - он получит пустые или частичные данные.
Нет эволюции схемы. Если вы добавляете колонку в Parquet-таблицу, нужно физически переписать все существующие файлы с новой схемой. Для терабайтного датасета - это часы работы кластера. Переименование колонки - аналогичная катастрофа.
Нет time travel. Нет возможности прочитать "вчерашнюю версию" данных, не делая отдельный бэкап.
Проблема мелких файлов при инкрементальной загрузке. Каждый INSERT создаёт новые файлы. После сотен инкрементальных загрузок директория содержит тысячи мелких файлов, и чтение становится медленным.
Проблема со schema mismatch. Если одна команда записала Parquet с 10 колонками, другая - с 11, читатель получит ошибку несовместимости схем или непредсказуемые null.
Это не теоретические проблемы - они преследовали дата-инженеров в production ежедневно. Именно для их решения появились table formats нового поколения: Apache Iceberg, Delta Lake, Apache Hudi.
Apache Iceberg: что это такое¶
Apache Iceberg - это спецификация открытого формата таблиц для больших аналитических датасетов. Он не хранит данные сам - данные по-прежнему лежат в Parquet, ORC или Avro файлах на S3/MinIO/HDFS. Iceberg добавляет поверх них метаданный слой: набор JSON-файлов, которые описывают, из каких именно файлов состоит таблица в каждый момент времени.
Ключевое отличие от Hive-подхода: в Hive метаданные таблицы хранятся в Hive Metastore (реляционная БД), а "что входит в партицию" определяется по имени директории на диске (event_date=2024-01-15/). В Iceberg метаданные хранятся рядом с данными в виде файлов, и каждое изменение таблицы - это создание нового snapshot (неизменяемого снимка состояния таблицы).
Роль dbt в экосистеме Iceberg¶
dbt управляет логикой трансформаций: что из чего получается, в каком порядке, с какой бизнес-логикой. Iceberg управляет физическим хранением: как атомарно записать результат, как эволюционировать схему, как сохранить историю версий.
Это идеальное разделение ответственности:
- Аналитик пишет SQL в dbt-модели → dbt компилирует SQL → Spark выполняет SQL → Iceberg фиксирует результат как новый snapshot
- Схема изменилась → dbt пересоздаёт модель с новой схемой → Iceberg применяет только метаданные, не перезаписывая терабайты файлов
- Нужно откатиться к вчерашним данным → Iceberg time travel через
VERSION AS OFилиTIMESTAMP AS OF
Анатомия Iceberg: как это работает внутри¶
Структура метаданных¶
Диаграмма показывает многоуровневую структуру Iceberg-таблицы. Разберём каждый уровень:
Catalog - точка входа. Хранит только один указатель: путь к текущему metadata-файлу таблицы. Это может быть Hive Metastore, REST Catalog, Nessie или Hadoop-каталог (файлы на S3).
Metadata File - JSON-файл, описывающий всё о таблице: схему, историю всех snapshot'ов, информацию о партиционировании. При каждом изменении создаётся новый metadata-файл, и каталог обновляет указатель атомарно.
Snapshot - неизменяемый снимок состояния таблицы в конкретный момент. Каждый INSERT, UPDATE, MERGE, DELETE создаёт новый snapshot. Snapshot содержит список всех manifest-файлов, которые описывают данные в этой версии таблицы.
Manifest List - файл, содержащий список всех manifest'ов для конкретного snapshot'а. Хранит статистику по каждому manifest'у для эффективного планирования запросов.
Manifest File - список конкретных data-файлов с их статистикой (min/max значения колонок, количество строк, размер). Именно благодаря этой статистике Iceberg может пропускать целые файлы при запросах с фильтрацией - без полного сканирования.
Data Files - реальные файлы с данными (Parquet, ORC, Avro). Они неизменяемы - Iceberg никогда не модифицирует существующие файлы, только добавляет новые.
Snapshot как атомарная транзакция¶
Иммутабельность и snapshot-подход дают Iceberg ACID-гарантии. Когда Spark пишет новую партицию:
- Все data-файлы записываются в объектное хранилище (S3 PUT)
- Создаётся manifest-файл, описывающий эти файлы
- Создаётся новый snapshot, ссылающийся на новый manifest
- Создаётся новый metadata-файл с обновлённым current-snapshot
- Атомарная операция: каталог переключает указатель на новый metadata-файл
Если что-то упало на шаге 1-4, таблица остаётся в предыдущем состоянии - шаг 5 не был выполнен. Это и есть атомарность: либо весь коммит применяется, либо не применяется ничего.
Настройка Spark для Iceberg¶
Spark configuration для Iceberg¶
Перед использованием Iceberg в Spark нужно зарегистрировать Iceberg каталог через конфигурацию SparkSession:
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.appName("dbt + Iceberg") \
.config("spark.sql.extensions",
"org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \
.config("spark.sql.catalog.iceberg",
"org.apache.iceberg.spark.SparkCatalog") \
.config("spark.sql.catalog.iceberg.type", "hive") \
.config("spark.sql.catalog.iceberg.uri",
"thrift://hive-metastore:9083") \
.config("spark.sql.catalog.iceberg.warehouse",
"s3a://datalake/iceberg-warehouse/") \
.config("spark.jars.packages",
"org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.4.3,"
"org.apache.hadoop:hadoop-aws:3.3.4") \
.getOrCreate()
spark.sql.catalog.iceberg - регистрация именованного каталога с именем iceberg. После этого таблицы в этом каталоге адресуются как iceberg.database.table.
spark.sql.catalog.iceberg.type - тип бэкенда каталога:
| Тип | Хранение метаданных | Когда использовать |
|---|---|---|
hive |
Hive Metastore (Thrift) | Shared cluster с существующим HMS |
hadoop |
Файлы в HDFS/S3 | Простые сценарии, нет HMS |
rest |
REST API (Apache Polaris, Tabular) | Cloud-native, multi-engine |
nessie |
Nessie server (git-like) | Версионирование данных как git |
glue |
AWS Glue Catalog | Managed AWS окружение |
Hadoop Catalog: Iceberg без HMS¶
Для локальной разработки и простых сценариев Hadoop Catalog проще - метаданные хранятся в файлах рядом с данными:
spark = SparkSession.builder \
.config("spark.sql.extensions",
"org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \
.config("spark.sql.catalog.local",
"org.apache.iceberg.spark.SparkCatalog") \
.config("spark.sql.catalog.local.type", "hadoop") \
.config("spark.sql.catalog.local.warehouse",
"s3a://datalake/iceberg/") \
.config("spark.hadoop.fs.s3a.endpoint", "http://localhost:9000") \
.config("spark.hadoop.fs.s3a.access.key", "minioadmin") \
.config("spark.hadoop.fs.s3a.secret.key", "minioadmin123") \
.config("spark.hadoop.fs.s3a.path.style.access", "true") \
.getOrCreate()
После этой конфигурации CREATE TABLE local.analytics.events USING iceberg создаст файлы по пути s3a://datalake/iceberg/analytics/events/.
profiles.yml для dbt + Iceberg¶
# ~/.dbt/profiles.yml
analytics_iceberg:
target: dev
outputs:
dev:
type: spark
method: session
schema: dbt_dev
threads: 2
server_side_parameters:
spark.master: "local[4]"
spark.driver.memory: "4g"
spark.sql.shuffle.partitions: "4"
spark.sql.adaptive.enabled: "true"
spark.jars.packages: >
org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.4.3,
org.apache.hadoop:hadoop-aws:3.3.4,
com.amazonaws:aws-java-sdk-bundle:1.12.262
spark.sql.extensions: >
org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions
spark.sql.catalog.iceberg: >
org.apache.iceberg.spark.SparkCatalog
spark.sql.catalog.iceberg.type: hadoop
spark.sql.catalog.iceberg.warehouse: "s3a://datalake/iceberg/"
spark.hadoop.fs.s3a.endpoint: "http://localhost:9000"
spark.hadoop.fs.s3a.access.key: "{{ env_var('MINIO_ACCESS_KEY', 'minioadmin') }}"
spark.hadoop.fs.s3a.secret.key: "{{ env_var('MINIO_SECRET_KEY', 'minioadmin123') }}"
spark.hadoop.fs.s3a.path.style.access: "true"
spark.hadoop.fs.s3a.impl: >
org.apache.hadoop.fs.s3a.S3AFileSystem
prod:
type: spark
method: thrift
host: "{{ env_var('SPARK_THRIFT_HOST') }}"
port: 10000
user: "{{ env_var('SPARK_USER') }}"
password: "{{ env_var('SPARK_PASSWORD') }}"
schema: analytics
threads: 8
server_side_parameters:
spark.sql.shuffle.partitions: "800"
spark.sql.adaptive.enabled: "true"
spark.sql.extensions: >
org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions
spark.sql.catalog.iceberg: >
org.apache.iceberg.spark.SparkCatalog
spark.sql.catalog.iceberg.type: hive
spark.sql.catalog.iceberg.uri: "thrift://hive-metastore:9083"
spark.sql.catalog.iceberg.warehouse: "s3a://datalake/iceberg/"
Iceberg-модели в dbt¶
Базовый синтаксис: table materialization¶
-- models/marts/fct_daily_events.sql
{{ config(
materialized='table',
file_format='iceberg',
location_root='s3a://datalake/iceberg/gold/'
) }}
SELECT
CAST(event_ts AS DATE) AS event_date,
channel,
campaign_id,
COUNT(*) AS event_count,
COUNT(DISTINCT user_id) AS unique_users,
SUM(revenue) AS total_revenue,
CURRENT_TIMESTAMP() AS updated_at
FROM {{ ref('stg_events') }}
GROUP BY CAST(event_ts AS DATE), channel, campaign_id
file_format='iceberg' - ключевой параметр, указывающий dbt использовать Iceberg вместо обычного Parquet. dbt сгенерирует DDL с USING iceberg.
location_root - корневой путь, под которым Iceberg создаст папку с именем таблицы. Итоговый путь: s3a://datalake/iceberg/gold/fct_daily_events/.
Что генерирует dbt под капотом:
CREATE TABLE IF NOT EXISTS dbt_dev.fct_daily_events
USING iceberg
LOCATION 's3a://datalake/iceberg/gold/fct_daily_events'
AS
SELECT
CAST(event_ts AS DATE) AS event_date,
channel,
campaign_id,
COUNT(*) AS event_count,
...
FROM dbt_dev.stg_events
GROUP BY ...
Скрытое партиционирование: Hidden Partitioning¶
В классическом Hive партиционирование требует явных колонок-партиций, которые существуют и в схеме таблицы, и в именах директорий: event_date=2024-01-15/part-00001.parquet. Это создаёт проблему: аналитик должен знать о партиции и явно указывать её в запросе для partition pruning.
Hidden Partitioning в Iceberg - принципиально другой подход. Партиция - это трансформация значения существующей колонки, а не отдельная колонка. Партиция не видна в схеме таблицы (отсюда "hidden"), но активно используется Spark для пропуска ненужных файлов.
Доступные трансформации:
| Трансформация | SQL | Что делает |
|---|---|---|
years(col) |
years(event_ts) |
Партиция по году |
months(col) |
months(event_ts) |
Партиция по году-месяцу |
days(col) |
days(event_ts) |
Партиция по дате |
hours(col) |
hours(event_ts) |
Партиция по году-месяц-дата-час |
bucket(N, col) |
bucket(16, user_id) |
N бакетов по хешу |
truncate(W, col) |
truncate(5, category) |
Усечение строки до W символов |
identity(col) |
identity(status) |
Точное значение (как Hive) |
-- models/marts/fct_events_partitioned.sql
{{ config(
materialized='incremental',
incremental_strategy='append',
file_format='iceberg',
location_root='s3a://datalake/iceberg/gold/',
partition_by={
'field': 'event_ts',
'data_type': 'timestamp',
'granularity': 'day'
}
) }}
SELECT
event_id,
user_id,
session_id,
event_type,
event_ts,
channel,
revenue,
CURRENT_TIMESTAMP() AS processed_at
FROM {{ ref('stg_events') }}
{% if is_incremental() %}
WHERE event_ts >= (SELECT MAX(event_ts) FROM {{ this }})
{% endif %}
Конфигурация partition_by в dbt-spark для Iceberg транслируется в DDL:
CREATE TABLE dbt_dev.fct_events_partitioned
USING iceberg
PARTITIONED BY (days(event_ts))
LOCATION 's3a://datalake/iceberg/gold/fct_events_partitioned'
Ключевое преимущество hidden partitioning: запрос вида WHERE event_ts >= '2024-01-15' автоматически получит partition pruning - Spark прочитает metadata manifest'ов и пропустит все файлы, не попадающие в диапазон дат. И это работает без изменения SQL: аналитик пишет WHERE event_ts >= '...', а не WHERE event_date = '...'.
Множественное партиционирование¶
{{ config(
materialized='table',
file_format='iceberg',
partition_by=[
{
'field': 'event_ts',
'data_type': 'timestamp',
'granularity': 'month'
},
{
'field': 'channel',
'data_type': 'string'
}
]
) }}
Это создаст партиционирование PARTITIONED BY (months(event_ts), identity(channel)). Spark будет эффективно пропускать файлы при фильтрации как по дате, так и по каналу.
TBLPROPERTIES: тонкая настройка таблицы¶
tblproperties в dbt - это словарь свойств Iceberg-таблицы, которые dbt передаёт через TBLPROPERTIES в DDL. Эти свойства определяют физическое поведение таблицы: формат файлов, размер, компрессию, политику хранения истории.
{{ config(
materialized='incremental',
incremental_strategy='merge',
unique_key=['user_id', 'event_date'],
file_format='iceberg',
location_root='s3a://datalake/iceberg/gold/',
partition_by={
'field': 'event_date',
'data_type': 'date',
'granularity': 'month'
},
tblproperties={
'write.format.default': 'parquet',
'write.parquet.compression-codec': 'zstd',
'write.target-file-size-bytes': '134217728',
'write.distribution-mode': 'hash',
'write.object-storage.enabled': 'true',
'write.object-storage.path': 's3a://datalake/iceberg/data/',
'history.expire.max-snapshot-age-ms': '604800000',
'history.expire.min-snapshots-to-keep': '3',
'commit.manifest.min-count-to-merge': '100',
'commit.manifest-merge.enabled': 'true',
'read.split.target-size': '134217728',
'read.split.planning-lookback': '10',
}
) }}
Разберём каждое свойство подробно:
Свойства записи (write.*)¶
write.format.default - формат данных внутри Iceberg-таблицы:
parquet- лучший выбор для аналитических нагрузок (колоночное хранение, эффективная компрессия)orc- хорошая альтернатива, особенно при интеграции с Hiveavro- для streaming-записи (строковое хранение, нет колоночного)
write.parquet.compression-codec - алгоритм сжатия:
snappy- быстрое сжатие/разжатие, умеренная степень сжатия. Стандарт для production аналитикиzstd- лучшее соотношение степени сжатия и скорости, рекомендуется для cold storagegzip- максимальное сжатие, но медленнее. Для архивных данныхlz4- быстрейшее разжатие, низкая степень сжатия. Для горячих данных с частым чтением
write.target-file-size-bytes - целевой размер файлов при записи. 134217728 = 128 МБ. Это баланс между количеством файлов и размером. Меньше 64 МБ - small files problem. Больше 512 МБ - медленная параллелизация чтения.
write.distribution-mode - как Iceberg распределяет строки между writer'ами:
none- каждый writer пишет в свои файлы без перераспределенияhash- строки с одинаковым partition hash пишутся одним writer'ом (меньше файлов на партицию)range- сортировка перед записью (лучшее сжатие в Parquet, но дороже)
write.object-storage.enabled и write.object-storage.path - важная настройка для S3/MinIO. По умолчанию Iceberg использует путь вида s3a://bucket/warehouse/db/table/data/part-001.parquet. Это детерминированный путь, и при интенсивной записи все запросы попадают в один S3 "prefix", что ограничивает IOPS (S3 throttles по prefix).
С write.object-storage.enabled=true Iceberg добавляет в путь хеш от содержимого файла: s3a://bucket/data/a1b2c3/part-001.parquet. Это случайные префиксы - запросы распределяются по сотням S3 серверов, устраняя bottleneck.
Свойства хранения истории (history.*)¶
history.expire.max-snapshot-age-ms - максимальный возраст snapshot'а в миллисекундах перед его удалением. 604800000 = 7 дней. После этого старые snapshot'ы станут кандидатами на удаление при вызове expire_snapshots.
Важно понимать: удаление snapshot'а не удаляет немедленно файлы данных - это только помечает snapshot как "expired". Физическое удаление файлов происходит при явном вызове CALL iceberg.system.remove_orphan_files(...).
history.expire.min-snapshots-to-keep - минимальное количество snapshot'ов, которые всегда сохраняются, даже если они старше max-snapshot-age-ms. Значение 3 означает: всегда держать как минимум последние 3 snapshot'а для возможности откатиться.
Свойства manifest'ов (commit.manifest.*)¶
commit.manifest.min-count-to-merge и commit.manifest-merge.enabled - управление автоматическим объединением manifest-файлов.
Каждый commit создаёт новый manifest-файл. После тысяч инкрементальных загрузок накапливаются тысячи manifest'ов, и планировщик запросов вынужден читать их все для построения плана. commit.manifest-merge.enabled=true включает автоматическое слияние мелких manifest'ов при каждом commit'е. min-count-to-merge=100 - слияние начнётся только когда накопится более 100 manifest'ов.
Свойства чтения (read.*)¶
read.split.target-size - целевой размер одного split при чтении (аналог maxPartitionBytes в Spark). 134217728 = 128 МБ.
read.split.planning-lookback - сколько manifest'ов анализировать для планирования splits. Увеличение ускоряет планирование, но может быть менее оптимальным при неравномерных файлах.
Schema Evolution: безопасное изменение схемы¶
Почему это работает: Column ID Tracking¶
Корень schema evolution в Iceberg - идентификаторы колонок (Column IDs). В обычном Parquet колонки идентифицируются по имени. Если вы переименуете name в full_name, старые файлы со схемой name стали несовместимы с новой схемой full_name.
В Iceberg каждая колонка имеет числовой ID (назначается при создании и никогда не меняется):
{
"schema": {
"type": "struct",
"fields": [
{"id": 1, "name": "user_id", "type": "long", "required": true},
{"id": 2, "name": "name", "type": "string", "required": false},
{"id": 3, "name": "email", "type": "string", "required": false},
{"id": 4, "name": "created_at","type": "timestamp","required": false}
]
}
}
Когда вы переименуете name → full_name, metadata-файл обновляется:
{"id": 2, "name": "full_name", "type": "string", "required": false}
Старые Parquet-файлы остаются нетронутыми - они по-прежнему хранят данные в поле с ID=2. При чтении Spark использует ID для сопоставления полей, игнорируя изменение имени. Данные читаются корректно, файлы не перезаписываются.
Добавление колонки (ADD COLUMN)¶
-- dbt model изменился: добавлена колонка loyalty_score
{{ config(
materialized='table',
file_format='iceberg',
location_root='s3a://datalake/iceberg/gold/'
) }}
SELECT
user_id,
full_name,
email,
created_at,
loyalty_score,
country_code
FROM {{ ref('stg_users') }}
Когда dbt обнаруживает новую колонку loyalty_score (которой не было в предыдущей версии таблицы), для Iceberg-таблицы он выполняет:
ALTER TABLE dbt_dev.dim_users ADD COLUMN loyalty_score INT;
Это операция только с метаданными - добавляется новая запись в JSON-схему. Существующие Parquet-файлы не изменяются. При чтении старых файлов, где loyalty_score отсутствует, Spark вернёт null для этой колонки.
Переименование колонки (RENAME COLUMN)¶
SELECT
user_id,
full_name,
email,
created_at
FROM {{ ref('stg_users') }}
Было name, стало full_name. dbt выполняет:
ALTER TABLE dbt_dev.dim_users RENAME COLUMN name TO full_name;
Снова - только метаданные. Старые файлы хранят данные с Column ID=2 (который раньше назывался name, теперь full_name). Читаются корректно.
Внимание: если другие dbt-модели или BI-инструменты обращаются к этой таблице по имени колонки name - они сломаются. Schema evolution в Iceberg физически безопасна, но не отменяет необходимость координации с потребителями данных.
Расширение типа (Type Promotion)¶
Iceberg поддерживает безопасные расширения типов:
| Исходный тип | Целевой тип | Безопасно? |
|---|---|---|
INT |
LONG |
Да - диапазон расширяется |
FLOAT |
DOUBLE |
Да - точность увеличивается |
DECIMAL(9,2) |
DECIMAL(18,4) |
Да - масштаб может расширяться |
STRING |
любой | Нет - нет implicit cast |
LONG |
INT |
Нет - потеря данных |
DATE |
TIMESTAMP |
Да (с caveat'ами) |
ALTER TABLE dbt_dev.fct_events
ALTER COLUMN click_count TYPE BIGINT;
Существующие файлы хранят INT значения. Iceberg при чтении автоматически приводит их к BIGINT без перезаписи. Это возможно, потому что Iceberg знает исходный тип (через Column ID) и умеет безопасно конвертировать при десериализации.
Удаление колонки (DROP COLUMN)¶
ALTER TABLE dbt_dev.dim_users DROP COLUMN legacy_field;
В Iceberg удаление колонки - логическая операция: колонка помечается как удалённая в метаданных, но физически из файлов не удаляется. При чтении данная колонка просто не возвращается. Это значит, что:
- Существующие файлы не перезаписываются (экономия ресурсов)
- Если колонку восстановить с тем же Column ID - данные будут доступны из старых файлов
- Физическое удаление произойдёт только при compaction (rewrite data files)
Эволюция партиционирования¶
Уникальная возможность Iceberg - изменить стратегию партиционирования без переписывания существующих данных. В Hive смена схемы партиционирования требует полного пересоздания таблицы.
-- Было: партиционирование по месяцам
-- Стало: партиционирование по дням (данные растут, нужна более мелкая партиция)
ALTER TABLE iceberg.analytics.fct_events
ADD PARTITION FIELD days(event_ts);
ALTER TABLE iceberg.analytics.fct_events
DROP PARTITION FIELD months(event_ts);
После этого:
- Старые данные остаются в файлах с месячным партиционированием
- Новые данные записываются с дневным партиционированием
- Iceberg понимает оба формата и корректно обрабатывает смешанные данные
Это называется Partition Evolution - ещё одна killer-фича Iceberg, которой нет в Hive и Delta Lake.
Incremental-модели с Iceberg в dbt¶
append: только добавление¶
-- models/staging/stg_raw_events.sql
{{ config(
materialized='incremental',
incremental_strategy='append',
file_format='iceberg',
location_root='s3a://datalake/iceberg/silver/',
partition_by={
'field': 'event_ts',
'data_type': 'timestamp',
'granularity': 'day'
},
tblproperties={
'write.format.default': 'parquet',
'write.parquet.compression-codec': 'zstd',
'write.target-file-size-bytes': '134217728',
'write.object-storage.enabled': 'true',
'history.expire.max-snapshot-age-ms': '2592000000',
}
) }}
SELECT
event_id,
user_id,
session_id,
event_type,
event_ts,
channel,
COALESCE(revenue, 0.0) AS revenue,
CURRENT_TIMESTAMP() AS processed_at
FROM {{ source('bronze', 'raw_events') }}
{% if is_incremental() %}
WHERE event_ts > (SELECT MAX(event_ts) FROM {{ this }})
{% endif %}
Стратегия append транслируется в INSERT INTO:
INSERT INTO dbt_dev.stg_raw_events
SELECT ... FROM bronze.raw_events
WHERE event_ts > (SELECT MAX(event_ts) FROM dbt_dev.stg_raw_events)
Каждый вызов создаёт новый Iceberg snapshot с добавленными файлами.
merge: UPSERT через MERGE INTO¶
-- models/marts/dim_users.sql
{{ config(
materialized='incremental',
incremental_strategy='merge',
unique_key='user_id',
file_format='iceberg',
location_root='s3a://datalake/iceberg/gold/',
partition_by={
'field': 'signup_date',
'data_type': 'date',
'granularity': 'month'
},
tblproperties={
'write.format.default': 'parquet',
'write.target-file-size-bytes': '67108864',
'write.distribution-mode': 'hash',
'write.object-storage.enabled': 'true',
'history.expire.max-snapshot-age-ms': '604800000',
'history.expire.min-snapshots-to-keep': '5',
}
) }}
WITH source AS (
SELECT * FROM {{ ref('stg_users') }}
{% if is_incremental() %}
WHERE updated_at >= (SELECT MAX(updated_at) FROM {{ this }})
{% endif %}
)
SELECT
user_id,
LOWER(email) AS email,
first_name,
last_name,
CAST(created_at AS DATE) AS signup_date,
country_code,
account_type,
loyalty_score,
CURRENT_TIMESTAMP() AS dbt_updated_at
FROM source
dbt с Iceberg и стратегией merge генерирует MERGE INTO:
MERGE INTO dbt_dev.dim_users AS DBT_INTERNAL_DEST
USING (
SELECT
user_id, email, first_name, last_name,
signup_date, country_code, account_type,
loyalty_score, dbt_updated_at
FROM dbt_dev.dim_users__dbt_tmp
) AS DBT_INTERNAL_SOURCE
ON (DBT_INTERNAL_SOURCE.user_id = DBT_INTERNAL_DEST.user_id)
WHEN MATCHED THEN UPDATE SET
email = DBT_INTERNAL_SOURCE.email,
first_name = DBT_INTERNAL_SOURCE.first_name,
loyalty_score = DBT_INTERNAL_SOURCE.loyalty_score,
dbt_updated_at = DBT_INTERNAL_SOURCE.dbt_updated_at
WHEN NOT MATCHED THEN INSERT (
user_id, email, first_name, last_name,
signup_date, country_code, account_type,
loyalty_score, dbt_updated_at
) VALUES (
DBT_INTERNAL_SOURCE.user_id, ...
)
insert_overwrite: перезапись партиций¶
-- models/marts/fct_daily_metrics.sql
{{ config(
materialized='incremental',
incremental_strategy='insert_overwrite',
file_format='iceberg',
location_root='s3a://datalake/iceberg/gold/',
partition_by={
'field': 'event_date',
'data_type': 'date',
'granularity': 'day'
}
) }}
SELECT
user_id,
CAST(event_ts AS DATE) AS event_date,
channel,
COUNT(DISTINCT session_id) AS sessions,
SUM(revenue) AS total_revenue,
COUNT(*) AS event_count
FROM {{ ref('stg_events') }}
{% if is_incremental() %}
WHERE CAST(event_ts AS DATE) >= DATE_ADD(CURRENT_DATE(), -7)
{% endif %}
GROUP BY user_id, CAST(event_ts AS DATE), channel
Time Travel: работа с историческими версиями¶
Чтение по номеру snapshot¶
SELECT *
FROM iceberg.analytics.fct_daily_metrics
VERSION AS OF 1003;
VERSION AS OF <snapshot_id> читает состояние таблицы в момент snapshot'а с этим ID.
Чтение по времени¶
SELECT *
FROM iceberg.analytics.fct_daily_metrics
TIMESTAMP AS OF '2024-01-15 12:00:00';
SELECT *
FROM iceberg.analytics.fct_daily_metrics
FOR SYSTEM_TIME AS OF TIMESTAMP_SECONDS(1705320000);
TIMESTAMP AS OF находит ближайший snapshot до указанного момента времени и читает данные из него.
Просмотр истории snapshot'ов¶
SELECT
committed_at,
snapshot_id,
parent_id,
operation,
summary
FROM iceberg.analytics.fct_daily_metrics.snapshots
ORDER BY committed_at DESC
LIMIT 10;
+--------------------+-------------------+---------+-----------+----------------------------------+
| committed_at | snapshot_id | parent_id| operation | summary|
+--------------------+-------------------+---------+-----------+----------------------------------+
|2024-01-15 14:30:00| 8834746592304881| 881... | append |{added-files: 4, added-records: ...|
|2024-01-15 10:15:00| 8813488294892744| 777... | merge |{changed-partition-count: 3, ... |
|2024-01-14 23:45:00| 7774839201048823| 521... | append |{added-files: 2, added-records: ...|
Поле operation показывает тип операции: append, overwrite, delete, merge. summary содержит статистику изменений: сколько файлов добавлено/удалено, сколько строк.
Metadata таблицы Iceberg¶
SELECT snapshot_id, added_data_files_count, deleted_data_files_count, total_data_files_count
FROM iceberg.analytics.fct_daily_metrics.snapshots;
SELECT file_path, file_format, record_count, file_size_in_bytes, partition
FROM iceberg.analytics.fct_daily_metrics.files
WHERE record_count < 10000;
SELECT snapshot_id, path, added_snapshot_id, deleted_snapshot_id
FROM iceberg.analytics.fct_daily_metrics.manifests;
SELECT *
FROM iceberg.analytics.fct_daily_metrics.history
ORDER BY made_current_at DESC;
Metadata-таблицы - мощный инструмент для диагностики и мониторинга. files-таблица позволяет найти все мелкие файлы (small files problem). manifests показывает количество накопленных manifest'ов (индикатор необходимости compaction).
Обслуживание таблицы: Compaction и Cleanup¶
Проблема накопления мелких файлов¶
Каждая инкрементальная загрузка через append создаёт новые файлы. Если загрузки происходят каждый час, за месяц накапливается 720 файлов. При чтении Spark должен открыть все 720 файлов - это медленно.
SELECT
partition,
COUNT(*) AS file_count,
SUM(file_size_in_bytes) / 1024 / 1024 AS total_mb,
AVG(file_size_in_bytes) / 1024 / 1024 AS avg_mb,
SUM(record_count) AS total_records
FROM iceberg.analytics.fct_daily_events.files
GROUP BY partition
ORDER BY file_count DESC
LIMIT 20;
Если avg_mb ниже 50 МБ и file_count выше 10 на партицию - время для compaction.
Compaction через Spark Procedures¶
CALL iceberg.system.rewrite_data_files(
table => 'analytics.fct_daily_events',
strategy => 'binpack',
options => map(
'target-file-size-bytes', '134217728',
'min-file-size-bytes', '67108864',
'max-file-size-bytes', '268435456'
)
);
rewrite_data_files - Spark SQL Procedure для compaction. binpack стратегия объединяет мелкие файлы в один целевого размера без перераспределения данных. Есть также sort стратегия - для Z-ordering данных для лучшей компрессии и predicate pushdown.
Удаление старых snapshot'ов¶
CALL iceberg.system.expire_snapshots(
table => 'analytics.fct_daily_events',
older_than => TIMESTAMP '2024-01-01 00:00:00',
retain_last => 5
);
expire_snapshots помечает snapshot'ы старше older_than как expired, сохраняя минимум retain_last snapshot'ов. Это освобождает ссылки, но не удаляет сами файлы данных.
CALL iceberg.system.remove_orphan_files(
table => 'analytics.fct_daily_events',
older_than => TIMESTAMP '2024-01-08 00:00:00'
);
remove_orphan_files физически удаляет файлы, которые не привязаны ни к одному живому snapshot'у. Запускайте после expire_snapshots.
Compaction manifest'ов¶
CALL iceberg.system.rewrite_manifests(
table => 'analytics.fct_daily_events'
);
Объединяет мелкие manifest-файлы в несколько крупных. Ускоряет планирование запросов за счёт уменьшения числа файлов метаданных для чтения.
Операции обслуживания в dbt¶
-- macros/iceberg_maintenance.sql
{% macro run_iceberg_maintenance(table_name, retain_snapshots=5, days_to_keep=7) %}
{%- set expire_ts = modules.datetime.datetime.now() - modules.datetime.timedelta(days=days_to_keep) -%}
CALL iceberg.system.expire_snapshots(
table => '{{ table_name }}',
older_than => TIMESTAMP '{{ expire_ts.strftime("%Y-%m-%d %H:%M:%S") }}',
retain_last => {{ retain_snapshots }}
);
CALL iceberg.system.remove_orphan_files(
table => '{{ table_name }}',
older_than => TIMESTAMP '{{ expire_ts.strftime("%Y-%m-%d %H:%M:%S") }}'
);
CALL iceberg.system.rewrite_data_files(
table => '{{ table_name }}',
strategy => 'binpack',
options => map('target-file-size-bytes', '134217728')
);
{% endmacro %}
-- models/maintenance/iceberg_maintenance.sql
{{ config(materialized='table') }}
{{ run_iceberg_maintenance('analytics.fct_daily_events', retain_snapshots=5, days_to_keep=7) }}
Лабораторная практика: Schema Drift в production¶
Бизнес-кейс¶
Продуктовая команда добавляет новую метрику loyalty_score в профили пользователей. Это поле появится в источнике данных - нужно расширить схему Gold-витрины без остановки аналитических дашбордов.
Шаг 1: Базовая модель (версия 1)¶
-- models/marts/dim_users.sql - версия 1
{{ config(
materialized='table',
file_format='iceberg',
location_root='s3a://datalake/iceberg/gold/',
tblproperties={
'write.format.default': 'parquet',
'write.parquet.compression-codec': 'zstd',
'write.object-storage.enabled': 'true',
'write.object-storage.path': 's3a://datalake/iceberg/data/',
'history.expire.max-snapshot-age-ms': '604800000',
'history.expire.min-snapshots-to-keep': '3',
}
) }}
SELECT
user_id,
name,
email,
CAST(created_at AS DATE) AS signup_date,
country_code
FROM {{ ref('stg_users') }}
WHERE user_id IS NOT NULL
dbt run --select dim_users
12:00:10 Running with dbt=1.7.3
12:00:15 1 of 1 START table model dbt_dev.dim_users [RUN]
12:00:28 1 of 1 OK created table model dbt_dev.dim_users [CREATE TABLE in 13.4s]
Шаг 2: Проверка начального состояния¶
DESCRIBE TABLE EXTENDED iceberg.dbt_dev.dim_users;
SELECT snapshot_id, committed_at, operation
FROM iceberg.dbt_dev.dim_users.snapshots
ORDER BY committed_at DESC;
+------------------+---------------------+-----------+
| snapshot_id | committed_at | operation |
+------------------+---------------------+-----------+
| 8834746592304881 | 2024-01-15 12:00:28 | append |
+------------------+---------------------+-----------+
Schema:
user_id BIGINT
name STRING
email STRING
signup_date DATE
country_code STRING
Шаг 3: Изменение схемы (Schema Drift)¶
Бизнес-изменение: поле name переименовать в full_name, добавить loyalty_score (INT), изменить тип внутренней метрики с INT на BIGINT (в источнике):
-- models/marts/dim_users.sql - версия 2 (с изменениями схемы)
{{ config(
materialized='table',
file_format='iceberg',
location_root='s3a://datalake/iceberg/gold/',
tblproperties={
'write.format.default': 'parquet',
'write.parquet.compression-codec': 'zstd',
'write.object-storage.enabled': 'true',
'write.object-storage.path': 's3a://datalake/iceberg/data/',
'history.expire.max-snapshot-age-ms': '604800000',
'history.expire.min-snapshots-to-keep': '3',
}
) }}
SELECT
user_id,
CONCAT(first_name, ' ', last_name) AS full_name,
email,
CAST(created_at AS DATE) AS signup_date,
country_code,
COALESCE(CAST(loyalty_score AS INT), 0) AS loyalty_score,
CAST(session_count AS BIGINT) AS session_count
FROM {{ ref('stg_users') }}
WHERE user_id IS NOT NULL
Шаг 4: Второй запуск - применение Schema Drift¶
dbt run --select dim_users
12:30:00 1 of 1 START table model dbt_dev.dim_users [RUN]
12:30:02 Applying schema changes to iceberg table...
12:30:02 ALTER TABLE dbt_dev.dim_users RENAME COLUMN name TO full_name
12:30:02 ALTER TABLE dbt_dev.dim_users ADD COLUMN loyalty_score INT
12:30:02 ALTER TABLE dbt_dev.dim_users ADD COLUMN session_count BIGINT
12:30:15 1 of 1 OK created table model dbt_dev.dim_users [ALTER + INSERT in 13.2s]
Шаг 5: Анализ результата¶
SELECT snapshot_id, committed_at, operation, summary
FROM iceberg.dbt_dev.dim_users.snapshots
ORDER BY committed_at DESC;
+------------------+---------------------+-----------+------------------------------------+
| snapshot_id | committed_at | operation | summary |
+------------------+---------------------+-----------+------------------------------------+
| 9945837201048823 | 2024-01-15 12:30:15 | overwrite| {added-files: 3, deleted-files: 2} |
| 8834746592304881 | 2024-01-15 12:00:28 | append | {added-files: 2, added-records: ...}|
+------------------+---------------------+-----------+------------------------------------+
Snapshot 9945837201048823 с операцией overwrite - это результат обновления схемы. В summary: added-files: 3, deleted-files: 2 - новые файлы записаны, старые помечены к удалению (но физически пока не удалены!).
SELECT *
FROM iceberg.dbt_dev.dim_users
TIMESTAMP AS OF '2024-01-15 12:00:00';
+--------+-------+------------------+------------+-------------+
|user_id | name | email |signup_date |country_code |
+--------+-------+------------------+------------+-------------+
| 1 | Alice | alice@example.com| 2023-01-01 | US |
Time travel к моменту до изменения схемы: видим старую схему с колонкой name (не full_name) - данные сохранены исторически.
SELECT user_id, full_name, loyalty_score, session_count
FROM iceberg.dbt_dev.dim_users
LIMIT 5;
+--------+-------------+-------------+---------------+
|user_id | full_name | loyalty_score| session_count|
+--------+-------------+-------------+---------------+
| 1 | Alice Smith | 75 | 1250|
| 2 | Bob Johnson | 0 | 320|
Текущая версия - уже с новой схемой.
Мульти-движковая совместимость¶
Чтение Iceberg-таблиц из Trino¶
-- Trino
SELECT
user_id,
full_name,
loyalty_score
FROM hive.analytics.dim_users
WHERE signup_date >= DATE '2024-01-01'
AND country_code = 'US';
Trino читает те же Iceberg-файлы что и Spark, через тот же Hive Metastore. Schema evolution прозрачна - Trino видит актуальную схему без какой-либо синхронизации.
Чтение из Flink¶
-- Flink SQL
CREATE TABLE dim_users_flink (
user_id BIGINT,
full_name STRING,
loyalty_score INT
) WITH (
'connector' = 'iceberg',
'catalog-name' = 'hive_catalog',
'catalog-type' = 'hive',
'uri' = 'thrift://hive-metastore:9083',
'database-name' = 'analytics',
'table-name' = 'dim_users'
);
SELECT * FROM dim_users_flink WHERE loyalty_score > 50;
Flink читает те же данные, что записал dbt через Spark. Это и есть мульти-движковая совместимость Iceberg - одни данные, разные движки.
Антипаттерны¶
Антипаттерн 1: Оставить snapshot retention по умолчанию¶
Без явного history.expire.max-snapshot-age-ms snapshot'ы накапливаются бесконечно. Через год таблица хранит тысячи snapshot'ов - это гигабайты метаданных и замедленное планирование запросов. Всегда задавайте retention.
Антипаттерн 2: Hive-style мышление в Iceberg¶
-- ПЛОХО: Hive-style partition col
{{ config(
partition_by='event_date'
) }}
SELECT
...,
event_date AS event_date,
...
FROM source
В Hive нужна явная колонка-партиция в схеме таблицы. В Iceberg - используйте hidden partition transforms на существующей колонке:
{{ config(
partition_by={
'field': 'event_ts',
'data_type': 'timestamp',
'granularity': 'day'
}
) }}
Колонка event_ts уже есть в данных - дополнительная колонка event_date не нужна.
Антипаттерн 3: Не запускать compaction для append-моделей¶
Если модель со стратегией append запускается каждый час без compaction - через месяц у каждой партиции по 720 мелких файлов. Запросы замедляются в разы. Планируйте регулярный запуск rewrite_data_files.
Антипаттерн 4: Переименовывать колонки без координации с downstream¶
Schema evolution в Iceberg физически безопасна, но не магия. Если колонку переименовали в схеме dbt-модели, а Trino-запросы или Superset-дашборды обращаются к ней по старому имени - они сломаются. Координируйте schema evolution с командами-потребителями данных.
Антипаттерн 5: Использовать write.format.default=avro для аналитических таблиц¶
Avro - строковый формат, оптимальный для streaming и CDC. Для аналитических запросов (агрегации, фильтрации по колонкам) Parquet или ORC значительно эффективнее из-за колоночного хранения и predicate pushdown. Всегда используйте parquet для Gold-витрин.
Антипаттерн 6: Игнорировать object storage optimization¶
tblproperties:
write.format.default: 'parquet'
# write.object-storage.enabled НЕ задан
Без write.object-storage.enabled=true Iceberg пишет файлы в детерминированный путь. При высокой нагрузке все запросы попадают в один S3-prefix → throttling → замедление записи. Для production на S3/MinIO всегда включайте object storage optimization.
Чеклист¶
При создании Iceberg-модели:
file_format='iceberg'в configlocation_rootуказывает на правильное место в S3/MinIO- Hidden partition transform задан через
partition_byсgranularity tblpropertiesсодержит минимально необходимые свойства (format, compression, target-file-size, snapshot retention)write.object-storage.enabled=trueдля S3/MinIO
Для incremental-моделей:
- Выбрана правильная стратегия (append/merge/insert_overwrite)
- Watermark в блоке
{% if is_incremental() %}покрывает window late-arriving data unique_keyпри merge включает все поля уникальности строки
Для production:
history.expire.max-snapshot-age-msзадан (7–30 дней)history.expire.min-snapshots-to-keepзадан (3–5)- Запланирован регулярный compaction (
rewrite_data_files) - Запланирован
expire_snapshots+remove_orphan_files
Домашнее задание¶
Дан dbt-проект с моделью fct_orders (материализация table, Parquet). Модель содержит агрегаты заказов по (customer_id, order_date).
Задание 1. Переведите модель на Iceberg:
- Настройте скрытое партиционирование по
order_dateс гранулярностьюmonth - Настройте
tblpropertiesдля production (compression, target file size, object storage optimization, snapshot retention 7 дней)
Задание 2. Измените материализацию с table на incremental со стратегией insert_overwrite:
- Напишите правильный watermark в
{% if is_incremental() %}- учтите, что заказы могут обновляться (статус меняется) в течение 14 дней - Обоснуйте выбор
insert_overwriteа неmerge
Задание 3. Добавьте в схему новую колонку discount_amount DECIMAL(10,2) и измените тип order_count с INT на BIGINT. Запустите dbt run и приложите:
- Вывод команды dbt run
- Результат
SELECT * FROM iceberg.analytics.fct_orders.snapshots
Задание 4. Напишите SQL-запрос, который использует time travel для сравнения агрегатов в текущей версии таблицы и версии неделю назад:
SELECT
'current' AS version,
SUM(total_revenue) AS revenue
FROM iceberg.analytics.fct_orders
UNION ALL
SELECT
'last_week' AS version,
SUM(total_revenue) AS revenue
FROM iceberg.analytics.fct_orders
TIMESTAMP AS OF ...