dbt Incremental Models: append, insert_overwrite и merge с unique_key

Три стратегии инкрементальной загрузки в dbt: append для иммутабельных логов, insert_overwrite для партиционированных агрегатов, merge для CDC и SCD на Delta/Iceberg.

platform

Концепция инкрементальной загрузки

Экономика вычислений в 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 только если:

  1. Таблица уже существует в базе данных (создана при предыдущем запуске)
  2. Запуск не является --full-refresh
  3. Материализация модели - 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 идемпотентна по своей природе. Повторный запуск за тот же период:

  1. Читает те же данные (те же 7 дней)
  2. Агрегирует - получает те же результаты
  3. Перезаписывает те же партиции теми же данными

Итог: таблица в точности том же состоянии. Это делает 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 - транзакционная операция. Она требует:

  1. Атомарности: либо все изменения применены, либо ни одно
  2. Изоляции: читатели не видят частично изменённые данные
  3. Консистентности: после 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'а:

  1. Job 0: чтение источника + фильтрация + вычисление MAX(event_ts) → запись в __dbt_tmp
  2. Job 1: INSERT INTO - чтение __dbt_tmp + запись в целевую таблицу

В Spark UI Job 1 должен быть быстрым: простая копия данных без shuffle (если не нужна сортировка для Parquet-файлов). Если видите тяжёлый shuffle в Job 1 - возможно, нужна настройка write.distribution-mode.

Как выглядит insert_overwrite в Spark UI

insert_overwrite выполняет:

  1. Job 0: чтение источника за N дней + фильтрация → запись в __dbt_tmp
  2. 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 данных

Бизнес-кейс

Пайплайн обработки заказов интернет-магазина. Особенность: статус заказа меняется - pendingprocessingshippeddelivered. Заказ создан вчера, оплачен сегодня. Нужна витрина с актуальным статусом каждого заказа.

Шаг 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)

Сформулируйте вывод: на сколько процентов сократилось время выполнения при инкрементальной загрузке?