dbt Incremental Models: append, insert_overwrite и merge с unique_key
Три стратегии инкрементальной загрузки в dbt: append для иммутабельных логов, insert_overwrite для партиционированных агрегатов, merge для CDC и SCD на Delta/Iceberg.
Концепция инкрементальной загрузки¶
Экономика вычислений в Big Data¶
Представьте пайплайн, который каждый день считает агрегаты по кликстриму за весь год: 365 дней × 100 ГБ в день = 36.5 ТБ данных. При полном пересчёте (materialized='table') каждый запуск сканирует всю историю - 36.5 ТБ shuffle, часы работы кластера, тысячи дорогостоящих compute-часов.
Но задумайтесь: сколько из этих 36.5 ТБ изменилось за вчера? Только один день - 100 ГБ. Если добавились новые события за вчера, агрегаты за предыдущие 364 дня остались теми же. Зачем пересчитывать всё?
Именно здесь и нужна инкрементальная загрузка (incremental processing): обрабатываем только новые или изменённые данные, а результат добавляем к уже существующей таблице. Вместо 36.5 ТБ - обрабатываем 100 ГБ. Экономия в 365 раз.
Идемпотентность: золотое правило¶
Инкрементальные пайплайны должны быть идемпотентными: повторный запуск за тот же период времени не должен создавать дублей и не должен искажать данные. Это критически важное свойство для production-систем.
Почему это важно: оркестраторы (Airflow, Dagster, Prefect) иногда перезапускают задачи при сбоях. Если пайплайн за вчера упал в середине - оркестратор запустит его снова. Если пайплайн не идемпотентен, данные за вчера окажутся записаны дважды.
Пример проблемы: append-стратегия с watermark WHERE event_ts > MAX(event_ts) выглядит идемпотентной, но при сбое в процессе записи - после частичной записи MAX(event_ts) уже сдвинулся, и при повторном запуске часть событий будет пропущена.
Идемпотентность требует продуманного дизайна для каждой стратегии:
append- идемпотентна только для строго монотонных источников без возможности дублей при перезапускеinsert_overwrite- идемпотентна по определению: перезапись партиции всегда даёт одинаковый результатmerge- идемпотентна при правильномunique_key: повторная вставка той же строки обновит, но не продублирует
is_incremental(): как dbt понимает режим выполнения¶
Макрос {% if is_incremental() %} - ключевой механизм инкрементальных моделей. Он возвращает true только если:
- Таблица уже существует в базе данных (создана при предыдущем запуске)
- Запуск не является
--full-refresh - Материализация модели -
incremental
При первом запуске (таблица не существует) is_incremental() возвращает false, и модель выполняет полный запрос без фильтрации. Это и есть "первоначальная загрузка". При всех последующих запусках - true, и активируется блок с фильтрацией новых данных.
{% if is_incremental() %}
WHERE event_ts > (SELECT MAX(event_ts) FROM {{ this }})
{% endif %}
Здесь {{ this }} - ссылка на саму текущую таблицу. При первом запуске этого блока нет - в таблице ещё нет данных, фильтровать нечего. При последующих - MAX(event_ts) возвращает последнее известное значение, и мы берём только более свежие данные.
Жизненный цикл выполнения incremental-модели¶
Диаграмма показывает полный жизненный цикл. При первом запуске dbt проверяет существование таблицы - если её нет, выполняется полная загрузка (CTAS). При последующих запусках - инкрементальная: создаётся временная таблица model__dbt_tmp с отфильтрованными новыми данными, затем применяется выбранная стратегия. Результат фиксируется как новый snapshot (Iceberg) или commit (Delta Lake).
Стратегия append: максимальная скорость¶
Механика append¶
Стратегия append - простейшая инкрементальная стратегия. dbt просто добавляет новые строки в конец существующей таблицы через INSERT INTO. Никакого сравнения с существующими данными, никакой дедупликации, никакого поиска совпадений.
Что генерирует dbt при append:
-- Шаг 1: Создание временной таблицы с новыми данными
CREATE TABLE dbt_dev.stg_events__dbt_tmp AS
SELECT
event_id,
user_id,
event_ts,
event_type,
revenue
FROM bronze.raw_events
WHERE event_ts > (SELECT MAX(event_ts) FROM dbt_dev.stg_events);
-- Шаг 2: Вставка в целевую таблицу
INSERT INTO dbt_dev.stg_events
SELECT * FROM dbt_dev.stg_events__dbt_tmp;
-- Шаг 3: Удаление временной таблицы
DROP TABLE IF EXISTS dbt_dev.stg_events__dbt_tmp;
Для Iceberg это выглядит немного иначе - INSERT INTO создаёт новый snapshot с добавленными файлами, не изменяя существующих.
Конфигурация append-модели¶
-- models/staging/stg_events.sql
{{ config(
materialized='incremental',
incremental_strategy='append',
file_format='iceberg',
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',
}
) }}
SELECT
event_id,
CAST(user_id AS BIGINT) AS user_id,
session_id,
event_type,
CAST(event_ts AS TIMESTAMP) AS event_ts,
CAST(event_ts AS DATE) AS event_date,
LOWER(TRIM(channel)) AS channel,
COALESCE(CAST(revenue AS DOUBLE), 0.0) AS revenue,
CURRENT_TIMESTAMP() AS processed_at
FROM {{ source('bronze', 'raw_events') }}
WHERE event_id IS NOT NULL
AND user_id IS NOT NULL
{% if is_incremental() %}
AND event_ts > (
SELECT COALESCE(MAX(event_ts), TIMESTAMP '1970-01-01')
FROM {{ this }}
)
{% endif %}
Важная деталь: COALESCE(MAX(event_ts), TIMESTAMP '1970-01-01'). Это защита от ситуации, когда таблица пустая (например, после ручной очистки). Без COALESCE подзапрос вернёт null, и условие event_ts > null никогда не выполнится - все строки будут отфильтрованы.
Watermark: разные паттерны фильтрации¶
Watermark - это способ определить "граница между старым и новым". Выбор watermark зависит от природы данных:
Фильтрация по timestamp (монотонный источник):
{% if is_incremental() %}
AND event_ts > (SELECT MAX(event_ts) FROM {{ this }})
{% endif %}
Подходит для источников, где event_ts строго возрастает и события никогда не приходят с опозданием. Для Kafka-топиков с гарантиями at-least-once - опасно: при повторной доставке сообщения появятся дубли.
Фильтрация по дате с перекрытием (late-arriving data):
{% if is_incremental() %}
AND event_date >= DATE_ADD(CURRENT_DATE(), -3)
{% endif %}
Всегда перезагружаем последние 3 дня. Это обрабатывает события с задержкой (late-arriving data), но при append-стратегии создаёт дубли! Подходит только для insert_overwrite и merge.
Фильтрация по целочисленному ID (append-only события):
{% if is_incremental() %}
AND event_id > (SELECT COALESCE(MAX(event_id), 0) FROM {{ this }})
{% endif %}
Подходит для источников с автоинкрементным ID. Гарантирует: каждый event_id обрабатывается ровно один раз.
Фильтрация по partition ID (для партиционированных источников):
{% if is_incremental() %}
AND partition_date > (
SELECT COALESCE(MAX(partition_date), DATE '1970-01-01')
FROM {{ this }}
)
{% endif %}
Если Bronze-данные партиционированы по дате, читаем только новые партиции целиком.
Когда использовать append¶
Append - идеальная стратегия для иммутабельных событий: данных, которые после создания никогда не изменяются. Классические примеры:
- Clickstream: каждый клик - уникальное событие с уникальным ID. Клик не может "измениться" - он произошёл или не произошёл
- IoT-метрики: показания датчиков - неизменяемые временные точки
- Платёжные транзакции: сам факт платежа неизменен (хотя его статус может меняться - это уже другая история)
- Логи запросов: каждая строка лога - уникальное событие
Append не подходит для:
- Данных с обновляемым статусом (
order_status: pending → completed) - SCD (Slowly Changing Dimensions) - таблиц пользователей, кампаний, продуктов
- Агрегатов, которые нужно пересчитать при приходе поздних данных
Риск дублей при append¶
Главная опасность append - дублирование при повторных запусках. Рассмотрим сценарий:
Запуск 1 (успешный):
Записаны события с event_ts <= 2024-01-15 10:00:00
MAX(event_ts) в таблице = 2024-01-15 10:00:00
Запуск 2 (падение при записи):
Прочитаны события с event_ts > 2024-01-15 10:00:00
Записана ЧАСТЬ событий (10:00–11:00) - Spark упал
MAX(event_ts) в таблице = 2024-01-15 11:00:00 (из частично записанных)
Запуск 3 (повторный):
Фильтр: event_ts > 2024-01-15 11:00:00
События за 10:00–11:00 ПРОПУЩЕНЫ (watermark уже после них)
Это не дубли, а потеря данных. Для append стратегии критична транзакционность хранилища - Delta Lake и Iceberg гарантируют атомарность записи: либо все файлы записаны, либо ни один (неудачный commit не меняет MAX).
Стратегия insert_overwrite: партиционированная замена¶
Механика insert_overwrite¶
insert_overwrite работает иначе: она смотрит, какие партиции присутствуют в новых данных, и полностью заменяет эти партиции в целевой таблице. Строки старых данных в затрагиваемых партициях удаляются, новые записываются.
Что генерирует dbt:
-- В Spark с dynamic partition overwrite:
INSERT OVERWRITE TABLE dbt_dev.fct_daily_aggregates
PARTITION (event_date)
SELECT
event_date,
channel,
SUM(revenue) AS total_revenue,
COUNT(*) AS event_count
FROM dbt_dev.stg_events__dbt_tmp
GROUP BY event_date, channel;
Spark с настройкой spark.sql.sources.partitionOverwriteMode=dynamic сам определяет, какие партиции event_date присутствуют в данных, и перезаписывает только их. Другие партиции остаются нетронутыми.
Настройка dynamic partition overwrite¶
Ключевая настройка для insert_overwrite в Spark:
# profiles.yml
server_side_parameters:
spark.sql.sources.partitionOverwriteMode: "dynamic"
Без этой настройки Spark использует статический режим перезаписи: INSERT OVERWRITE затрагивает всю таблицу, удаляя данные из всех партиций. С dynamic - только те партиции, которые явно указаны в новых данных.
Проверить режим:
SET spark.sql.sources.partitionOverwriteMode;
-- spark.sql.sources.partitionOverwriteMode dynamic
Конфигурация insert_overwrite-модели¶
-- models/marts/fct_daily_aggregates.sql
{{ config(
materialized='incremental',
incremental_strategy='insert_overwrite',
file_format='delta',
partition_by=['event_date'],
tblproperties={
'delta.autoOptimize.optimizeWrite': 'true',
'delta.autoOptimize.autoCompact': 'true',
}
) }}
WITH events AS (
SELECT * FROM {{ ref('stg_events') }}
{% if is_incremental() %}
WHERE event_date >= DATE_ADD(CURRENT_DATE(), -7)
{% endif %}
)
SELECT
event_date,
channel,
campaign_id,
COUNT(DISTINCT user_id) AS unique_users,
COUNT(DISTINCT session_id) AS sessions,
SUM(revenue) AS total_revenue,
COUNT(*) AS event_count,
CURRENT_TIMESTAMP() AS updated_at
FROM events
GROUP BY event_date, channel, campaign_id
Здесь WHERE event_date >= DATE_ADD(CURRENT_DATE(), -7) - окно в 7 дней. dbt прочитает только данные за последние 7 дней, агрегирует их, и перезапишет партиции event_date за эти 7 дней в целевой таблице. Данные за более ранние даты останутся нетронутыми.
Почему 7 дней, а не 1? Late-arriving data: события могут приходить с задержкой до 3–7 дней (мобильные приложения в offline-режиме, ретри-очереди, CDC с задержкой). Перезапись 7 дней гарантирует, что все поздние события попадут в правильную дату.
Идемпотентность insert_overwrite¶
insert_overwrite идемпотентна по своей природе. Повторный запуск за тот же период:
- Читает те же данные (те же 7 дней)
- Агрегирует - получает те же результаты
- Перезаписывает те же партиции теми же данными
Итог: таблица в точности том же состоянии. Это делает insert_overwrite самой надёжной стратегией с точки зрения идемпотентности.
Когда использовать insert_overwrite¶
Идеальные сценарии:
- Ежедневные/еженедельные агрегаты: дневная выручка по кампаниям, недельные когорты. Исторические данные за прошлые месяцы не меняются, но последние N дней могут корректироваться из-за late data
- ETL с партиционированием по дате: Silver-слой с партиционированием по
event_date, куда поступают данные с задержкой - Пересчёт агрегатов при изменении бизнес-логики: изменили формулу метрики - перезапустите с
--full-refreshили перезапишите нужные партиции вручную
Ограничения:
- Нет построчного UPDATE: нельзя обновить одну строку внутри партиции - только заменить всю партицию
- Гранулярность - уровень партиции: если партиция
event_dateсодержит миллионы строк, а нужно обновить одну - придётся перезаписывать всю партицию - Требует партиционирования: без
partition_byв конфиге нет смысла в этой стратегии
Стратегия merge: точечный UPSERT¶
Что такое MERGE INTO¶
MERGE - SQL-операция, которая объединяет логику INSERT и UPDATE в одном выражении. Она сравнивает строки источника и целевой таблицы по ключу и:
- Если строка найдена в обоих - выполняет UPDATE (обновление)
- Если строка только в источнике - выполняет INSERT (вставка)
- Опционально: если строка только в цели - выполняет DELETE
Это называется UPSERT (UPDATE + INSERT). В Spark SQL это MERGE INTO:
MERGE INTO dbt_dev.dim_users AS target
USING dbt_dev.dim_users__dbt_tmp AS source
ON (target.user_id = source.user_id)
WHEN MATCHED THEN
UPDATE SET
email = source.email,
loyalty_score = source.loyalty_score,
country_code = source.country_code,
updated_at = source.updated_at
WHEN NOT MATCHED THEN
INSERT (user_id, email, loyalty_score, country_code, created_at, updated_at)
VALUES (source.user_id, source.email, source.loyalty_score,
source.country_code, source.created_at, source.updated_at);
Ключевая строка: ON (target.user_id = source.user_id) - это условие сопоставления. Строка из источника и строка из цели считаются "одной и той же" если user_id совпадает.
Зависимость от формата: Delta Lake и Iceberg¶
MERGE - транзакционная операция. Она требует:
- Атомарности: либо все изменения применены, либо ни одно
- Изоляции: читатели не видят частично изменённые данные
- Консистентности: после merge таблица в валидном состоянии
Обычный Parquet в HDFS/S3 не поддерживает MERGE - нет механизма атомарной замены строк. Только Delta Lake и Apache Iceberg реализуют MERGE через свои транзакционные механизмы:
- Delta Lake: Transaction Log (delta/_delta_log/) фиксирует все изменения как JSON-записи; каждый merge - атомарный commit в этом логе
- Iceberg: Snapshot - новая неизменяемая версия таблицы; merge создаёт новый snapshot с обновлёнными файлами
Попытка использовать merge с file_format='parquet' приведёт к ошибке.
Механика merge под капотом¶
Когда dbt выполняет merge, происходит следующее:
Шаг 1: dbt создаёт временную таблицу model__dbt_tmp с новыми/изменёнными данными:
CREATE TABLE dbt_dev.dim_users__dbt_tmp AS
SELECT * FROM dbt_dev.stg_users
WHERE updated_at >= (SELECT MAX(updated_at) FROM dbt_dev.dim_users);
Шаг 2: dbt выполняет MERGE INTO:
MERGE INTO dbt_dev.dim_users AS DBT_INTERNAL_DEST
USING 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 ...
WHEN NOT MATCHED THEN INSERT ...;
Шаг 3: Delta/Iceberg создаёт новый snapshot/commit.
Шаг 4: dbt удаляет временную таблицу.
Механизмы записи: Copy-on-Write vs Merge-on-Read¶
Это важная концепция для понимания производительности merge. Оба Delta Lake и Iceberg поддерживают оба механизма, но с разными компромиссами.
Copy-on-Write (CoW) - стандартный режим:
- При каждом UPDATE/DELETE/MERGE Spark читает все файлы, затронутые изменениями
- Переписывает их с применёнными изменениями
- Старые файлы помечаются как удалённые
- Новые файлы записываются
Результат: чтение быстрое (нет overhead), но запись дорогая - нужно переписывать целые Parquet-файлы из-за нескольких изменённых строк.
Merge-on-Read (MoR) - режим для частых обновлений (в Iceberg: delete files):
- При UPDATE/DELETE/MERGE создаётся "delete file" (список изменений)
- Сами data-файлы не переписываются
- При чтении Spark объединяет (merges) data files и delete files "на лету"
Результат: запись мгновенная, но чтение медленнее из-за объединения файлов при каждом запросе. Через некоторое время нужен compaction.
Для dbt gold-витрин с умеренными обновлениями (раз в день) - Copy-on-Write оптимален. Для CDC-пайплайнов с тысячами обновлений в минуту - рассмотрите MoR.
Конфигурация merge-модели¶
-- models/marts/dim_users.sql
{{ config(
materialized='incremental',
incremental_strategy='merge',
unique_key='user_id',
file_format='delta',
partition_by=['signup_month'],
tblproperties={
'delta.autoOptimize.optimizeWrite': 'true',
}
) }}
SELECT
user_id,
LOWER(email) AS email,
CONCAT(first_name, ' ', last_name) AS full_name,
CAST(created_at AS DATE) AS signup_date,
DATE_FORMAT(created_at, 'yyyy-MM') AS signup_month,
country_code,
account_type,
COALESCE(loyalty_score, 0) AS loyalty_score,
updated_at,
CURRENT_TIMESTAMP() AS dbt_updated_at
FROM {{ ref('stg_users') }}
{% if is_incremental() %}
WHERE updated_at > (
SELECT COALESCE(MAX(updated_at), TIMESTAMP '2000-01-01')
FROM {{ this }}
)
{% endif %}
unique_key='user_id' - dbt использует это поле в условии ON в MERGE. Каждый пользователь идентифицируется по user_id.
Составной unique_key¶
Для таблиц с составным первичным ключом:
{{ config(
materialized='incremental',
incremental_strategy='merge',
unique_key=['user_id', 'event_date'],
file_format='iceberg',
) }}
SELECT
user_id,
event_date,
channel,
SUM(revenue) AS total_revenue,
COUNT(*) AS event_count,
CURRENT_TIMESTAMP() AS updated_at
FROM {{ ref('stg_events') }}
{% if is_incremental() %}
WHERE event_date >= DATE_ADD(CURRENT_DATE(), -7)
{% endif %}
GROUP BY user_id, event_date, channel
С unique_key=['user_id', 'event_date'] dbt сгенерирует:
ON (
DBT_INTERNAL_SOURCE.user_id = DBT_INTERNAL_DEST.user_id
AND DBT_INTERNAL_SOURCE.event_date = DBT_INTERNAL_DEST.event_date
)
Строка считается "той же" только если совпадают оба поля.
Кастомное условие merge¶
В dbt 1.6+ появился параметр merge_condition для переопределения условия merge:
{{ config(
materialized='incremental',
incremental_strategy='merge',
unique_key='order_id',
merge_condition="DBT_INTERNAL_SOURCE.order_id = DBT_INTERNAL_DEST.order_id AND DBT_INTERNAL_DEST.status != 'completed'",
file_format='delta',
) }}
SELECT
order_id,
user_id,
status,
total_amount,
order_date,
updated_at
FROM {{ ref('stg_orders') }}
{% if is_incremental() %}
WHERE updated_at > (SELECT MAX(updated_at) FROM {{ this }})
{% endif %}
Это обновит только заказы, статус которых ещё не completed. Завершённые заказы не будут затронуты merge, даже если данные о них пришли в источнике.
merge для CDC-пайплайнов¶
CDC (Change Data Capture) - один из главных use case для merge. CDC-источник присылает не полный snapshot таблицы, а только изменения: INSERT, UPDATE, DELETE-операции.
Типичная структура CDC-записи:
{operation: "INSERT", user_id: 1, email: "alice@old.com", ts: "2024-01-01"}
{operation: "UPDATE", user_id: 1, email: "alice@new.com", ts: "2024-01-15"}
{operation: "DELETE", user_id: 1, ts: "2024-01-20"}
Обработка через dbt merge с дополнительным WHEN NOT MATCHED BY SOURCE:
-- Для Delta Lake: поддерживается DELETE в merge
{{ config(
materialized='incremental',
incremental_strategy='merge',
unique_key='user_id',
file_format='delta',
merge_exclude_columns=['created_at'],
) }}
{% if is_incremental() %}
WITH cdc_source AS (
SELECT
user_id,
email,
full_name,
operation,
updated_at,
ROW_NUMBER() OVER (
PARTITION BY user_id ORDER BY updated_at DESC
) AS rn
FROM {{ ref('stg_users_cdc') }}
WHERE updated_at > (SELECT MAX(updated_at) FROM {{ this }})
),
deduped AS (
SELECT * FROM cdc_source WHERE rn = 1
),
non_deleted AS (
SELECT * FROM deduped WHERE operation != 'DELETE'
)
SELECT user_id, email, full_name, updated_at FROM non_deleted
{% else %}
SELECT user_id, email, full_name, updated_at
FROM {{ ref('stg_users_cdc') }}
WHERE operation != 'DELETE'
{% endif %}
Обратите внимание на ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY updated_at DESC) - дедупликацию. Если в одном батче пришло несколько CDC-событий для одного пользователя, нужно взять только последнее. Без дедупликации merge упадёт с ошибкой "multiple rows matched".
Матрица выбора стратегии¶
Когда что использовать¶
| Сценарий | Данные | Стратегия | Почему |
|---|---|---|---|
| Clickstream, IoT, логи | Иммутабельные события, нет обновлений | append |
Максимальная скорость, нет overhead на merge |
| Дневные агрегаты | Могут пересчитаться при late data | insert_overwrite |
Идемпотентно, нет построчного сравнения |
| SCD Type 1 (текущий снимок) | Атрибуты сущностей меняются | merge |
Обновляет существующие строки |
| CDC из OLTP | INSERT/UPDATE/DELETE операции | merge |
Точная репликация изменений |
| Финансовые агрегаты | Нельзя допускать дубли, late data | insert_overwrite |
Идемпотентна по определению |
| Рекомендательные витрины | Пересчитываются полностью | table |
Нет инкрементальной логики, проще |
Матрица форматов¶
| Стратегия | Parquet | Delta Lake | Apache Iceberg |
|---|---|---|---|
append |
Да | Да | Да |
insert_overwrite |
Да (с dynamic mode) | Да | Да |
merge |
Нет | Да | Да |
Производительность incremental-моделей в Spark UI¶
Как выглядит append в Spark UI¶
При стратегии append Spark выполняет два Job'а:
- Job 0: чтение источника + фильтрация + вычисление
MAX(event_ts)→ запись в__dbt_tmp - Job 1:
INSERT INTO- чтение__dbt_tmp+ запись в целевую таблицу
В Spark UI Job 1 должен быть быстрым: простая копия данных без shuffle (если не нужна сортировка для Parquet-файлов). Если видите тяжёлый shuffle в Job 1 - возможно, нужна настройка write.distribution-mode.
Как выглядит insert_overwrite в Spark UI¶
insert_overwrite выполняет:
- Job 0: чтение источника за N дней + фильтрация → запись в
__dbt_tmp - Job 1: чтение
__dbt_tmp+ GROUP BY (если агрегация) + запись с перезаписью партиций
Job 1 с GROUP BY имеет характерный паттерн в Spark UI: Exchange operator (shuffle), HashAggregate. Тюнинг shuffle.partitions и AQE критичен для этого job'а.
Как выглядит merge в Spark UI¶
merge - самая тяжёлая операция в Spark UI. Типичный физический план:
MergeIntoTable
+-- SortMergeJoin [unique_key]
+-- Exchange (HashPartitioning) ← shuffle source
| +-- source table scan
+-- Exchange (HashPartitioning) ← shuffle target
+-- target table scan (затронутые файлы)
В Spark UI видно два Scan'а (источник и часть целевой таблицы), два Exchange'а (shuffle по unique_key), и SortMergeJoin. Обратите внимание на target table scan - Spark не сканирует всю целевую таблицу, а только файлы, потенциально затронутые merge (благодаря file-skipping в Delta/Iceberg).
Ключевые метрики в Spark UI для merge:
- Input Size у source scan - сколько новых данных обрабатывается
- Input Size у target scan - сколько исторических данных нужно прочитать для сравнения
- Shuffle Write/Read - объём shuffle по unique_key
- Duration - общее время merge (должно быть пропорционально объёму target scan)
Если target scan читает значительно больше данных, чем source - рассмотрите партиционирование по полю, совпадающему с watermark, чтобы file-skipping работал эффективнее.
Практика: сравнение стратегий на e-commerce данных¶
Бизнес-кейс¶
Пайплайн обработки заказов интернет-магазина. Особенность: статус заказа меняется - pending → processing → shipped → delivered. Заказ создан вчера, оплачен сегодня. Нужна витрина с актуальным статусом каждого заказа.
Шаг 1: Источник - Bronze-таблица заказов¶
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, current_timestamp
spark = SparkSession.builder \
.config("spark.sql.extensions",
"io.delta.sql.DeltaSparkSessionExtension") \
.config("spark.sql.catalog.spark_catalog",
"org.apache.spark.sql.delta.catalog.DeltaCatalog") \
.getOrCreate()
orders_v1 = spark.createDataFrame([
(1, 101, "2024-01-15", 250.00, "pending", "2024-01-15 09:00:00"),
(2, 102, "2024-01-15", 180.50, "pending", "2024-01-15 10:30:00"),
(3, 103, "2024-01-15", 420.00, "processing", "2024-01-15 11:00:00"),
(4, 104, "2024-01-14", 99.99, "shipped", "2024-01-14 14:00:00"),
(5, 105, "2024-01-14", 760.00, "delivered", "2024-01-14 08:00:00"),
], ["order_id", "user_id", "order_date", "amount", "status", "updated_at"])
orders_v1.write \
.format("delta") \
.mode("overwrite") \
.saveAsTable("bronze.raw_orders")
Шаг 2: Демонстрация проблемы с append¶
-- models/marts/fct_orders_append.sql (ПЛОХАЯ СТРАТЕГИЯ ДЛЯ ЭТИХ ДАННЫХ)
{{ config(
materialized='incremental',
incremental_strategy='append',
file_format='delta',
) }}
SELECT
order_id,
user_id,
CAST(order_date AS DATE) AS order_date,
amount,
status,
CAST(updated_at AS TIMESTAMP) AS updated_at,
CURRENT_TIMESTAMP() AS processed_at
FROM {{ source('bronze', 'raw_orders') }}
{% if is_incremental() %}
WHERE CAST(updated_at AS TIMESTAMP) > (
SELECT MAX(updated_at) FROM {{ this }}
)
{% endif %}
Запуск 1: все 5 заказов попали в таблицу. Хорошо.
Теперь статус заказов меняется - источник обновился:
orders_v2 = spark.createDataFrame([
(1, 101, "2024-01-15", 250.00, "processing", "2024-01-15 12:00:00"),
(2, 102, "2024-01-15", 180.50, "delivered", "2024-01-15 15:00:00"),
], ["order_id", "user_id", "order_date", "amount", "status", "updated_at"])
orders_v2.write.format("delta").mode("append").saveAsTable("bronze.raw_orders")
Запуск 2: заказы 1 и 2 появятся дважды: с pending и с processing/delivered. Теперь для заказа 1 в таблице два ряда!
SELECT order_id, status, COUNT(*) as cnt
FROM dbt_dev.fct_orders_append
GROUP BY order_id, status
ORDER BY order_id;
+--------+------------+---+
|order_id| status|cnt|
+--------+------------+---+
| 1| pending| 1|
| 1| processing| 1| ← ДУБЛЬ!
| 2| pending| 1|
| 2| delivered| 1| ← ДУБЛЬ!
| 3| processing| 1|
| 4| shipped| 1|
| 5| delivered| 1|
+--------+------------+---+
Это и есть проблема append для данных с обновлениями.
Шаг 3: Правильное решение - merge¶
-- models/marts/fct_orders.sql
{{ config(
materialized='incremental',
incremental_strategy='merge',
unique_key='order_id',
file_format='delta',
partition_by=['order_date'],
tblproperties={
'delta.autoOptimize.optimizeWrite': 'true',
}
) }}
SELECT
order_id,
user_id,
CAST(order_date AS DATE) AS order_date,
amount,
status,
CAST(updated_at AS TIMESTAMP) AS updated_at,
CURRENT_TIMESTAMP() AS dbt_processed_at
FROM {{ source('bronze', 'raw_orders') }}
{% if is_incremental() %}
WHERE CAST(updated_at AS TIMESTAMP) > (
SELECT COALESCE(MAX(updated_at), TIMESTAMP '2000-01-01')
FROM {{ this }}
)
{% endif %}
Запуск 1: все 5 заказов INSERT (новые строки).
Запуск 2 (после обновления статусов): заказы 1 и 2 UPDATE (существующие строки обновляются), новые заказы - INSERT.
SELECT order_id, status, COUNT(*) as cnt
FROM dbt_dev.fct_orders
GROUP BY order_id, status
ORDER BY order_id;
+--------+------------+---+
|order_id| status|cnt|
+--------+------------+---+
| 1| processing| 1| ← Обновлён (был pending)
| 2| delivered| 1| ← Обновлён (был pending)
| 3| processing| 1|
| 4| shipped| 1|
| 5| delivered| 1|
+--------+------------+---+
Нет дублей, актуальный статус. Каждый order_id - ровно одна строка.
Шаг 4: Анализ MERGE INTO в логах dbt¶
cat logs/dbt.log | grep -A 20 "MERGE INTO"
MERGE INTO dbt_dev.fct_orders AS DBT_INTERNAL_DEST
USING dbt_dev.fct_orders__dbt_tmp AS DBT_INTERNAL_SOURCE
ON (DBT_INTERNAL_SOURCE.order_id = DBT_INTERNAL_DEST.order_id)
WHEN MATCHED THEN UPDATE SET
`user_id` = DBT_INTERNAL_SOURCE.`user_id`,
`order_date` = DBT_INTERNAL_SOURCE.`order_date`,
`amount` = DBT_INTERNAL_SOURCE.`amount`,
`status` = DBT_INTERNAL_SOURCE.`status`,
`updated_at` = DBT_INTERNAL_SOURCE.`updated_at`,
`dbt_processed_at` = DBT_INTERNAL_SOURCE.`dbt_processed_at`
WHEN NOT MATCHED THEN INSERT (
`order_id`, `user_id`, `order_date`, `amount`,
`status`, `updated_at`, `dbt_processed_at`
) VALUES (
DBT_INTERNAL_SOURCE.`order_id`,
...
)
Тестирование incremental-моделей¶
Тест на идемпотентность¶
-- tests/assert_fct_orders_idempotent.sql
-- Запустить dbt run дважды подряд, потом проверить количество строк
-- Если модель идемпотентна, COUNT до и после повторного запуска одинаков
SELECT
order_id,
COUNT(*) AS cnt
FROM {{ ref('fct_orders') }}
GROUP BY order_id
HAVING COUNT(*) > 1
Если этот запрос возвращает строки - в таблице есть дубли по order_id. Тест должен всегда возвращать 0 строк.
dbt test --select fct_orders
Тест на полноту данных¶
-- tests/assert_orders_no_missing_days.sql
-- Проверка: нет пропущенных дней в диапазоне дат
WITH expected_dates AS (
SELECT EXPLODE(SEQUENCE(
DATE '2024-01-01',
CURRENT_DATE(),
INTERVAL 1 DAY
)) AS expected_date
),
actual_dates AS (
SELECT DISTINCT order_date FROM {{ ref('fct_orders') }}
)
SELECT e.expected_date
FROM expected_dates e
LEFT JOIN actual_dates a ON e.expected_date = a.order_date
WHERE a.order_date IS NULL
AND e.expected_date <= CURRENT_DATE()
Тест на freshness¶
-- tests/assert_orders_are_fresh.sql
SELECT
MAX(updated_at) AS max_updated_at,
CURRENT_TIMESTAMP() AS now,
TIMESTAMPDIFF(HOUR, MAX(updated_at), CURRENT_TIMESTAMP()) AS hours_lag
FROM {{ ref('fct_orders') }}
HAVING TIMESTAMPDIFF(HOUR, MAX(updated_at), CURRENT_TIMESTAMP()) > 4
Если данные старше 4 часов - тест проваливается.
Обработка late-arriving data¶
Проблема: события приходят с задержкой¶
Мобильные приложения могут работать в offline-режиме и отправлять события спустя дни после их возникновения. При append-стратегии с watermark WHERE event_ts > MAX(event_ts) такие события пропустятся навсегда - watermark уже "ушёл вперёд".
Решение для append: использовать поле загрузки¶
Вместо event_ts (время события) используйте loaded_at (время загрузки в Bronze):
{% if is_incremental() %}
WHERE loaded_at > (
SELECT COALESCE(MAX(loaded_at), TIMESTAMP '2000-01-01')
FROM {{ this }}
)
{% endif %}
Теперь все события, загруженные в Bronze после предыдущего запуска, попадут в таблицу - независимо от того, когда они реально произошли. Исторические агрегаты при этом могут быть неточными (событие попало в таблицу за "сегодня", хотя произошло "вчера"), но дубликатов нет.
Решение для insert_overwrite: скользящее окно¶
-- Пересчитываем последние 7 дней при каждом запуске
{% if is_incremental() %}
WHERE event_date >= DATE_ADD(CURRENT_DATE(), -7)
{% endif %}
Поздние события за последние 7 дней будут корректно учтены - их партиции перезапишутся. Для более старых поздних событий - отдельный процесс исторического пересчёта.
Решение для merge: широкое окно + unique_key¶
{% if is_incremental() %}
WHERE event_date >= DATE_ADD(CURRENT_DATE(), -30)
OR updated_at > (SELECT MAX(updated_at) FROM {{ this }})
{% endif %}
Merge обновит строки, которые изменились (по updated_at), и дополнительно перезапишет все строки за последние 30 дней (на случай если updated_at не обновляется корректно).
Compaction после incremental-загрузок¶
Проблема накопления мелких файлов¶
Каждый append создаёт новые файлы в таблице. Ежедневные загрузки → 365 файлов в год на партицию. Каждый merge может создавать новые файлы (rewrite) и помечать старые как удалённые. После множества merge'ей таблица содержит много мелких файлов.
Compaction для Delta Lake¶
-- macros/delta_optimize.sql
{% macro delta_optimize(table_name, zorder_columns=[]) %}
OPTIMIZE {{ table_name }}
{% if zorder_columns %}
ZORDER BY ({{ zorder_columns | join(', ') }})
{% endif %};
VACUUM {{ table_name }} RETAIN 168 HOURS;
{% endmacro %}
-- models/maintenance/optimize_fct_orders.sql
{{ config(materialized='table') }}
{{ delta_optimize('dbt_dev.fct_orders', ['user_id', 'order_date']) }}
OPTIMIZE объединяет мелкие файлы в файлы ~1 ГБ. ZORDER BY (user_id, order_date) применяет Z-order clustering - данные сортируются так, чтобы строки с близкими значениями user_id и order_date находились в одном файле, что ускоряет point-lookup запросы.
Compaction для Iceberg¶
CALL iceberg.system.rewrite_data_files(
table => 'analytics.fct_orders',
strategy => 'binpack',
options => map(
'target-file-size-bytes', '134217728',
'min-file-size-bytes', '67108864',
'max-file-size-bytes', '268435456',
'max-concurrent-file-group-rewrites', '10'
)
);
CALL iceberg.system.expire_snapshots(
table => 'analytics.fct_orders',
older_than => TIMESTAMP subtractinterval(CURRENT_TIMESTAMP(), 7, 'day'),
retain_last => 5
);
Схема эволюции incremental-моделей¶
Добавление колонки к incremental-модели¶
При добавлении новой колонки в SELECT incremental-модели поведение dbt зависит от формата:
- Iceberg: dbt выполняет
ALTER TABLE ADD COLUMN- безопасно, без перезаписи - Delta Lake: dbt также выполняет
ALTER TABLE ADD COLUMN- безопасно
Для старых форматов (Parquet) добавление колонки требует --full-refresh.
{{ config(
materialized='incremental',
incremental_strategy='merge',
unique_key='order_id',
file_format='delta',
on_schema_change='append_new_columns',
) }}
on_schema_change - параметр dbt для управления поведением при изменении схемы:
fail(default) - ошибка при несоответствии схемыappend_new_columns- добавить новые колонки, не трогая существующиеsync_all_columns- синхронизировать все колонки (добавить новые, удалить отсутствующие)ignore- игнорировать изменения схемы
Для production рекомендуется append_new_columns или fail - последний требует явного --full-refresh при изменении схемы, что предотвращает случайные breaking changes.
Антипаттерны¶
Антипаттерн 1: Полный скан source внутри is_incremental()¶
{% if is_incremental() %}
WHERE event_ts > (
SELECT MAX(event_ts) FROM {{ source('bronze', 'raw_events') }}
)
{% endif %}
Здесь watermark берётся из источника (Bronze), а не из целевой таблицы. Это означает: если источник обновился раньше, чем выполнился dbt, watermark уже сдвинулся, и часть данных пропустится. Всегда берите watermark из {{ this }}.
Антипаттерн 2: Неправильный unique_key при merge¶
{{ config(
unique_key='id',
) }}
SELECT
order_id AS id,
user_id,
order_date,
...
Alias id используется как unique_key, но в MERGE INTO используется оригинальное имя. Если оригинальное имя колонки order_id, а unique_key задан как id - merge не найдёт совпадений и превратится в INSERT для всех строк. Всегда используйте исходное имя колонки, не alias.
Антипаттерн 3: merge без партиционирования на большой таблице¶
{{ config(
materialized='incremental',
incremental_strategy='merge',
unique_key='user_id',
file_format='delta',
-- partition_by НЕ задан
) }}
Без партиционирования Delta/Iceberg не может применить file-skipping для scan целевой таблицы. Merge прочитает всю таблицу, чтобы найти совпадения по user_id. При терабайтной таблице это часы работы. Всегда добавляйте partition_by по полю, коррелирующему с watermark.
Антипаттерн 4: append для данных с обновлениями¶
{{ config(
incremental_strategy='append',
) }}
SELECT order_id, status, updated_at
FROM orders_source
{% if is_incremental() %}
WHERE updated_at > (SELECT MAX(updated_at) FROM {{ this }})
{% endif %}
Заказы с обновлённым статусом появятся в таблице дважды - со старым и новым статусом. Для данных с обновлениями нужен merge.
Антипаттерн 5: is_incremental() с оконными функциями по всей таблице¶
{% if is_incremental() %}
AND user_id IN (
SELECT DISTINCT user_id
FROM {{ this }}
WHERE event_date >= DATE_ADD(CURRENT_DATE(), -7)
)
{% endif %}
SELECT DISTINCT user_id FROM {{ this }} читает всю таблицу истории. Даже если нужны только последние 7 дней. Это лишний full scan. Используйте прямой watermark, а не подзапросы с IN.
Антипаттерн 6: Отсутствие COALESCE в watermark-подзапросе¶
{% if is_incremental() %}
WHERE event_ts > (SELECT MAX(event_ts) FROM {{ this }})
{% endif %}
Если таблица пустая (после ручной очистки), MAX(event_ts) = null. Условие event_ts > null никогда не выполняется → при первом запуске после очистки никакие данные не записываются.
Правильно:
WHERE event_ts > (
SELECT COALESCE(MAX(event_ts), TIMESTAMP '1970-01-01')
FROM {{ this }}
)
Чеклист выбора стратегии¶
Задайте себе вопросы:
1. Изменяются ли данные в источнике после создания?
- Нет →
append - Да (статус, атрибуты) →
mergeилиinsert_overwrite
2. Нужно ли обновлять конкретные строки или достаточно перезаписать партицию?
- Строки →
merge - Партиции →
insert_overwrite
3. Поддерживает ли формат MERGE?
- Delta Lake или Iceberg →
mergeвозможен - Parquet → только
appendилиinsert_overwrite
4. Насколько часто приходят поздние данные (late-arriving)?
- Редко →
appendс загрузочным watermark - Регулярно, до N дней →
insert_overwriteс окном N дней - Всегда →
mergeс широким окном
5. Критична ли идемпотентность?
- Да →
insert_overwriteилиmerge(оба идемпотентны при правильном unique_key) - Приемлем retry с дублями →
append(но нужна стратегия дедупликации)
Домашнее задание¶
Дан dbt-проект с моделью fct_user_sessions, работающей в режиме materialized='table'. Модель агрегирует сессии пользователей за год (365 ГБ данных), полный пересчёт занимает 4 часа.
-- Текущая модель (НЕ оптимальная)
{{ config(materialized='table') }}
SELECT
user_id,
session_id,
DATE(session_start) AS session_date,
session_start,
session_end,
TIMESTAMPDIFF(SECOND, session_start, session_end) AS duration_sec,
page_views,
clicks,
revenue
FROM {{ ref('stg_sessions') }}
Задание 1. Переведите модель на incremental с insert_overwrite:
- Настройте партиционирование по
session_date - Напишите фильтр в
{% if is_incremental() %}с окном 3 дня (для компенсации late data) - Добавьте конфиг для Delta Lake с оптимизацией
Задание 2. Как изменится время выполнения? Рассчитайте ожидаемый объём сканируемых данных при ежедневном запуске (365 ГБ / 365 дней × 3 дня окно = ?).
Задание 3. Опишите сценарий, когда insert_overwrite для этой модели не подойдёт, и нужен merge. Что должно измениться в бизнес-требованиях, чтобы insert_overwrite стал недостаточным?
Задание 4. Напишите dbt-тест на идемпотентность: запустите dbt run дважды подряд и проверьте, что количество строк и суммы агрегатов не изменились.
Задание 5. Откройте Spark UI после первого (полного) и второго (инкрементального) запусков. Сравните:
- Input Size (сколько данных прочитано)
- Shuffle Write (объём shuffle)
- Task Duration (сколько длился самый долгий task)
Сформулируйте вывод: на сколько процентов сократилось время выполнения при инкрементальной загрузке?